mirror of
https://github.com/dadosfera/maestro.git
synced 2026-08-31 19:58:21 +00:00
Compare commits
347
Commits
release/07-29
...
v1.89.0
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
d54f381998 | ||
|
|
3fd586753e | ||
|
|
9b3894f9e4 | ||
|
|
4a5f8f679a | ||
|
|
a444f5e5ec | ||
|
|
0d30c1cf83 | ||
|
|
c70abfc826 | ||
|
|
db2c7d6c02 | ||
|
|
0a99ce1aa4 | ||
|
|
39c66f030a | ||
|
|
1223e21ac4 | ||
|
|
a5d78a97ae | ||
|
|
2b33c22149 | ||
|
|
e32787baff | ||
|
|
f449d8ebd9 | ||
|
|
02627d023a | ||
|
|
cf94c73648 | ||
|
|
6b3241281b | ||
|
|
fe1caa003e | ||
|
|
e980c58507 | ||
|
|
ff5f8672e2 | ||
|
|
d51456ecf7 | ||
|
|
0cb9da13d0 | ||
|
|
f198449c16 | ||
|
|
0d627b0451 | ||
|
|
ff8eb3197d | ||
|
|
abd4e6f90a | ||
|
|
901518d26c | ||
|
|
17a4c0ef82 | ||
|
|
4ab8147149 | ||
|
|
0086528d51 | ||
|
|
bf12598915 | ||
|
|
efb0b58648 | ||
|
|
8a9b58b612 | ||
|
|
837a9d7265 | ||
|
|
d0122a9c20 | ||
|
|
968f75b688 | ||
|
|
877cb9d281 | ||
|
|
59efdb6272 | ||
|
|
7b1049224c | ||
|
|
0369f10b5c | ||
|
|
ecca9106f0 | ||
|
|
efc1f49d92 | ||
|
|
2d46ac3213 | ||
|
|
271174176b | ||
|
|
dcf7aed51c | ||
|
|
7fbce8b5be | ||
|
|
049eca9700 | ||
|
|
3bcd78bb32 | ||
|
|
f603b64679 | ||
|
|
b640465624 | ||
|
|
98282471a9 | ||
|
|
8983430889 | ||
|
|
d3aeca12cd | ||
|
|
ce619942a3 | ||
|
|
71b3f278a5 | ||
|
|
c92542ed90 | ||
|
|
69a9d78642 | ||
|
|
bf4f3cfd8c | ||
|
|
01afa69fcb | ||
|
|
80f3913fd2 | ||
|
|
10293ad6a9 | ||
|
|
e4d0c9c3e6 | ||
|
|
21c57e5620 | ||
|
|
959210e354 | ||
|
|
f10953e949 | ||
|
|
13903b9bb9 | ||
|
|
2e3a13d421 | ||
|
|
fa17fc3001 | ||
|
|
b7171556b8 | ||
|
|
254a638392 | ||
|
|
74b3bd6b46 | ||
|
|
9a29ef5401 | ||
|
|
3f21faaa66 | ||
|
|
bf19d29a1d | ||
|
|
6523f707e3 | ||
|
|
8b0bf84d34 | ||
|
|
9ef4c51ba1 | ||
|
|
19521489fa | ||
|
|
2ce9aad005 | ||
|
|
051fb6e4dd | ||
|
|
2125884c6c | ||
|
|
e3099aa2b2 | ||
|
|
f4c9226ef9 | ||
|
|
12c61d9b5d | ||
|
|
ff3999a6aa | ||
|
|
209470482a | ||
|
|
99c2a9ecf5 | ||
|
|
c0f75d241f | ||
|
|
1e0fb78dff | ||
|
|
c26194554c | ||
|
|
b70d37423d | ||
|
|
4ddd5edcfd | ||
|
|
3bcbba9581 | ||
|
|
67d47a9642 | ||
|
|
6f9c967c96 | ||
|
|
b8bdc5beea | ||
|
|
269f70b309 | ||
|
|
93f452ae05 | ||
|
|
e39378229f | ||
|
|
61a4f724ef | ||
|
|
38a9e21f5f | ||
|
|
1a2153d62f | ||
|
|
d25bfd147c | ||
|
|
d4ba45fd03 | ||
|
|
82f1035a5e | ||
|
|
6c57bac235 | ||
|
|
cf8eed35a3 | ||
|
|
55fc85c544 | ||
|
|
72ed637640 | ||
|
|
890364f597 | ||
|
|
30eec733b2 | ||
|
|
30a41ba144 | ||
|
|
2dc032e7e7 | ||
|
|
7f5981731f | ||
|
|
66309c7bbe | ||
|
|
15048eaf8a | ||
|
|
0e169a3cbc | ||
|
|
b5d933eaf3 | ||
|
|
6fa9bf861a | ||
|
|
9a217dff57 | ||
|
|
e08734c97f | ||
|
|
cd4382c1ff | ||
|
|
54b75ce11b | ||
|
|
0e038d0b12 | ||
|
|
85234fe0dd | ||
|
|
c97a02cb17 | ||
|
|
878ec977b8 | ||
|
|
8006867bc2 | ||
|
|
784b0ef090 | ||
|
|
b736cddf07 | ||
|
|
3d6328fb0b | ||
|
|
4246c7495e | ||
|
|
3d40746ccb | ||
|
|
3ab2f8f27d | ||
|
|
b44d23552c | ||
|
|
8e63757738 | ||
|
|
3b84409003 | ||
|
|
0ce3822300 | ||
|
|
a2c7ce00db | ||
|
|
121e30ce45 | ||
|
|
85c8a4937d | ||
|
|
53920f2f3f | ||
|
|
8c34914806 | ||
|
|
1881d07c4a | ||
|
|
3b8310fdca | ||
|
|
52bda8ebe2 | ||
|
|
275a53dbd1 | ||
|
|
9cefdb226d | ||
|
|
007f3911ff | ||
|
|
8ac0a8a79f | ||
|
|
285de97375 | ||
|
|
e2a7d2b92b | ||
|
|
da23ad76db | ||
|
|
55b0961b82 | ||
|
|
d3c5c0fa63 | ||
|
|
5dbc644d1d | ||
|
|
8eddb9e1bf | ||
|
|
76485f929d | ||
|
|
8e0182aa50 | ||
|
|
a18bdccc09 | ||
|
|
b27298501d | ||
|
|
33ebc91826 | ||
|
|
6948156693 | ||
|
|
0aaa4384c3 | ||
|
|
1f9d0c29ec | ||
|
|
039c652b28 | ||
|
|
653d4f53b5 | ||
|
|
31f8c2c1a6 | ||
|
|
8b93d4e97b | ||
|
|
adeb022818 | ||
|
|
f61c241dde | ||
|
|
e326cab44d | ||
|
|
61109f8ae9 | ||
|
|
3f8dc5cabe | ||
|
|
7f5d157739 | ||
|
|
bce73fb11f | ||
|
|
9473e65deb | ||
|
|
9c55c22230 | ||
|
|
77b9acd2d0 | ||
|
|
e99306adba | ||
|
|
96b947ebdc | ||
|
|
fb521f53cd | ||
|
|
6f7436f33f | ||
|
|
865140e681 | ||
|
|
c4a664572a | ||
|
|
acb631e33d | ||
|
|
7a10f88113 | ||
|
|
e616061c21 | ||
|
|
466f8fb8cc | ||
|
|
e03b9e7a14 | ||
|
|
5989822263 | ||
|
|
2cc8f46418 | ||
|
|
64a3e2652e | ||
|
|
ea44a1cbb6 | ||
|
|
fc9c0b0991 | ||
|
|
a4b5a44e44 | ||
|
|
d99a6aa322 | ||
|
|
3f910f851a | ||
|
|
bd231382eb | ||
|
|
295f1f86ca | ||
|
|
0a5e8001f9 | ||
|
|
9c1979e17a | ||
|
|
b2700d4bb0 | ||
|
|
01c1087e07 | ||
|
|
fcf7fb054e | ||
|
|
288aaeabc0 | ||
|
|
918c3d7416 | ||
|
|
7dedb3bd33 | ||
|
|
1ff5589a2e | ||
|
|
d24e9a1d80 | ||
|
|
7d3ef1ef92 | ||
|
|
df3f2489f9 | ||
|
|
53245b0067 | ||
|
|
dd699614ae | ||
|
|
6919a2a8d0 | ||
|
|
5c29e07450 | ||
|
|
00304262ba | ||
|
|
bc931c6dd8 | ||
|
|
c444d6e956 | ||
|
|
f7efb757bf | ||
|
|
2f140d213a | ||
|
|
df674dd441 | ||
|
|
09dced9fcd | ||
|
|
165b533172 | ||
|
|
b1dc567394 | ||
|
|
a0a1303515 | ||
|
|
c21d977f7c | ||
|
|
399d3492d3 | ||
|
|
b0b557246e | ||
|
|
bb671a90d6 | ||
|
|
f6ababbe7a | ||
|
|
82290285d0 | ||
|
|
c490814a98 | ||
|
|
7616b1e32c | ||
|
|
3ba2c91893 | ||
|
|
e789076ed4 | ||
|
|
3b6ddaaee6 | ||
|
|
8073194604 | ||
|
|
e7f410831f | ||
|
|
cd2c53b5c5 | ||
|
|
adf1b3b97e | ||
|
|
d84b5e184b | ||
|
|
6d608a0457 | ||
|
|
b8be2c7803 | ||
|
|
24fce721e3 | ||
|
|
23a9a27db1 | ||
|
|
0f5ed50af9 | ||
|
|
5f8f6a64ab | ||
|
|
8800ac2736 | ||
|
|
bdb82c2ce4 | ||
|
|
6d9ecc3568 | ||
|
|
7764447adc | ||
|
|
576fdecf89 | ||
|
|
3de1e90fa8 | ||
|
|
5e90950660 | ||
|
|
3ed26e648f | ||
|
|
fff3523152 | ||
|
|
19e4daeea4 | ||
|
|
8624d3f016 | ||
|
|
4f4da5bebe | ||
|
|
f00d2bf41d | ||
|
|
dc1d1400d8 | ||
|
|
d4451153a3 | ||
|
|
06ba759ea3 | ||
|
|
dd0d08ad6e | ||
|
|
ec81082877 | ||
|
|
7f58dc7090 | ||
|
|
43aa379c06 | ||
|
|
c284f8753c | ||
|
|
3a8f2495c4 | ||
|
|
af248716ef | ||
|
|
a57ad41ad4 | ||
|
|
ffddceec3b | ||
|
|
b04d5bb402 | ||
|
|
53df3caf3f | ||
|
|
ab2e35b54f | ||
|
|
bfa77d8f7b | ||
|
|
5002d147ad | ||
|
|
cd55dc0dc4 | ||
|
|
f1d56e4c2c | ||
|
|
bc49f79cb3 | ||
|
|
00163ad894 | ||
|
|
a2fbeb97cc | ||
|
|
288a796f46 | ||
|
|
b5a1e93770 | ||
|
|
1ee49ceab8 | ||
|
|
6eaf9cf6d0 | ||
|
|
075ca747df | ||
|
|
9ead4588c1 | ||
|
|
92b64bb362 | ||
|
|
cb8798d871 | ||
|
|
be70f4da09 | ||
|
|
b436d7de7a | ||
|
|
6780167f2b | ||
|
|
2a3ab4228f | ||
|
|
caeaf62a9d | ||
|
|
8e605aa361 | ||
|
|
fa4267f54a | ||
|
|
47d8ca2760 | ||
|
|
f3c51e5328 | ||
|
|
dcf70fdaae | ||
|
|
70f663374d | ||
|
|
38f2317bb6 | ||
|
|
31e0ca6c91 | ||
|
|
dc063cbf3f | ||
|
|
a4c82ae4a8 | ||
|
|
11a65e11d1 | ||
|
|
73d47f0ff5 | ||
|
|
d36e532c8e | ||
|
|
b2a24e3a49 | ||
|
|
301b6e98ec | ||
|
|
381a401ebb | ||
|
|
ba53934068 | ||
|
|
1b6846532f | ||
|
|
216349f303 | ||
|
|
878c9cf450 | ||
|
|
a654baef13 | ||
|
|
df852449b2 | ||
|
|
04d761ca14 | ||
|
|
9b168612a4 | ||
|
|
a772804dcd | ||
|
|
8a399fbbcf | ||
|
|
94049d6e8e | ||
|
|
c28cc58ac7 | ||
|
|
56a655bb90 | ||
|
|
b3d4a6b028 | ||
|
|
c7bd1f761a | ||
|
|
09c045455e | ||
|
|
52d651126e | ||
|
|
2f11b1e78f | ||
|
|
885e9b71bb | ||
|
|
193031dfdd | ||
|
|
f626a9bb5e | ||
|
|
ac26ad9f44 | ||
|
|
fbe00af6dc | ||
|
|
bfa9929b69 | ||
|
|
bada7000f6 | ||
|
|
2c908d8d95 | ||
|
|
da35633d42 | ||
|
|
cc7e06013f | ||
|
|
5ec36a05b5 | ||
|
|
f61e46c7c1 | ||
|
|
5267fa6d2f | ||
|
|
b590c9402d | ||
|
|
8efa6790d6 | ||
|
|
a381f50e88 |
@@ -0,0 +1,12 @@
|
||||
node_modules
|
||||
dist
|
||||
.git
|
||||
*.log
|
||||
npm-debug.log*
|
||||
.DS_Store
|
||||
.env
|
||||
.env.*
|
||||
coverage
|
||||
.nyc_output
|
||||
*.tgz
|
||||
!protospack.tgz
|
||||
@@ -142,45 +142,11 @@ jobs:
|
||||
docker system prune --volumes -a -f
|
||||
docker system df
|
||||
|
||||
k8s-setup:
|
||||
needs: [extract_environment]
|
||||
uses: ./.github/workflows/k8s-setup.yml
|
||||
k8s-deploy:
|
||||
needs: [extract_environment, semantic_release, build_ecr_image]
|
||||
uses: ./.github/workflows/k8s-deploy.yml
|
||||
with:
|
||||
cloud: azure
|
||||
cloud: 'oracle'
|
||||
environment: ${{ needs.extract_environment.outputs.environment }}
|
||||
image: ${{ needs.semantic_release.outputs.new_release_version }}
|
||||
secrets: inherit
|
||||
|
||||
helmfile-deploy:
|
||||
needs: [extract_environment, semantic_release, build_ecr_image, k8s-setup]
|
||||
runs-on: [self-hosted, "prd-azure"]
|
||||
environment: ${{ needs.extract_environment.outputs.environment }}
|
||||
|
||||
steps:
|
||||
- name: Checkout code
|
||||
uses: actions/checkout@v3
|
||||
|
||||
- name: Set up Helm
|
||||
uses: azure/setup-helm@v1
|
||||
with:
|
||||
version: 'v3.9.0'
|
||||
|
||||
- name: Set up Python
|
||||
uses: actions/setup-python@v4
|
||||
with:
|
||||
python-version: '3.8'
|
||||
|
||||
- name: Install Helmfile
|
||||
run: |
|
||||
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
|
||||
mv helmfile /usr/local/bin/
|
||||
helmfile --version
|
||||
|
||||
- name: Install Helm Diff Plugin
|
||||
run: helm plugin install https://github.com/databus23/helm-diff || true
|
||||
|
||||
- name: Run Helmfile Apply
|
||||
env:
|
||||
ENV: ${{ needs.extract_environment.outputs.environment }}
|
||||
IMAGE_TAG: ${{ needs.semantic_release.outputs.new_release_version }}
|
||||
run: helmfile -f deploy/helmfiles/${ENV}.yaml sync --set image.tag=$IMAGE_TAG
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
name : K8s Setup kube config to deploy
|
||||
name : K8s deploy
|
||||
|
||||
on:
|
||||
workflow_call:
|
||||
@@ -13,12 +13,24 @@ on:
|
||||
required: true
|
||||
default: "prd"
|
||||
type: string
|
||||
image:
|
||||
description: "Image Tag"
|
||||
required: true
|
||||
type: string
|
||||
|
||||
jobs:
|
||||
azure:
|
||||
if: inputs.cloud == 'azure'
|
||||
runs-on: [self-hosted, "prd-azure"]
|
||||
steps:
|
||||
- name: Checkout code
|
||||
uses: actions/checkout@v3
|
||||
|
||||
- name: Set up Helm
|
||||
uses: azure/setup-helm@v1
|
||||
with:
|
||||
version: 'v3.9.0'
|
||||
|
||||
- name: Install Azure ClI
|
||||
run: |
|
||||
curl -sL https://aka.ms/InstallAzureCLIDeb | bash
|
||||
@@ -36,12 +48,45 @@ jobs:
|
||||
uses: azure/setup-kubectl@v1
|
||||
with:
|
||||
version: 'v1.30.1'
|
||||
|
||||
- name: Set up Python
|
||||
uses: actions/setup-python@v4
|
||||
with:
|
||||
python-version: '3.8'
|
||||
|
||||
oci:
|
||||
if: inputs.cloud == 'oci'
|
||||
- name: Install Helmfile
|
||||
run: |
|
||||
curl -fsSLO 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 || true
|
||||
|
||||
- name: Run Helmfile Apply
|
||||
env:
|
||||
ENV: ${{ inputs.environment }}
|
||||
IMAGE_TAG: ${{ inputs.image }}
|
||||
run: helmfile -f deploy/helmfiles/${ENV}.yaml sync --set image.tag=$IMAGE_TAG
|
||||
|
||||
oracle:
|
||||
if: inputs.cloud == 'oracle'
|
||||
runs-on: [self-hosted, "prd-oracle"]
|
||||
env:
|
||||
HOME: /home/runner
|
||||
steps:
|
||||
- name: Checkout code
|
||||
uses: actions/checkout@v3
|
||||
|
||||
- name: Set up Helm
|
||||
uses: azure/setup-helm@v1
|
||||
with:
|
||||
version: 'v3.9.0'
|
||||
|
||||
- name: Install OCI CLI
|
||||
env:
|
||||
HOME: /home/runner
|
||||
run: |
|
||||
bash -c "$(curl -L https://raw.githubusercontent.com/oracle/oci-cli/master/scripts/install/install.sh)" -- --accept-all-defaults
|
||||
echo "$HOME/bin" >> $GITHUB_PATH
|
||||
@@ -52,13 +97,27 @@ jobs:
|
||||
echo "${{ secrets.OCI_CONFIG }}" > ~/.oci/config
|
||||
echo "${{ secrets.OCI_PRIVATE_KEY }}" > ~/.oci/oci_api_key.pem
|
||||
chmod 600 ~/.oci/oci_api_key.pem
|
||||
|
||||
|
||||
- name: Set up Python
|
||||
uses: actions/setup-python@v4
|
||||
with:
|
||||
python-version: '3.8'
|
||||
|
||||
- name: Install Helmfile
|
||||
run: |
|
||||
curl -fsSLO 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 || true
|
||||
|
||||
- name: Authenticate with OKE cluster
|
||||
env:
|
||||
ENV: ${{ inputs.environment }}
|
||||
STG_CLUSTER_ID: "ocid1.cluster.oc1.sa-saopaulo-1.aaaaaaaagh3jvln52a3ebm3dodx6emmhv5bmfs7i7sv2k4zkbcbrzcl6v37q"
|
||||
PRD_CLUSTER_ID: "ocid1.cluster.oc1.sa-saopaulo-1.aaaaaaaanf3vptl6hc2tzd4enfd2hfpsht3wikxww5xejc3l7cwfm6l3sndq"
|
||||
HOME: /home/runner
|
||||
run: |
|
||||
if [ "$ENV" = "stg" ]; then
|
||||
CLUSTER_ID=$STG_CLUSTER_ID
|
||||
@@ -70,8 +129,9 @@ jobs:
|
||||
fi
|
||||
|
||||
oci ce cluster create-kubeconfig --cluster-id ${CLUSTER_ID} --file $HOME/.kube/config --region sa-saopaulo-1 --token-version 2.0.0 --kube-endpoint PRIVATE_ENDPOINT
|
||||
|
||||
- name: Setup kubectl
|
||||
uses: azure/setup-kubectl@v1
|
||||
with:
|
||||
version: 'v1.30.1'
|
||||
|
||||
- name: Run Helmfile Apply
|
||||
env:
|
||||
ENV: ${{ inputs.environment }}
|
||||
IMAGE_TAG: ${{ inputs.image }}
|
||||
run: helmfile -f deploy/helmfiles/${ENV}.yaml sync --set image.tag=$IMAGE_TAG
|
||||
@@ -66,13 +66,16 @@ jobs:
|
||||
|
||||
- name: Install Helmfile
|
||||
run: |
|
||||
wget https://github.com/helmfile/helmfile/releases/download/v0.148.0/helmfile_0.148.0_linux_amd64.tar.gz
|
||||
curl -fsSLO 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 || true
|
||||
- name: Debug Helm env
|
||||
run: |
|
||||
helm env
|
||||
echo "HOME=$HOME"
|
||||
ls -R $HOME/.local/share/helm || true
|
||||
|
||||
- name: Authenticate with OKE cluster
|
||||
env:
|
||||
|
||||
+4
-3
@@ -1,4 +1,5 @@
|
||||
FROM node:18.17-alpine AS base_image
|
||||
FROM node:20-alpine AS base_image
|
||||
RUN npm install -g npm@10.8.2
|
||||
|
||||
FROM base_image AS build_base
|
||||
WORKDIR /app
|
||||
@@ -21,7 +22,7 @@ ENV PUPPETEER_SKIP_CHROMIUM_DOWNLOAD=true \
|
||||
# run aws cli without mounting secret, because CI already has AWS credentials
|
||||
FROM build_base AS ci_image
|
||||
RUN aws codeartifact login --tool npm --namespace @dadosfera --repository dadosfera-npm --domain dadosfera --domain-owner 611330257153 --region us-east-1
|
||||
RUN npm ci
|
||||
RUN npm ci --ignore-scripts
|
||||
COPY . .
|
||||
|
||||
|
||||
@@ -36,7 +37,7 @@ FROM build_base AS dev
|
||||
RUN --mount=type=secret,id=aws,target=/root/.aws/credentials \
|
||||
aws codeartifact login --tool npm --namespace @dadosfera --repository dadosfera-npm --domain dadosfera --domain-owner 611330257153 --region us-east-1
|
||||
# flag --build-from-source is required to force-build sqlite3
|
||||
RUN npm ci
|
||||
RUN npm ci --ignore-scripts
|
||||
COPY . .
|
||||
ENTRYPOINT npm run start:dev
|
||||
|
||||
|
||||
@@ -0,0 +1,47 @@
|
||||
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,9 +2,9 @@
|
||||
<image src="./assets/maestro.svg" style="width:10rem">
|
||||
</p>
|
||||
|
||||
|
||||
# Maestro
|
||||
|
||||
|
||||
Maestro é a API principal da Dadosfera. É responsável pela comunicação do Frontend com nossos microsserviços.
|
||||
|
||||
```mermaid
|
||||
|
||||
@@ -48,6 +48,9 @@ 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
|
||||
@@ -74,6 +77,8 @@ spec:
|
||||
value: "logstash-pipelines.dadosfera.ai"
|
||||
- name: LOGGER_GELF_PORT
|
||||
value: "{{ .Values.maestro.logger_gelf_port }}"
|
||||
- name: LOGGER_CONSOLE_EXTRA
|
||||
value: "true"
|
||||
- name: NIMBUS_BASE_URL
|
||||
value: "http://nimbus-api"
|
||||
- name: NPM_TOKEN
|
||||
@@ -94,12 +99,22 @@ spec:
|
||||
value: {{ .Values.maestro.open_group_id }}
|
||||
- name: DEDICATED_PROXY
|
||||
value: {{ .Values.maestro.dedicated_proxy }}
|
||||
- name: COOKIE_SECRET
|
||||
value: {{ .Values.maestro.cookie_secret }}
|
||||
- name: REDIS_DATABASE
|
||||
value: "{{ .Values.maestro.redis_database }}"
|
||||
- name: REDIS_HOST
|
||||
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:
|
||||
@@ -120,3 +135,14 @@ 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,9 +4,18 @@ 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:
|
||||
|
||||
@@ -9,6 +9,9 @@ metadata:
|
||||
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 }}
|
||||
|
||||
@@ -38,3 +38,15 @@ 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
|
||||
|
||||
|
||||
@@ -0,0 +1,19 @@
|
||||
maestro:
|
||||
env: stg
|
||||
duc_url: duc.stg.dadosfera.ai
|
||||
pi_factory_url: pi-factory.stg.dadosfera.ai
|
||||
in_factory_url: in-factory.stg.dadosfera.ai
|
||||
tr_factory_url: in-factory.stg.dadosfera.ai
|
||||
open_customer_id: b3e3dfe5-b992-4586-a73c-c0b0c00f615d
|
||||
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
|
||||
|
||||
replicaCount: 1
|
||||
|
||||
affinity: null
|
||||
@@ -27,6 +27,9 @@ 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
|
||||
@@ -43,11 +46,16 @@ 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: "product-redis-prd.z4xvqj.0001.use1.cache.amazonaws.com"
|
||||
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
|
||||
@@ -60,7 +68,7 @@ affinity:
|
||||
requiredDuringSchedulingIgnoredDuringExecution:
|
||||
nodeSelectorTerms:
|
||||
- matchExpressions:
|
||||
- key: application
|
||||
- key: name
|
||||
operator: In
|
||||
values:
|
||||
- general
|
||||
- product
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
charts:
|
||||
releases:
|
||||
- name: maestro
|
||||
chart: ../helm-chart
|
||||
values:
|
||||
@@ -49,5 +49,8 @@ charts:
|
||||
value: dea2c27f-0973-4588-a2e0-9e31b64c7ffd
|
||||
- name: replicaCount
|
||||
value: 1
|
||||
# 10.70.0.0/16 internal network
|
||||
# 137.131.167.254/32 loadbalancer
|
||||
# 159.112.184.81/32 cluster ip for the uptime request ingest
|
||||
- name: maestro.restricted_ip
|
||||
value: "177.52.172.0/24,189.84.160.157/32,186.237.171.146/32,57.151.113.140/30"
|
||||
value: "177.52.172.0/24, 189.84.160.157/32, 186.237.171.146/32, 137.131.167.254/32, 10.70.0.0/16, 159.112.184.81/32, 10.244.0.0/16"
|
||||
|
||||
+23
-51
@@ -3,56 +3,28 @@ charts:
|
||||
chart: ../helm-chart
|
||||
values:
|
||||
- ../helm-chart/values.yaml
|
||||
set:
|
||||
- name: maestro.duc_url
|
||||
value: duc.stg.dadosfera.ai
|
||||
- name: hostname
|
||||
value: maestro.stg.dadosfera.ai
|
||||
- name: maestro.pi_factory_url
|
||||
value: pi-factory.stg.dadosfera.ai
|
||||
- name: maestro.in_factory_url
|
||||
value: in-factory.stg.dadosfera.ai
|
||||
- name: maestro.tr_factory_url
|
||||
value: in-factory.stg.dadosfera.ai
|
||||
- name: maestro.open_customer_id
|
||||
value: b3e3dfe5-b992-4586-a73c-c0b0c00f615d
|
||||
- name: maestro.open_group_id
|
||||
value: e3f98a2f-7748-4981-8505-7695c8ca8218
|
||||
- name: replicaCount
|
||||
value: 1
|
||||
- ../helm-chart/values-stg.yaml
|
||||
|
||||
|
||||
# Environment to test Network Policies
|
||||
# - name: private-maestro
|
||||
# chart: ../helm-chart
|
||||
# values:
|
||||
# - ../helm-chart/values.yaml
|
||||
# set:
|
||||
# - name: app_name
|
||||
# value: maestro-private
|
||||
# - name: maestro.env
|
||||
# value: stg
|
||||
# - name: maestro.duc_url
|
||||
# value: duc.stg.dadosfera.ai
|
||||
# - name: hostname
|
||||
# value: private-maestro.stg.dadosfera.ai
|
||||
# - name: maestro.pi_factory_url
|
||||
# value: pi-factory.dadosfera.ai
|
||||
# - name: maestro.in_factory_url
|
||||
# value: in-factory.stg.dadosfera.ai
|
||||
# - name: maestro.tr_factory_url
|
||||
# value: in-factory.dadosfera.ai
|
||||
# - name: maestro.open_customer_id
|
||||
# value: b3e3dfe5-b992-4586-a73c-c0b0c00f615d
|
||||
# - name: maestro.open_group_id
|
||||
# value: e3f98a2f-7748-4981-8505-7695c8ca8218
|
||||
# # Customer id
|
||||
# - name: maestro.dedicated_proxy
|
||||
# value: 14d52fd4-d83d-4cdd-be34-bf11cc28b3bd
|
||||
# - name: replicaCount
|
||||
# value: 1
|
||||
# - name: affinity
|
||||
# value: null
|
||||
# - name: resources
|
||||
# value: null
|
||||
# - name: maestro.restricted_ip
|
||||
# value: "57.151.113.140/30"
|
||||
- name: private-maestro
|
||||
chart: ../helm-chart
|
||||
values:
|
||||
- ../helm-chart/values.yaml
|
||||
- ../helm-chart/values-stg.yaml
|
||||
set:
|
||||
- name: app_name
|
||||
value: maestro-private
|
||||
- name: hostname
|
||||
value: private-maestro.stg.dadosfera.ai
|
||||
# Customer id
|
||||
- name: maestro.dedicated_proxy
|
||||
value: 14d52fd4-d83d-4cdd-be34-bf11cc28b3bd
|
||||
- name: replicaCount
|
||||
value: 1
|
||||
- name: affinity
|
||||
value: null
|
||||
- name: resources
|
||||
value: null
|
||||
- name: maestro.restricted_ip
|
||||
value: "137.131.167.254/32, 10.70.0.0/16, 159.112.184.81/32, 10.244.0.0/16"
|
||||
|
||||
+2742
-864
File diff suppressed because it is too large
Load Diff
Vendored
+2
@@ -15,6 +15,8 @@ declare global {
|
||||
OPEN_GROUP_ID: string;
|
||||
OPEN_CUSTOMER_ID: string;
|
||||
DEDICATED_PROXY: string;
|
||||
COOKIE_SECRET: string;
|
||||
REDIS_TLS?: string;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Generated
+2614
-1772
File diff suppressed because it is too large
Load Diff
+21
-5
@@ -10,7 +10,7 @@
|
||||
},
|
||||
"scripts": {
|
||||
"co:login": "aws codeartifact login --tool npm --namespace @dadosfera --repository dadosfera-npm --domain dadosfera --domain-owner 611330257153 --region us-east-1",
|
||||
"proto-update": "npm i @dadosfera/protospack-v2@latest --save-exact",
|
||||
"proto-update": "npm i @dadosfera/protospack-v2@v3.40.0-beta.1 --save-exact",
|
||||
"prebuild": "rimraf dist",
|
||||
"build": "nest build",
|
||||
"format": "prettier --write \"src/**/*.ts\" \"test/**/*.ts\"",
|
||||
@@ -27,10 +27,14 @@
|
||||
"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": "2.5.3",
|
||||
"@dadosfera/protospack-v2": "3.38.0-beta.10",
|
||||
"@dadosfera/protospack-v2": "^3.40.0-beta.10",
|
||||
"@grpc/grpc-js": "^1.9.3",
|
||||
"@grpc/proto-loader": "^0.7.9",
|
||||
"@nestjs/cli": "^9.5.0",
|
||||
@@ -44,11 +48,12 @@
|
||||
"@nestjs/schematics": "^9.2.0",
|
||||
"@nestjs/swagger": "^6.3.0",
|
||||
"@nestjs/testing": "^9.4.3",
|
||||
"axios": "^0.27.2",
|
||||
"axios": "0.30.3",
|
||||
"cache-manager": "^5.1.4",
|
||||
"cache-manager-ioredis-yet": "^1.1.0",
|
||||
"class-transformer": "^0.5.1",
|
||||
"class-validator": "^0.14.0",
|
||||
"cookie-parser": "^1.4.7",
|
||||
"cron-parser": "^4.9.0",
|
||||
"csv": "^6.3.11",
|
||||
"dotenv": "^14.3.2",
|
||||
@@ -59,6 +64,7 @@
|
||||
"jwk-to-pem": "^2.0.5",
|
||||
"mixpanel": "^0.17.0",
|
||||
"ms": "^3.0.0-canary.1",
|
||||
"multer": "^2.0.2",
|
||||
"openid-client": "^5.7.1",
|
||||
"passport": "^0.6.0",
|
||||
"passport-facebook": "^3.0.0",
|
||||
@@ -74,10 +80,17 @@
|
||||
"swagger-ui-express": "^4.6.3"
|
||||
},
|
||||
"overrides": {
|
||||
"multer": "1.4.5-lts.1"
|
||||
"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"
|
||||
},
|
||||
"devDependencies": {
|
||||
"@types/cache-manager": "^4.0.6",
|
||||
"@types/cookie-parser": "^1.4.9",
|
||||
"@types/express": "^4.17.17",
|
||||
"@types/express-session": "^1.18.1",
|
||||
"@types/jest": "27.0.2",
|
||||
@@ -104,5 +117,8 @@
|
||||
"ts-node": "^10.9.1",
|
||||
"tsconfig-paths": "^3.14.2",
|
||||
"typescript": "^4.9.5"
|
||||
},
|
||||
"resolutions": {
|
||||
"axios": "0.30.3"
|
||||
}
|
||||
}
|
||||
|
||||
+7
-2
@@ -17,7 +17,6 @@ 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';
|
||||
@@ -33,6 +32,10 @@ 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: [
|
||||
@@ -56,7 +59,6 @@ import { ApiKeyModule } from './modules/api-key/api-key.module';
|
||||
PermissionsModule,
|
||||
TermsOfUseModule,
|
||||
ConnectionTestModule,
|
||||
PipelinesModule,
|
||||
TransformationsModule,
|
||||
UsersModule,
|
||||
RolesModule,
|
||||
@@ -73,8 +75,11 @@ import { ApiKeyModule } from './modules/api-key/api-key.module';
|
||||
ApiKeyModule,
|
||||
IdentityProviderModule,
|
||||
NetworkPolicyModule,
|
||||
PlatformApiModule,
|
||||
StorageExplorerModule,
|
||||
//Always leave HealthModule last, so it is on the bottom of swagger
|
||||
HealthModule,
|
||||
ReleaseNoteModule,
|
||||
],
|
||||
})
|
||||
export class AppModule {}
|
||||
|
||||
@@ -153,6 +153,7 @@ 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,
|
||||
|
||||
@@ -0,0 +1,21 @@
|
||||
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,6 +116,44 @@ 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',
|
||||
@@ -319,6 +357,16 @@ 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',
|
||||
@@ -352,6 +400,25 @@ 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',
|
||||
@@ -611,6 +678,35 @@ 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;
|
||||
@@ -625,6 +721,7 @@ export const DADOSFERA_MODULES_KEYS = {
|
||||
DANGER_ZONE: 'danger-zone',
|
||||
PII: 'pii',
|
||||
EMBED: 'embedded-analytics',
|
||||
EMBED_ASSIGNED: 'embed-assigned',
|
||||
}
|
||||
|
||||
export const DADOSFERA_MODULES: Array<DadosferaModule> = [
|
||||
|
||||
@@ -12,6 +12,7 @@ export interface RequestUser {
|
||||
customer_tier: string;
|
||||
access_token: string;
|
||||
customer_modules: string[];
|
||||
roles: string[];
|
||||
}
|
||||
|
||||
export const User: (options?: { required?: boolean }) => ParameterDecorator =
|
||||
|
||||
@@ -0,0 +1,65 @@
|
||||
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);
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
+27
-4
@@ -9,6 +9,7 @@ import { AppModule } from './app.module';
|
||||
import { writeFileSync } from 'fs';
|
||||
import { execSync } from 'child_process';
|
||||
import { INestApplication } from '@nestjs/common';
|
||||
import cookieParser from 'cookie-parser';
|
||||
|
||||
async function bootstrap() {
|
||||
DadosferaLogger.setupLogger({
|
||||
@@ -17,20 +18,41 @@ 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: '*',
|
||||
origin: corsOrigins,
|
||||
methods: 'GET,HEAD,PUT,PATCH,POST,DELETE',
|
||||
preflightContinue: false,
|
||||
optionsSuccessStatus: 204,
|
||||
credentials: true,
|
||||
},
|
||||
});
|
||||
app.use(helmet());
|
||||
|
||||
if (process.env.ENV === 'prd') {
|
||||
app.use(helmet());
|
||||
app.use(cookieParser(process.env.COOKIE_SECRET));
|
||||
|
||||
if (process.env.ENV !== 'local') {
|
||||
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);
|
||||
@@ -89,3 +111,4 @@ function configureSwagger(app: INestApplication) {
|
||||
);
|
||||
}
|
||||
bootstrap();
|
||||
|
||||
|
||||
@@ -1,4 +1,4 @@
|
||||
import { Controller, Post, Body, Put, Get} from '@nestjs/common';
|
||||
import { Controller, Body, Put, Get, NotFoundException} from '@nestjs/common';
|
||||
import { AssignService } from './assign.service';
|
||||
import { CreateAssignDto } from './dto/create-assign.dto';
|
||||
import { Authenticated, RequireModule, RequireSomePermission } from 'src/decorators/authentication.decorator';
|
||||
@@ -15,7 +15,7 @@ export class AssignController {
|
||||
@RequireSomePermission(
|
||||
PERMISSIONS_GROUPS.USERS.permissions.ADMIN
|
||||
)
|
||||
@RequireModule(DADOSFERA_MODULES_KEYS.EMBED)
|
||||
@RequireModule(DADOSFERA_MODULES_KEYS.EMBED_ASSIGNED)
|
||||
create(@Body() createAssignDto: CreateAssignDto, @User() user: RequestUser) {
|
||||
const metadata = PackTheMetadata(user);
|
||||
return this.assignService.create(createAssignDto, metadata);
|
||||
@@ -25,9 +25,13 @@ export class AssignController {
|
||||
@RequireSomePermission(
|
||||
PERMISSIONS_GROUPS.USERS.permissions.ADMIN
|
||||
)
|
||||
@RequireModule(DADOSFERA_MODULES_KEYS.EMBED)
|
||||
@RequireModule(DADOSFERA_MODULES_KEYS.EMBED_ASSIGNED)
|
||||
async get(@User() user: RequestUser) {
|
||||
const metadata = PackTheMetadata(user);
|
||||
return await this.assignService.get(metadata);
|
||||
try {
|
||||
return await this.assignService.get(metadata);
|
||||
} catch (error) {
|
||||
throw new NotFoundException(error.message)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,6 +12,8 @@ import {
|
||||
Redirect,
|
||||
Req,
|
||||
Param,
|
||||
Res,
|
||||
UnauthorizedException,
|
||||
} from '@nestjs/common';
|
||||
import {
|
||||
ApiHeaders,
|
||||
@@ -47,12 +49,19 @@ import {
|
||||
} from './dtos/login';
|
||||
import { PackTheMetadata } from 'src/utils/PackTheMetadata';
|
||||
import { AuthGuard } from '@nestjs/passport';
|
||||
import { Request } from 'express';
|
||||
import { Request, Response } from 'express';
|
||||
import ErrorCodes, { OauthErrors } from 'src/utils/errorCodes';
|
||||
import jwt from 'jsonwebtoken';
|
||||
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';
|
||||
|
||||
type CookiesValues = {
|
||||
accessToken?: string;
|
||||
refreshToken?: string;
|
||||
userId?: string
|
||||
}
|
||||
|
||||
@ApiTags('Auth')
|
||||
@ApiHeaders([{ name: 'dadosfera-lang', enum: LanguageEnum, required: false }])
|
||||
@@ -66,6 +75,7 @@ export class AuthController {
|
||||
@Inject(DadosferaLogger)
|
||||
dadosferaLogger: DadosferaLogger,
|
||||
private authClient: AuthClientService,
|
||||
private apiKeyService: ApiKeyService,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
|
||||
@@ -86,11 +96,46 @@ export class AuthController {
|
||||
async signIn(
|
||||
@Body() { username, password, totp }: AuthSignInReq,
|
||||
@Language() language: LanguageEnum,
|
||||
): Promise<AuthSignInRes> {
|
||||
this.logger.info('/auth - SignIn');
|
||||
const metadata = PackTheMetadata({ language });
|
||||
this.logger.info('metadata: ' + JSON.stringify(metadata.toJSON()));
|
||||
return this.authClient.signIn({ username, password, totp }, metadata);
|
||||
@Res() res: Response,
|
||||
) {
|
||||
try {
|
||||
this.logger.info('/auth - SignIn');
|
||||
const metadata = PackTheMetadata({ language });
|
||||
this.logger.info('metadata: ' + JSON.stringify(metadata.toJSON()));
|
||||
const data = await this.authClient.signIn({ username, password, totp }, metadata);
|
||||
|
||||
if (data.tokens) {
|
||||
this.authClient.writeAuthSession(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);
|
||||
throw error;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@Post('sign-out')
|
||||
@HttpCode(HttpStatus.NO_CONTENT)
|
||||
async signOut(
|
||||
@Language() language: LanguageEnum,
|
||||
@Res() res: Response,
|
||||
) {
|
||||
try {
|
||||
this.logger.info('/auth - SignOut');
|
||||
|
||||
this.authClient.cleanUpAuthSession(res);
|
||||
|
||||
return res.send();
|
||||
} catch (error) {
|
||||
this.logger.error('/auth - SignIn - ERROR', error);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@Post('refresh-access-token')
|
||||
@@ -100,16 +145,26 @@ export class AuthController {
|
||||
@Body() body: AuthRefreshAccessTokenReq,
|
||||
@Language() language: LanguageEnum,
|
||||
@Headers('origin') origin: string,
|
||||
@Res() res: Response,
|
||||
) {
|
||||
this.logger.info('/auth - RefreshAccessToken');
|
||||
const { refreshToken, userId } = body;
|
||||
const frontHost = origin.replace(/^https?:\/\//, '');
|
||||
const { refreshToken, userId } = body;
|
||||
|
||||
const metadata = PackTheMetadata({
|
||||
language,
|
||||
custom_host: frontHost,
|
||||
});
|
||||
|
||||
return this.authClient.refreshAccessToken({ refreshToken, userId }, metadata);
|
||||
const data = await this.authClient.refreshAccessToken({ refreshToken, userId }, metadata);
|
||||
|
||||
this.authClient.writeAuthSession(res, {
|
||||
accessToken: data.accessToken,
|
||||
refreshToken: data.refreshToken,
|
||||
userId
|
||||
});
|
||||
|
||||
return res.send(data);
|
||||
}
|
||||
|
||||
@ApiInternalOnlyEndpoint()
|
||||
@@ -137,13 +192,14 @@ export class AuthController {
|
||||
) {
|
||||
this.logger.info('/auth - change-password');
|
||||
|
||||
const { oldPassword, newPassword } = body;
|
||||
const { oldPassword, newPassword, totpCode } = body;
|
||||
const { authorization: accessToken } = headers;
|
||||
|
||||
return this.authClient.changePassword({
|
||||
accessToken,
|
||||
oldPassword,
|
||||
newPassword,
|
||||
totpCode,
|
||||
});
|
||||
}
|
||||
|
||||
@@ -161,7 +217,8 @@ export class AuthController {
|
||||
|
||||
const { username } = body;
|
||||
|
||||
return this.authClient.resetPassword({ username }, metadata);
|
||||
await this.authClient.resetPassword({ username }, metadata);
|
||||
return { authProvider: process.env.AUTH_PROVIDER || 'cognito' };
|
||||
}
|
||||
|
||||
@ApiInternalOnlyEndpoint()
|
||||
@@ -418,4 +475,60 @@ export class AuthController {
|
||||
return this.authClient.resetUsers(body.users, metadata);
|
||||
}
|
||||
|
||||
@Get('me')
|
||||
async getMe(@Req() req: Request, @Res() res: Response) {
|
||||
this.logger.info('GET /auth/me ')
|
||||
this.logger.info(JSON.stringify(req.headers));
|
||||
|
||||
// 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,
|
||||
customer: {
|
||||
id: api_key.customer_id,
|
||||
name: api_key.customer_name,
|
||||
tier: api_key.customer_tier,
|
||||
}
|
||||
};
|
||||
|
||||
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.error('Invalid refresh token or customer name');
|
||||
throw new UnauthorizedException("Invalid refresh token or customer name");
|
||||
};
|
||||
|
||||
const {
|
||||
authSession,
|
||||
user
|
||||
} = await this.authClient.refreshUserSession(refreshToken, userId, resourceHost);
|
||||
this.authClient.writeAuthSession(res, authSession);
|
||||
return res.status(200).json(user);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -8,10 +8,11 @@ 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])],
|
||||
imports: [ClientsModule.register([client.providerOptions]), ApiKeyModule],
|
||||
controllers: [AuthController],
|
||||
providers: [
|
||||
AuthClientService,
|
||||
|
||||
@@ -1,10 +1,20 @@
|
||||
import { OnModuleInit, Inject, Injectable, ForbiddenException } from '@nestjs/common';
|
||||
import {
|
||||
OnModuleInit,
|
||||
Inject,
|
||||
Injectable,
|
||||
ForbiddenException,
|
||||
HttpException,
|
||||
HttpStatus,
|
||||
} 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, IdentityProviderProtoService } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/write-service';
|
||||
import {
|
||||
AuthProtoService as AuthServiceInterface,
|
||||
UsersProtoService,
|
||||
} from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/write-service';
|
||||
import {
|
||||
AuthSnowflakeSignInRequest,
|
||||
AuthSignInRequest,
|
||||
@@ -21,18 +31,24 @@ import {
|
||||
} from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/messages';
|
||||
import { DucClient } from '../duc/client.config';
|
||||
import { Metadata } from '@grpc/grpc-js';
|
||||
import { BulkEditResponse } from './dtos/login';
|
||||
import { BulkEditResponse, UserDTO } from './dtos/login';
|
||||
import jwt, { JwtPayload } from 'jsonwebtoken';
|
||||
import { PackTheMetadata } from 'src/utils/PackTheMetadata';
|
||||
import { Request, Response } from 'express';
|
||||
|
||||
type AuthSession = {
|
||||
accessToken?: string;
|
||||
refreshToken?: string;
|
||||
userId?: string;
|
||||
};
|
||||
|
||||
@Injectable()
|
||||
export class AuthClientService implements OnModuleInit {
|
||||
|
||||
|
||||
logger: DadosferaLogger;
|
||||
|
||||
|
||||
private authService: AuthServiceInterface;
|
||||
private identityProviderService: IdentityProviderProtoService;
|
||||
private userService: UsersProtoService;
|
||||
|
||||
constructor(
|
||||
@Inject(DadosferaLogger)
|
||||
dadosferaLogger: DadosferaLogger,
|
||||
@@ -46,8 +62,8 @@ export class AuthClientService implements OnModuleInit {
|
||||
ProtoServices.AuthProtoService,
|
||||
);
|
||||
|
||||
this.identityProviderService = this.grpcClient.getService<IdentityProviderProtoService>(
|
||||
ProtoServices.IdentityProviderProtoService,
|
||||
this.userService = this.grpcClient.getService<UsersProtoService>(
|
||||
ProtoServices.UsersProtoService,
|
||||
);
|
||||
}
|
||||
|
||||
@@ -63,11 +79,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) {
|
||||
@@ -75,7 +91,9 @@ 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();
|
||||
@@ -94,7 +112,6 @@ 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);
|
||||
@@ -105,7 +122,7 @@ export class AuthClientService implements OnModuleInit {
|
||||
this.checkDedicatedProxy(result);
|
||||
}
|
||||
|
||||
return result
|
||||
return result;
|
||||
}
|
||||
|
||||
async refreshAccessToken(
|
||||
@@ -115,7 +132,10 @@ export class AuthClientService implements OnModuleInit {
|
||||
this.logger.info('RefreshAccessToken');
|
||||
|
||||
return lastValueFrom(
|
||||
this.authService.AuthRefreshAccessToken({ refreshToken, userId }, metadata),
|
||||
this.authService.AuthRefreshAccessToken(
|
||||
{ refreshToken, userId },
|
||||
metadata,
|
||||
),
|
||||
);
|
||||
}
|
||||
|
||||
@@ -123,6 +143,7 @@ export class AuthClientService implements OnModuleInit {
|
||||
accessToken,
|
||||
oldPassword,
|
||||
newPassword,
|
||||
totpCode,
|
||||
}: AuthChangePasswordRequest) {
|
||||
this.logger.info('ChangePassword');
|
||||
|
||||
@@ -131,6 +152,7 @@ export class AuthClientService implements OnModuleInit {
|
||||
accessToken,
|
||||
oldPassword,
|
||||
newPassword,
|
||||
totpCode,
|
||||
}),
|
||||
);
|
||||
}
|
||||
@@ -283,4 +305,191 @@ 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,3 +140,17 @@ 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,6 +17,7 @@ import {
|
||||
HttpStatus,
|
||||
Res,
|
||||
} from '@nestjs/common';
|
||||
import { ValidationPipe } from '../../pipes/object-validation.pipe';
|
||||
import {
|
||||
ApiCreatedResponse,
|
||||
ApiHeaders,
|
||||
@@ -46,6 +47,7 @@ import {
|
||||
IMakeAComment,
|
||||
IOneDataAsset,
|
||||
IPreviewResponse,
|
||||
IUpdateCertificationStatusRequest,
|
||||
IUpdateDataRequest,
|
||||
TriggerCatalogReq,
|
||||
TriggerCatalogRes,
|
||||
@@ -241,6 +243,34 @@ 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('data-asset/:id')
|
||||
@RequireSomePermission(
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
|
||||
@@ -423,6 +453,7 @@ 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;
|
||||
|
||||
@@ -439,7 +470,7 @@ export class CatalogController {
|
||||
language,
|
||||
});
|
||||
|
||||
const docs = await this.catalogService.getDataDocs(id, metadata);
|
||||
const docs = await this.catalogService.getDataDocs(id, asset_type, metadata);
|
||||
|
||||
return { docs };
|
||||
}
|
||||
@@ -464,6 +495,8 @@ export class CatalogController {
|
||||
language,
|
||||
});
|
||||
|
||||
delete (body as any).certification_status;
|
||||
|
||||
const result = await this.catalogService.updateOneDataAsset({
|
||||
body,
|
||||
data_asset_id,
|
||||
@@ -477,6 +510,33 @@ 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,
|
||||
@@ -487,21 +547,32 @@ 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 } = user;
|
||||
const { user_id, customer_name, customer_id, username } = user;
|
||||
|
||||
const metadata = PackTheMetadata({
|
||||
customer_id,
|
||||
customer_name,
|
||||
user_id,
|
||||
username,
|
||||
});
|
||||
|
||||
this.logger.info(`/catalog - ON GET DATA DOCS ROUTE`, {
|
||||
this.logger.info(`/catalog - ON POST DATA DOCS ROUTE`, {
|
||||
user_id,
|
||||
customer_name,
|
||||
});
|
||||
|
||||
const res = await this.catalogService.createDataDocs({
|
||||
const body = {
|
||||
table_id,
|
||||
docs,
|
||||
asset_type,
|
||||
info: {
|
||||
customer: customer_name,
|
||||
},
|
||||
});
|
||||
}
|
||||
|
||||
const res = await this.catalogService.createDataDocs(body, metadata);
|
||||
|
||||
return res;
|
||||
}
|
||||
|
||||
@@ -5,20 +5,17 @@ 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,
|
||||
|
||||
@@ -24,9 +24,12 @@ import { CatalogClientConfiguration } from './catalog-client';
|
||||
import { UsersService } from '../users/users.service';
|
||||
import { RolesService } from '../roles/roles.service';
|
||||
import { Metadata } from '@grpc/grpc-js';
|
||||
import { PackTheMetadata } from 'src/utils/PackTheMetadata';
|
||||
import {
|
||||
AssetReporter,
|
||||
BatchRemoveRlsRulesRequest,
|
||||
CreateDataDocsDTO,
|
||||
IUpdateCertificationStatusRequest,
|
||||
IUpdateDataRequest,
|
||||
TriggerCatalogReq,
|
||||
} from './dtos';
|
||||
@@ -85,25 +88,23 @@ class CatalogService implements OnModuleInit {
|
||||
}
|
||||
|
||||
async getPiiReporter(metadata: Metadata, type: TypeParser) {
|
||||
this.logger.info('getPiiReporter: ' + type)
|
||||
this.logger.info('getPiiReporter: ' + type);
|
||||
try {
|
||||
const {
|
||||
data
|
||||
} = await lastValueFrom(
|
||||
this.catalogWriteService.GetPiiReporter({}, metadata)
|
||||
)
|
||||
this.logger.info("Finish grpc call")
|
||||
const { data } = await lastValueFrom(
|
||||
this.catalogWriteService.GetPiiReporter({}, metadata),
|
||||
);
|
||||
this.logger.info('Finish grpc call');
|
||||
|
||||
const parser = ParserBuilder.build<PiiMetadata>(type);
|
||||
|
||||
this.logger.info('parser file to: ' + type)
|
||||
const file = await parser.parse(data)
|
||||
this.logger.info('finish parser')
|
||||
this.logger.info('parser file to: ' + type);
|
||||
const file = await parser.parse(data);
|
||||
this.logger.info('finish parser');
|
||||
const mimeTypes: Record<TypeParser, string> = {
|
||||
'csv': 'text/csv',
|
||||
'html': 'text/html',
|
||||
'pdf': 'application/pdf'
|
||||
}
|
||||
csv: 'text/csv',
|
||||
html: 'text/html',
|
||||
pdf: 'application/pdf',
|
||||
};
|
||||
|
||||
const timestamp = new Date().toISOString().replace(/[:.]/g, '-');
|
||||
const filename = `relatorio-pii-${timestamp}.${type}`;
|
||||
@@ -111,13 +112,12 @@ class CatalogService implements OnModuleInit {
|
||||
return {
|
||||
file,
|
||||
filename: filename,
|
||||
type: mimeTypes[type]
|
||||
}
|
||||
type: mimeTypes[type],
|
||||
};
|
||||
} catch (error) {
|
||||
this.logger.error(error.message);
|
||||
throw error;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
async createDataAsset(data: Messages.CreateDataAssetRequest, metadata) {
|
||||
@@ -197,9 +197,11 @@ class CatalogService implements OnModuleInit {
|
||||
async getUserRolesIds(userId: string) {
|
||||
const result = await this.userService.findOneById(userId).catch(() => null);
|
||||
|
||||
const roles_ids = result.user.roles.map((role) => role.id);
|
||||
if (result) {
|
||||
return result.user.roles.map((role) => role.id);
|
||||
}
|
||||
|
||||
return roles_ids;
|
||||
return [];
|
||||
}
|
||||
|
||||
async searchDataAssets(
|
||||
@@ -207,10 +209,81 @@ class CatalogService implements OnModuleInit {
|
||||
metadata: Metadata,
|
||||
customer_id: string,
|
||||
) {
|
||||
this.logger.info('CatalogService - searchDataAssets');
|
||||
this.logger.info('CatalogService - searchDataAssets', { query });
|
||||
|
||||
const { search, page, size, sort_by, order, ...filters } = query;
|
||||
|
||||
this.logger.debug('Extracted filters:', { filters });
|
||||
|
||||
if (
|
||||
filters.manually !== undefined &&
|
||||
filters.manually !== null &&
|
||||
filters.manually !== ''
|
||||
) {
|
||||
filters.manually = Number(filters.manually); // 1 ou 0
|
||||
} else {
|
||||
delete filters.manually;
|
||||
}
|
||||
|
||||
if (filters.owner) {
|
||||
const { users: customer_users } =
|
||||
await this.userService.findAllUsersByCustomerId(customer_id);
|
||||
|
||||
this.logger.info('Available users in database count:', {
|
||||
count: customer_users.length,
|
||||
});
|
||||
this.logger.info('First 5 users:', {
|
||||
users: customer_users
|
||||
.slice(0, 5)
|
||||
.map((u) => ({ id: u.id, email: u.email, name: u.name })),
|
||||
});
|
||||
|
||||
const ownerValues = Array.isArray(filters.owner)
|
||||
? filters.owner
|
||||
: typeof filters.owner === 'string' && filters.owner.includes(',')
|
||||
? filters.owner.split(',').map((o: string) => o.trim())
|
||||
: [filters.owner];
|
||||
|
||||
this.logger.info('Owner values to convert:', {
|
||||
ownerValues,
|
||||
ownerFiltersOriginal: filters.owner,
|
||||
});
|
||||
|
||||
const ownerIds = ownerValues
|
||||
.map((ownerValue: string) => {
|
||||
const normalizedOwner = ownerValue.replace(/\s/g, '+');
|
||||
const user = customer_users.find((u) => {
|
||||
const isIdMatch = u.id === ownerValue;
|
||||
const isEmailMatch =
|
||||
u.email === ownerValue || u.email === normalizedOwner;
|
||||
const isNameMatch =
|
||||
u.name === ownerValue || u.name === normalizedOwner;
|
||||
this.logger.info('Comparing:', {
|
||||
userId: u.id,
|
||||
userEmail: u.email,
|
||||
userName: u.name,
|
||||
filterValue: ownerValue,
|
||||
normalizedFilter: normalizedOwner,
|
||||
idMatch: isIdMatch,
|
||||
emailMatch: isEmailMatch,
|
||||
nameMatch: isNameMatch,
|
||||
});
|
||||
return isIdMatch || isEmailMatch || isNameMatch;
|
||||
});
|
||||
this.logger.info('Looking for owner result:', {
|
||||
ownerValue,
|
||||
found: !!user,
|
||||
userId: user?.id,
|
||||
});
|
||||
return user?.id || ownerValue;
|
||||
})
|
||||
.filter((id: string) => id);
|
||||
|
||||
if (ownerIds.length > 0) {
|
||||
filters.owner = ownerIds;
|
||||
}
|
||||
}
|
||||
|
||||
const { data_assets, total } = await lastValueFrom(
|
||||
this.catalogReadService.GetAllDataAssets(
|
||||
{
|
||||
@@ -225,6 +298,8 @@ class CatalogService implements OnModuleInit {
|
||||
),
|
||||
);
|
||||
|
||||
console.log('MAESTRO RECEBEU RESPOSTA DO PI-FACTORY');
|
||||
|
||||
const result = JSON.parse(data_assets);
|
||||
|
||||
const response = await this.getAssetsUsersAndRoles(
|
||||
@@ -242,13 +317,13 @@ class CatalogService implements OnModuleInit {
|
||||
) {
|
||||
const data = await this.searchDataAssets(query, metadata, customer_id);
|
||||
|
||||
const formatData = data.data_assets.map(asset => ({
|
||||
const formatData = data.data_assets.map((asset) => ({
|
||||
id: asset.id,
|
||||
display_name: asset.display_name,
|
||||
data_asset_type: asset.data_asset_type,
|
||||
created_at: asset.created_at,
|
||||
tags: '[' + asset.tags.join(', ') + ']'
|
||||
}))
|
||||
tags: '[' + asset.tags.join(', ') + ']',
|
||||
}));
|
||||
|
||||
const parser = ParserBuilder.build<AssetReporter>('csv');
|
||||
|
||||
@@ -259,8 +334,8 @@ class CatalogService implements OnModuleInit {
|
||||
|
||||
return {
|
||||
file,
|
||||
filename
|
||||
}
|
||||
filename,
|
||||
};
|
||||
}
|
||||
|
||||
async getOneDataAsset(data: {
|
||||
@@ -310,6 +385,28 @@ class CatalogService implements OnModuleInit {
|
||||
return { data_asset: asset[0] };
|
||||
}
|
||||
|
||||
async updateCertificationStatus(data: {
|
||||
data_asset_id: string;
|
||||
body: IUpdateCertificationStatusRequest;
|
||||
metadata: Metadata;
|
||||
}) {
|
||||
const { body, data_asset_id, metadata } = data;
|
||||
|
||||
await lastValueFrom(
|
||||
this.catalogWriteService.UpdateDataAsset(
|
||||
{
|
||||
id: data_asset_id,
|
||||
changes: JSON.stringify({
|
||||
certification_status: body.certification_status,
|
||||
}),
|
||||
},
|
||||
metadata,
|
||||
),
|
||||
);
|
||||
|
||||
return { certification_status: body.certification_status };
|
||||
}
|
||||
|
||||
async updateOneDataAsset(data: {
|
||||
data_asset_id: string;
|
||||
customer_id: string;
|
||||
@@ -335,11 +432,11 @@ class CatalogService implements OnModuleInit {
|
||||
return { data_asset: asset[0] };
|
||||
}
|
||||
|
||||
async getDataDocs(id: string, metadata: Metadata) {
|
||||
async getDataDocs(id: string, assetType: string, metadata: Metadata) {
|
||||
const { documentation } = await lastValueFrom(
|
||||
this.catalogReadService.GetDatasetDoc({ id, type: undefined }, metadata),
|
||||
this.catalogReadService.GetDatasetDoc({ id }, metadata),
|
||||
);
|
||||
|
||||
|
||||
const docs = JSON.parse(documentation);
|
||||
return docs;
|
||||
}
|
||||
@@ -366,7 +463,16 @@ class CatalogService implements OnModuleInit {
|
||||
return result;
|
||||
}
|
||||
|
||||
async createDataDocs(body) {
|
||||
async createDataDocs(body: CreateDataDocsDTO, metadata: Metadata) {
|
||||
if (body.asset_type === 'table' || body.asset_type === 'view') {
|
||||
return this.createDataDocsViaNimbus(body);
|
||||
}
|
||||
|
||||
return this.createDataDocsViaGrpc(body, metadata);
|
||||
}
|
||||
|
||||
private async createDataDocsViaNimbus(body: CreateDataDocsDTO) {
|
||||
this.logger.info('Creating data docs via Nimbus for table/view');
|
||||
const nimbusUrl = this._getNimbusUrl(body);
|
||||
const { data } = await axios.post(
|
||||
`${nimbusUrl}/api/catalog/data-docs/`,
|
||||
@@ -375,6 +481,29 @@ class CatalogService implements OnModuleInit {
|
||||
return data;
|
||||
}
|
||||
|
||||
private async createDataDocsViaGrpc(body: CreateDataDocsDTO, metadata: Metadata) {
|
||||
this.logger.info('Creating data docs via gRPC for other asset types');
|
||||
try {
|
||||
const response: any = await lastValueFrom(
|
||||
this.catalogWriteService.UpdateDataAssetDoc(
|
||||
{
|
||||
id: body.table_id,
|
||||
docs: body.docs,
|
||||
},
|
||||
metadata,
|
||||
),
|
||||
);
|
||||
|
||||
return response;
|
||||
} catch (error) {
|
||||
this.logger.error('Error creating data asset docs:', error);
|
||||
throw new HttpException(
|
||||
'Failed to create data asset documentation',
|
||||
HttpStatus.INTERNAL_SERVER_ERROR,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
async findAllTags(data, metadata) {
|
||||
this.logger.info('CatalogService - findAllCustomerTags');
|
||||
|
||||
@@ -392,6 +521,23 @@ class CatalogService implements OnModuleInit {
|
||||
|
||||
return response;
|
||||
}
|
||||
|
||||
async findSchemas(metadata: Metadata) {
|
||||
this.logger.info('CatalogService - findSchemas');
|
||||
|
||||
try {
|
||||
const response = await lastValueFrom(
|
||||
this.catalogReadService.GetSchemas({}, metadata),
|
||||
);
|
||||
|
||||
return response;
|
||||
|
||||
} catch (error) {
|
||||
this.logger.error('Error fetching schemas:', error);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async getAssetsUsersAndRoles(data_assets: Array<any>, customer_id: string) {
|
||||
const { users: customer_users } =
|
||||
await this.userService.findAllUsersByCustomerId(customer_id);
|
||||
@@ -402,17 +548,20 @@ class CatalogService implements OnModuleInit {
|
||||
return data_assets.map((data_asset) => {
|
||||
const owner = customer_users.find(
|
||||
(u) => u.id === data_asset.owner,
|
||||
)?.username;
|
||||
)?.email;
|
||||
|
||||
const roles = [];
|
||||
const users = [];
|
||||
for (const role_id of data_asset.roles) {
|
||||
const data_asset_roles = data_asset?.roles || [];
|
||||
for (const role_id of data_asset_roles) {
|
||||
const role = customer_roles.find((r) => r.id === role_id);
|
||||
if (role) roles.push({ id: role.id, name: role.name });
|
||||
}
|
||||
for (const user_id of data_asset.users) {
|
||||
|
||||
const data_asset_users = data_asset?.users || [];
|
||||
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, username: user.username });
|
||||
if (user) users.push({ id: user.id, email: user.email });
|
||||
}
|
||||
return {
|
||||
...data_asset,
|
||||
@@ -502,27 +651,36 @@ class CatalogService implements OnModuleInit {
|
||||
|
||||
async createTableMetadata(body: any): Promise<number> {
|
||||
const nimbusUrl = this._getNimbusUrl(body);
|
||||
this.logger.info(`Nimbus URL: ${nimbusUrl}`, {...body.logMetadata});
|
||||
this.logger.info(`Nimbus URL: ${nimbusUrl}`, { ...body.logMetadata });
|
||||
|
||||
const endpoint = `${nimbusUrl}/api/catalog/table-metadata/`;
|
||||
|
||||
this.logger.info(`Creating table metadata for table ${body.table_metadata.table_name}`, {...body.logMetadata});
|
||||
this.logger.info(`Using endpoint: ${endpoint}`, {...body.logMetadata});
|
||||
this.logger.debug(`Payload: ${JSON.stringify(body.table_metadata)}`, {...body.logMetadata});
|
||||
this.logger.info(
|
||||
`Creating table metadata for table ${body.table_metadata.table_name}`,
|
||||
{ ...body.logMetadata },
|
||||
);
|
||||
this.logger.info(`Using endpoint: ${endpoint}`, { ...body.logMetadata });
|
||||
this.logger.debug(`Payload: ${JSON.stringify(body.table_metadata)}`, {
|
||||
...body.logMetadata,
|
||||
});
|
||||
|
||||
try {
|
||||
const { data, status } = await axios.post(endpoint, {...body.table_metadata});
|
||||
const { data, status } = await axios.post(endpoint, {
|
||||
...body.table_metadata,
|
||||
});
|
||||
|
||||
this.logger.info(
|
||||
`Table metadata created successfully with status ${status} for table ${body.table_metadata.table_name}`,
|
||||
{...body.logMetadata},
|
||||
{ ...body.logMetadata },
|
||||
);
|
||||
return data.id;
|
||||
} catch (error) {
|
||||
this.logger.error(
|
||||
`Failed to create table metadata for table ${body.table_metadata.table_name} failed with status ${
|
||||
error.response?.status
|
||||
} because of ${JSON.stringify(error.response?.data) || error.message}`, {...body.logMetadata});
|
||||
} because of ${JSON.stringify(error.response?.data) || error.message}`,
|
||||
{ ...body.logMetadata },
|
||||
);
|
||||
throw new Error(error.response?.data?.message || error.message);
|
||||
}
|
||||
}
|
||||
@@ -532,43 +690,58 @@ class CatalogService implements OnModuleInit {
|
||||
this.logger.info(`Nimbus URL: ${nimbusUrl}`, body.logMetadata);
|
||||
const endpoint = `${nimbusUrl}/api/catalog/column-metadata/`;
|
||||
|
||||
|
||||
|
||||
try {
|
||||
this.logger.info(`Creating column metadata for table ${body.column_metadata.table_name}`, {...body.logMetadata});
|
||||
this.logger.info(`Using endpoint: ${endpoint}`, {...body.logMetadata});
|
||||
this.logger.debug(`Payload: ${JSON.stringify(body.column_metadata)}`, {...body.logMetadata});
|
||||
const { data, status } = await axios.post(endpoint, body.column_metadata);
|
||||
this.logger.info(
|
||||
`Creating column metadata for table ${body.column_metadata.table_name}`,
|
||||
{ ...body.logMetadata },
|
||||
);
|
||||
this.logger.info(`Using endpoint: ${endpoint}`, { ...body.logMetadata });
|
||||
this.logger.debug(
|
||||
`Payload: ${JSON.stringify(body.column_metadata)}`,
|
||||
{ ...body.logMetadata },
|
||||
);
|
||||
const { data, status } = await axios.post(
|
||||
endpoint,
|
||||
body.column_metadata,
|
||||
);
|
||||
|
||||
this.logger.info(
|
||||
`Column metadata created successfully with status ${status} for table ${body.column_metadata.table_name}`,
|
||||
{...body.logMetadata},
|
||||
{ ...body.logMetadata },
|
||||
);
|
||||
return data.map((column) => column.id);
|
||||
} catch (error) {
|
||||
this.logger.error(
|
||||
`Failed to create column metadata failed with status for table ${body.column_metadata.table_name} ${
|
||||
error.response?.status
|
||||
} because of ${error.response?.data || error.message}`, {...body.logMetadata});
|
||||
} because of ${error.response?.data || error.message}`,
|
||||
{ ...body.logMetadata },
|
||||
);
|
||||
throw new Error(error.response?.data?.message || error.message);
|
||||
}
|
||||
}
|
||||
|
||||
async createDataPreview(body: any): Promise<number> {
|
||||
const nimbusUrl = this._getNimbusUrl(body);
|
||||
this.logger.info(`Nimbus URL: ${nimbusUrl}`, {...body.logMetadata});
|
||||
this.logger.info(`Nimbus URL: ${nimbusUrl}`, { ...body.logMetadata });
|
||||
const endpoint = `${nimbusUrl}/api/catalog/data-preview/`;
|
||||
|
||||
this.logger.info(`Creating data preview for table ${body.data_preview.table_name}`, {...body.logMetadata});
|
||||
this.logger.info(`Using endpoint: ${endpoint}`, {...body.logMetadata});
|
||||
this.logger.debug(`Payload: ${JSON.stringify(body.data_preview)}`, {...body.logMetadata});
|
||||
this.logger.info(
|
||||
`Creating data preview for table ${body.data_preview.table_name}`,
|
||||
{ ...body.logMetadata },
|
||||
);
|
||||
this.logger.info(`Using endpoint: ${endpoint}`, { ...body.logMetadata });
|
||||
this.logger.debug(
|
||||
`Payload: ${JSON.stringify(body.data_preview)}`,
|
||||
{ ...body.logMetadata },
|
||||
);
|
||||
|
||||
try {
|
||||
const { data, status } = await axios.post(endpoint, body.data_preview);
|
||||
|
||||
this.logger.info(
|
||||
`Data preview created successfully with status ${status} for table ${body.data_preview.table_name}`,
|
||||
{...body.logMetadata},
|
||||
{ ...body.logMetadata },
|
||||
);
|
||||
return data.id;
|
||||
} catch (error) {
|
||||
@@ -576,12 +749,70 @@ class CatalogService implements OnModuleInit {
|
||||
`Failed to create data preview for table ${body.data_preview.table_name} failed with status ${
|
||||
error.response?.status
|
||||
} because of ${error.response?.data || error.message}`,
|
||||
{...body.logMetadata},
|
||||
{ ...body.logMetadata },
|
||||
);
|
||||
throw new Error(error.response?.data?.message || error.message);
|
||||
}
|
||||
}
|
||||
|
||||
async renameTableOnNimbus(
|
||||
nimbusUrl: string,
|
||||
nimbusId: number,
|
||||
changes: { table_name?: string; table_schema?: string; display_name?: string },
|
||||
): Promise<void> {
|
||||
const endpoint = `${nimbusUrl}/api/catalog/table-metadata/${nimbusId}`;
|
||||
this.logger.info(`Renaming table-metadata ${nimbusId} on Nimbus`, { endpoint, changes });
|
||||
await axios.patch(endpoint, changes);
|
||||
}
|
||||
|
||||
async renameColumnMetadataOnNimbus(
|
||||
nimbusUrl: string,
|
||||
databaseName: string,
|
||||
oldTableName: string,
|
||||
oldTableSchema: string,
|
||||
newTableName: string,
|
||||
newTableSchema: string,
|
||||
): Promise<void> {
|
||||
const listEndpoint = `${nimbusUrl}/api/catalog/column-metadata/?database_name=${encodeURIComponent(databaseName)}&table_name=${encodeURIComponent(oldTableName)}&table_schema=${encodeURIComponent(oldTableSchema)}`;
|
||||
this.logger.info(`Fetching column-metadata records to rename`, { listEndpoint });
|
||||
const { data: columns } = await axios.get(listEndpoint);
|
||||
|
||||
const filtered = Array.isArray(columns) ? columns : [];
|
||||
|
||||
for (const column of filtered) {
|
||||
const patchEndpoint = `${nimbusUrl}/api/catalog/column-metadata/${column.id}`;
|
||||
await axios.patch(patchEndpoint, {
|
||||
table_name: newTableName,
|
||||
table_schema: newTableSchema,
|
||||
});
|
||||
}
|
||||
this.logger.info(`Renamed ${filtered.length} column-metadata records on Nimbus`);
|
||||
}
|
||||
|
||||
async renameDataPreviewOnNimbus(
|
||||
nimbusUrl: string,
|
||||
databaseName: string,
|
||||
oldTableName: string,
|
||||
oldTableSchema: string,
|
||||
newTableName: string,
|
||||
newTableSchema: string,
|
||||
): Promise<void> {
|
||||
const listEndpoint = `${nimbusUrl}/api/catalog/data-preview/?database_name=${encodeURIComponent(databaseName)}&table_name=${encodeURIComponent(oldTableName)}&table_schema=${encodeURIComponent(oldTableSchema)}`;
|
||||
this.logger.info(`Fetching data-preview records to rename`, { listEndpoint });
|
||||
const { data: previews } = await axios.get(listEndpoint);
|
||||
|
||||
const filtered = Array.isArray(previews) ? previews : [];
|
||||
|
||||
for (const preview of filtered) {
|
||||
const patchEndpoint = `${nimbusUrl}/api/catalog/data-preview/${preview.id}`;
|
||||
await axios.patch(patchEndpoint, {
|
||||
table_name: newTableName,
|
||||
table_schema: newTableSchema,
|
||||
});
|
||||
}
|
||||
this.logger.info(`Renamed ${filtered.length} data-preview records on Nimbus`);
|
||||
}
|
||||
|
||||
async catalogDatasetItem(table_metadata_id: number, metadata: Metadata) {
|
||||
const customer_name_raw = metadata.get('customer_name');
|
||||
|
||||
@@ -607,4 +838,4 @@ class CatalogService implements OnModuleInit {
|
||||
}
|
||||
}
|
||||
|
||||
export { CatalogService };
|
||||
export { CatalogService };
|
||||
|
||||
@@ -1,4 +1,5 @@
|
||||
import { ApiProperty, ApiPropertyOptional, PickType } from '@nestjs/swagger';
|
||||
import { IsEnum } from 'class-validator';
|
||||
import { CreateDataAssetRequest } from '@dadosfera/protospack-v2/dist/lib/Catalog/interfaces/messages';
|
||||
|
||||
export enum DataAssetShareType {
|
||||
@@ -6,6 +7,12 @@ 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',
|
||||
@@ -98,6 +105,8 @@ export class IDataAsset {
|
||||
embed?: EmbedObject;
|
||||
@ApiPropertyOptional({ enum: DataAssetShareType })
|
||||
share_type?: DataAssetShareType;
|
||||
@ApiPropertyOptional()
|
||||
docs?: string;
|
||||
}
|
||||
|
||||
export class IOneDataAsset {
|
||||
@@ -147,6 +156,24 @@ 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 {
|
||||
@@ -182,7 +209,16 @@ export class IUpdateDataRequest {
|
||||
embed: EmbedObject;
|
||||
@ApiPropertyOptional({ enum: DataAssetShareType })
|
||||
share_type?: DataAssetShareType;
|
||||
@ApiPropertyOptional()
|
||||
docs?: string;
|
||||
}
|
||||
|
||||
export class IUpdateCertificationStatusRequest {
|
||||
@ApiProperty({ enum: CertificationStatus })
|
||||
@IsEnum(CertificationStatus)
|
||||
certification_status: CertificationStatus;
|
||||
}
|
||||
|
||||
export class ICreateDataAsset implements CreateDataAssetRequest {
|
||||
@ApiProperty()
|
||||
display_name: string;
|
||||
@@ -196,6 +232,8 @@ export class ICreateDataAsset implements CreateDataAssetRequest {
|
||||
location: string;
|
||||
@ApiPropertyOptional()
|
||||
embed: EmbedObject;
|
||||
@ApiPropertyOptional()
|
||||
docs: string;
|
||||
}
|
||||
|
||||
export class IPreview {
|
||||
@@ -328,3 +366,10 @@ export type AssetReporter = {
|
||||
created_at: string;
|
||||
tags: string;
|
||||
}
|
||||
|
||||
export type CreateDataDocsDTO = {
|
||||
table_id: string;
|
||||
docs: string;
|
||||
asset_type: string;
|
||||
|
||||
}
|
||||
@@ -160,7 +160,7 @@ export class ShareService implements OnModuleInit {
|
||||
});
|
||||
|
||||
const { documentation } = await lastValueFrom(
|
||||
this.catalogReadService.GetDatasetDoc({ id, type: undefined }, metadata),
|
||||
this.catalogReadService.GetDatasetDoc({ id }, 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,
|
||||
)?.username;
|
||||
)?.email;
|
||||
|
||||
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, username: user.username });
|
||||
if (user) users.push({ id: user.id, email: user.email });
|
||||
}
|
||||
return {
|
||||
...data_asset,
|
||||
@@ -262,6 +262,7 @@ 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,
|
||||
|
||||
@@ -136,4 +136,32 @@ 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 { CustomerUpdateRequest } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/messages';
|
||||
import { CustomerSetLinksRequest } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/messages';
|
||||
import { CustomersProtoService } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/write-service';
|
||||
import { 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) {
|
||||
async getLinks(customerId: string): Promise<CustomerLinksConfig | null> {
|
||||
try {
|
||||
const result = await lastValueFrom(
|
||||
this.customerService.CustomerFindOneById({ id: customerId }),
|
||||
this.customerService.CustomerGetLinks({ customerId }),
|
||||
);
|
||||
return result.customer?.links || [];
|
||||
return (result.links as CustomerLinksConfig) || null;
|
||||
} catch (err) {
|
||||
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: Link[]) {
|
||||
async setLinks(customerId: string, links: CustomerLinksConfig) {
|
||||
if (!customerId || !links) {
|
||||
throw new HttpException(null, HttpStatus.BAD_REQUEST);
|
||||
}
|
||||
|
||||
try {
|
||||
return await firstValueFrom(
|
||||
this.customerService.CustomerUpdate({
|
||||
id: customerId,
|
||||
links,
|
||||
} as CustomerUpdateRequest),
|
||||
this.customerService.CustomerSetLinks({
|
||||
customerId,
|
||||
links: links as CustomerSetLinksRequest['links'],
|
||||
}),
|
||||
);
|
||||
} catch (err) {
|
||||
if (err.details === ErrorCodes.CUSTOMER.NOT_FOUND)
|
||||
@@ -223,4 +223,57 @@ 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,7 +1,6 @@
|
||||
import { Link } from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/entities';
|
||||
import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger';
|
||||
|
||||
export class CustomerLink implements Link {
|
||||
export class CustomerLinkItem {
|
||||
@ApiProperty()
|
||||
href: string;
|
||||
@ApiProperty()
|
||||
@@ -9,15 +8,59 @@ export class CustomerLink implements Link {
|
||||
@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: [CustomerLink] })
|
||||
links: CustomerLink[];
|
||||
@ApiProperty({ type: CustomerLinksConfig })
|
||||
links: CustomerLinksConfig;
|
||||
}
|
||||
|
||||
export class CustomerLinksResponse {
|
||||
@ApiProperty({ type: [CustomerLink] })
|
||||
links: CustomerLink[];
|
||||
}
|
||||
|
||||
@ApiPropertyOptional({ type: CustomerLinksConfig })
|
||||
links?: CustomerLinksConfig;
|
||||
}
|
||||
@@ -0,0 +1,27 @@
|
||||
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;
|
||||
}
|
||||
@@ -153,12 +153,27 @@ export class IdentityProviderController {
|
||||
throw new Error(ErrorCodes.IDENTITY_PROVIDER.INVALID_RESPONSE);
|
||||
}
|
||||
|
||||
const origin = req.headers['origin'] as string;
|
||||
try {
|
||||
const origin = req.headers['origin'] as string;
|
||||
this.logger.info('Header Origin: ' + origin);
|
||||
|
||||
const lang = language.substring(0, 2) + language.substring(2).toUpperCase();
|
||||
const callbackUrl = process.env.ENV !== "prd" ? `${origin}/auth/callback` : `${origin}/${lang}/auth/callback`;
|
||||
const lang =
|
||||
language.substring(0, 2) + language.substring(2).toUpperCase();
|
||||
const callbackUrl =
|
||||
process.env.ENV !== 'prd'
|
||||
? `${origin}/auth/callback`
|
||||
: `${origin}/${lang}/auth/callback`;
|
||||
|
||||
return await this.identityProviderService.getTokenByIdp(code, state, callbackUrl);
|
||||
this.logger.info('Callback URL: ' + callbackUrl);
|
||||
return await this.identityProviderService.getTokenByIdp(
|
||||
code,
|
||||
state,
|
||||
callbackUrl,
|
||||
);
|
||||
} catch (error) {
|
||||
this.logger.error(error);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
@Get('/links')
|
||||
@@ -166,41 +181,66 @@ export class IdentityProviderController {
|
||||
async providerLinks(@Req() req: Request) {
|
||||
this.logger.info('GET /identity-providers/links');
|
||||
|
||||
const frontDomain = req.headers['origin'] as string;
|
||||
try {
|
||||
const frontDomain = req.headers['origin'] as string;
|
||||
this.logger.info('Header Origin: ' + frontDomain);
|
||||
|
||||
if (!frontDomain) {
|
||||
throw new Error(ErrorCodes.IDENTITY_PROVIDER.INVALID_HEADER);
|
||||
if (!frontDomain) {
|
||||
this.logger.info('Not found front domain');
|
||||
throw new Error(ErrorCodes.IDENTITY_PROVIDER.INVALID_HEADER);
|
||||
}
|
||||
|
||||
const result =
|
||||
await this.identityProviderService.identityProvidersLinksPerDomain(
|
||||
frontDomain,
|
||||
);
|
||||
return result;
|
||||
} catch (error) {
|
||||
this.logger.error(error);
|
||||
throw error;
|
||||
}
|
||||
|
||||
const result =
|
||||
await this.identityProviderService.identityProvidersLinksPerDomain(
|
||||
frontDomain,
|
||||
);
|
||||
return result;
|
||||
}
|
||||
|
||||
@Get(':id')
|
||||
@HttpCode(HttpStatus.OK)
|
||||
@Redirect()
|
||||
async loginIdp(@Param('id') id: string, @Req() req: Request, @Language() language: LanguageEnum) {
|
||||
async loginIdp(
|
||||
@Param('id') id: string,
|
||||
@Req() req: Request,
|
||||
@Language() language: LanguageEnum,
|
||||
) {
|
||||
this.logger.info('GET /identity-providers/:id');
|
||||
|
||||
try {
|
||||
const frontDomain =
|
||||
(req.headers['origin'] as string) || (req.headers['referer'] as string);
|
||||
this.logger.info(`Front domain: ${frontDomain}`);
|
||||
const host =
|
||||
frontDomain.lastIndexOf('/') !== -1
|
||||
? frontDomain.substring(0, frontDomain.lastIndexOf('/'))
|
||||
: frontDomain;
|
||||
|
||||
const frontDomain = req.headers['origin'] as string || req.headers['referer'] as string;
|
||||
this.logger.info(`Front domain: ${frontDomain}`);
|
||||
const host = frontDomain.lastIndexOf('/') !== -1
|
||||
? frontDomain.substring(0, frontDomain.lastIndexOf('/'))
|
||||
: frontDomain;
|
||||
const lang =
|
||||
language.substring(0, 2) + language.substring(2).toUpperCase();
|
||||
const callbackUrl =
|
||||
process.env.ENV !== 'prd'
|
||||
? `${host}/auth/callback`
|
||||
: `${host}/${lang}/auth/callback`;
|
||||
|
||||
const lang = language.substring(0, 2) + language.substring(2).toUpperCase();
|
||||
const callbackUrl = process.env.ENV !== "prd" ? `${host}/auth/callback` : `${host}/${lang}/auth/callback`;
|
||||
this.logger.info('Callback URL: ' + callbackUrl);
|
||||
const redirectUrl =
|
||||
await this.identityProviderService.loginIdentityProvider(
|
||||
id,
|
||||
callbackUrl,
|
||||
);
|
||||
|
||||
const redirectUrl =
|
||||
await this.identityProviderService.loginIdentityProvider(id, callbackUrl);
|
||||
|
||||
this.logger.info(`Redirecting to: ${redirectUrl}`);
|
||||
return {
|
||||
url: redirectUrl,
|
||||
};
|
||||
this.logger.info(`Redirecting to: ${redirectUrl}`);
|
||||
return {
|
||||
url: redirectUrl,
|
||||
};
|
||||
} catch (error) {
|
||||
this.logger.error(error);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -14,16 +14,21 @@ import { Metadata } from '@grpc/grpc-js';
|
||||
import { Issuer, generators } from 'openid-client';
|
||||
import { SsoSignInDto } from './dto/sso-signin.dto';
|
||||
import { CacheService } from 'src/services/cache.service';
|
||||
import { Request } from 'express';
|
||||
import { CreateIdentityProvider } from './dto/identity-provider.dto';
|
||||
import DadosferaLogger from '@dadosfera/dadosfera-logs';
|
||||
|
||||
@Injectable()
|
||||
export class IdentityProviderService implements OnModuleInit {
|
||||
private logger: DadosferaLogger;
|
||||
private identityProviderService: IdentityProviderProtoService;
|
||||
constructor(
|
||||
@Inject(DucClient.name) private readonly grpcClient: ClientGrpc,
|
||||
private readonly cacheService: CacheService<SsoSignInDto>,
|
||||
) {}
|
||||
@Inject(DadosferaLogger)
|
||||
private dadosferaLoggger: DadosferaLogger
|
||||
) {
|
||||
this.logger = dadosferaLoggger.logger;
|
||||
}
|
||||
|
||||
onModuleInit() {
|
||||
this.identityProviderService =
|
||||
@@ -33,22 +38,26 @@ export class IdentityProviderService implements OnModuleInit {
|
||||
}
|
||||
|
||||
async create(body: IdentityProviderRequest, metadata: Metadata) {
|
||||
this.logger.info("Call IdentityProvider GRPC Create")
|
||||
return await lastValueFrom(
|
||||
this.identityProviderService.Create(body, metadata),
|
||||
);
|
||||
}
|
||||
|
||||
async getList(metadata: Metadata) {
|
||||
this.logger.info("Call IdentityProvider GRPC GetList")
|
||||
return await lastValueFrom(
|
||||
this.identityProviderService.GetList({}, metadata),
|
||||
);
|
||||
}
|
||||
|
||||
async loginIdentityProvider(id: string, callbackUrl: string) {
|
||||
this.logger.info("Call IdentityProvider GRPC FindIdentityProvider with: " + id);
|
||||
const idp = await lastValueFrom(
|
||||
this.identityProviderService.FindIdentityProvider({ id }),
|
||||
);
|
||||
|
||||
this.logger.info("Discovery issueURL: " + idp.issuerUrl)
|
||||
const issuer = await Issuer.discover(idp.issuerUrl);
|
||||
const client = new issuer.Client({
|
||||
client_id: idp.clientId,
|
||||
@@ -57,14 +66,19 @@ export class IdentityProviderService implements OnModuleInit {
|
||||
response_types: ['code'],
|
||||
});
|
||||
|
||||
this.logger.info("Generate Challenge")
|
||||
const code_verifier: string = generators.codeVerifier();
|
||||
const code_challenge: string = generators.codeChallenge(code_verifier);
|
||||
|
||||
this.logger.info("Generate State")
|
||||
const state = generators.state();
|
||||
|
||||
this.logger.info("Generate Nonce")
|
||||
const nonce = generators.nonce();
|
||||
|
||||
// Using state because it is returned in the callback
|
||||
// and we can use it to retrieve the code_verifier and nonce
|
||||
this.logger.info("Save Login parameters in redis")
|
||||
await this.cacheService.set(state, {
|
||||
codeVerifier: code_verifier,
|
||||
nonce,
|
||||
@@ -76,6 +90,7 @@ export class IdentityProviderService implements OnModuleInit {
|
||||
redirectUrls: idp.redirectUrls,
|
||||
});
|
||||
|
||||
this.logger.info("Generate Authorization URL")
|
||||
const url = client.authorizationUrl({
|
||||
scope: 'openid email',
|
||||
response_type: 'code',
|
||||
@@ -86,16 +101,20 @@ export class IdentityProviderService implements OnModuleInit {
|
||||
redirect_uri: callbackUrl,
|
||||
});
|
||||
const idpUrl = url + '&identity_provider=' + idp.name;
|
||||
this.logger.info(idpUrl)
|
||||
return idpUrl;
|
||||
}
|
||||
|
||||
async getTokenByIdp(code: string, state: string, callbackUrl: string) {
|
||||
this.logger.info("Get login parameters in redis")
|
||||
const ssoSign = await this.cacheService.get(state);
|
||||
|
||||
if (!ssoSign) {
|
||||
this.logger.info("Login Parameters Not Found")
|
||||
throw new BadRequestException('SSO sign-in is expired or not found');
|
||||
}
|
||||
|
||||
this.logger.info("Discovery Issue URL: " + ssoSign.issuerUrl)
|
||||
const issuer = await Issuer.discover(ssoSign.issuerUrl);
|
||||
const client = new issuer.Client({
|
||||
client_id: ssoSign.clientId,
|
||||
@@ -103,22 +122,22 @@ export class IdentityProviderService implements OnModuleInit {
|
||||
redirect_uris: ssoSign.redirectUrls,
|
||||
});
|
||||
|
||||
if (ssoSign === null) {
|
||||
throw new BadRequestException('SSO sign-in is expired or not found');
|
||||
}
|
||||
|
||||
const params = client.callbackParams(
|
||||
`${callbackUrl}?code=${code}&state=${state}`,
|
||||
);
|
||||
try {
|
||||
|
||||
this.logger.info("Get Token Set");
|
||||
const tokenSet = await client.callback(callbackUrl, params, {
|
||||
nonce: ssoSign.nonce,
|
||||
code_verifier: ssoSign.codeVerifier,
|
||||
state: ssoSign.state
|
||||
});
|
||||
|
||||
this.logger.info("Delete parameters in redis");
|
||||
await this.cacheService.delete(ssoSign.state);
|
||||
|
||||
this.logger.info("Call IdentityProvider GRPC SignInUser");
|
||||
return await lastValueFrom(
|
||||
this.identityProviderService.SignInUser({
|
||||
accessToken: tokenSet.access_token,
|
||||
@@ -128,12 +147,13 @@ export class IdentityProviderService implements OnModuleInit {
|
||||
}),
|
||||
);
|
||||
} catch (error) {
|
||||
console.log(error);
|
||||
this.logger.error(error);
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async deleteIdentityProvider(id: string, metadata: Metadata) {
|
||||
this.logger.info("Call IdentityProvider GRPC Delete with: " + id)
|
||||
return await lastValueFrom(
|
||||
this.identityProviderService.DeleteIdentityProvider({ id }, metadata),
|
||||
);
|
||||
@@ -144,6 +164,7 @@ export class IdentityProviderService implements OnModuleInit {
|
||||
body: CreateIdentityProvider,
|
||||
metadata: Metadata,
|
||||
) {
|
||||
this.logger.info("Call IdentityProvider GRPC Update with: " + id)
|
||||
return await lastValueFrom(
|
||||
this.identityProviderService.UpdateIdentityProvider(
|
||||
{
|
||||
@@ -156,6 +177,7 @@ export class IdentityProviderService implements OnModuleInit {
|
||||
}
|
||||
|
||||
async identityProvidersLinksPerDomain(frontDomain: string) {
|
||||
this.logger.info("Call IdentityProvider GRPC LinksPerDomain with: " + frontDomain)
|
||||
return await lastValueFrom(
|
||||
this.identityProviderService.GetProviderLinksFromDomain({ frontDomain }),
|
||||
);
|
||||
|
||||
@@ -11,10 +11,19 @@ export class TableColumns {
|
||||
name: string;
|
||||
@ApiProperty()
|
||||
columns: string[];
|
||||
@ApiProperty()
|
||||
@ApiPropertyOptional({ type: [Column] })
|
||||
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,4 +1,8 @@
|
||||
import { Info } from '@dadosfera/protospack/dist/lib/interfaces';
|
||||
export interface Info {
|
||||
user_id: string;
|
||||
customer_id: string;
|
||||
customer: string;
|
||||
}
|
||||
|
||||
interface Values {
|
||||
jdbc_user: string;
|
||||
|
||||
@@ -99,6 +99,7 @@ export class InputsController {
|
||||
customer: info.customer,
|
||||
});
|
||||
|
||||
this.logger.info(JSON.stringify(body))
|
||||
const response = await this.inputService.create({ body, info });
|
||||
|
||||
return response;
|
||||
|
||||
@@ -17,6 +17,8 @@ import {
|
||||
InputCreateGenericRequest,
|
||||
InputCreateS3Request,
|
||||
InputNewCreateRequest,
|
||||
InputUpdateResponse,
|
||||
RollbackInputRequest,
|
||||
TestConnectionRequest,
|
||||
} from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/messages';
|
||||
import { Info } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/entities';
|
||||
@@ -71,8 +73,8 @@ export class InputsService {
|
||||
objectCamelToSnake(createInputResponse);
|
||||
return createInputResponse;
|
||||
},
|
||||
update: async (updateInputDTO: UpdateInputRequest) => {
|
||||
this.logger.info('InputClientService - Update');
|
||||
update: async (updateInputDTO: UpdateInputRequest): Promise<InputUpdateResponse> => {
|
||||
this.logger.info('InputClientService - Update' + JSON.stringify(updateInputDTO));
|
||||
const updateInputResponse = await lastValueFrom(
|
||||
this.inputWriteService.InputUpdate(updateInputDTO),
|
||||
);
|
||||
@@ -164,6 +166,11 @@ 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,
|
||||
};
|
||||
@@ -200,23 +207,46 @@ export class InputsService {
|
||||
}
|
||||
|
||||
async update(id: string, data, info: Info) {
|
||||
this.validateCron({ ...data, info });
|
||||
// this.validateCron({ ...data, info });
|
||||
try {
|
||||
const updateInputResponse: any = await this.OLD_inputClient.update({
|
||||
const {
|
||||
tablesUpdate,
|
||||
dataAssetUpdate,
|
||||
input
|
||||
} = await this.OLD_inputClient.update({
|
||||
id,
|
||||
info,
|
||||
...data,
|
||||
info,
|
||||
});
|
||||
|
||||
updateInputResponse.input = this.adjustInputPayload(
|
||||
updateInputResponse?.input,
|
||||
const updateInputResponse = this.adjustInputPayload(
|
||||
input,
|
||||
);
|
||||
return updateInputResponse;
|
||||
return {
|
||||
input: updateInputResponse,
|
||||
tablesUpdate,
|
||||
dataAssetUpdate
|
||||
};
|
||||
} catch (err) {
|
||||
throw new HttpException(err.message, HttpStatus.NOT_FOUND);
|
||||
}
|
||||
}
|
||||
|
||||
async rollbackUpdate(
|
||||
data: RollbackInputRequest
|
||||
) {
|
||||
this.logger.info('PipelinesClientService - rollbackUpdate');
|
||||
this.logger.info('Rolling back input update with data: ' + JSON.stringify(data));
|
||||
const updatePipelineResponse = await lastValueFrom(
|
||||
this.inputWriteService.RollbackInputUpdate(
|
||||
data
|
||||
),
|
||||
);
|
||||
this.logger.info('Done');
|
||||
|
||||
return updatePipelineResponse;
|
||||
}
|
||||
|
||||
async remove(idRequest: IIdRequest) {
|
||||
return lastValueFrom(this.inputWriteService.InputRemove(idRequest));
|
||||
}
|
||||
@@ -258,4 +288,12 @@ export class InputsService {
|
||||
};
|
||||
return formatedPayload;
|
||||
}
|
||||
|
||||
async markTableDeleted(data: { input_id: string; table_name: string; info: Info }) {
|
||||
return lastValueFrom(this.inputWriteService.MarkTableDeleted(data));
|
||||
}
|
||||
|
||||
async unmarkTableDeleted(data: { input_id: string; table_name: string; info: Info }) {
|
||||
return lastValueFrom((this.inputWriteService as any).UnmarkTableDeleted(data));
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,29 +1,46 @@
|
||||
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, User } from 'src/decorators/user.decorator';
|
||||
import { RequestUser } 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(
|
||||
private mixpanelService: MixpanelService
|
||||
) {}
|
||||
@Inject(DadosferaLogger)
|
||||
dadosferaLogger: DadosferaLogger,
|
||||
private mixpanelService: MixpanelService,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
}
|
||||
|
||||
@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 };
|
||||
}
|
||||
|
||||
|
||||
}
|
||||
|
||||
@@ -1,78 +0,0 @@
|
||||
import { ConflictException, Inject, OnModuleInit } from '@nestjs/common';
|
||||
import { ClientGrpc } from '@nestjs/microservices';
|
||||
import {
|
||||
PipelineServicesNames,
|
||||
PipelinesServiceInterface,
|
||||
} from '@dadosfera/protospack';
|
||||
import { lastValueFrom } from 'rxjs';
|
||||
|
||||
import { IIdRequest } from './interfaces';
|
||||
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import { PipelinesClientConfiguration } from './pipelines-client';
|
||||
|
||||
export class PipelinesClientService implements OnModuleInit {
|
||||
private pipelineService: PipelinesServiceInterface;
|
||||
logger: DadosferaLogger;
|
||||
|
||||
constructor(
|
||||
@Inject(DadosferaLogger)
|
||||
dadosferaLogger: DadosferaLogger,
|
||||
@Inject(PipelinesClientConfiguration.name)
|
||||
private readonly grpcClient: ClientGrpc,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
}
|
||||
|
||||
onModuleInit() {
|
||||
this.pipelineService =
|
||||
this.grpcClient.getService<PipelinesServiceInterface>(
|
||||
PipelineServicesNames.PipelineService,
|
||||
);
|
||||
}
|
||||
|
||||
async getPipelineStatus(data) {
|
||||
this.logger.info('PipelinesClientService - GetPipelineStatus');
|
||||
|
||||
const statusPipelineResponse = await lastValueFrom(
|
||||
this.pipelineService.getPipelineStatus(data),
|
||||
)
|
||||
.then((res) => {
|
||||
const statusArray =
|
||||
res.status?.sort((a, b) => {
|
||||
if (a.id < b.id) {
|
||||
return 1;
|
||||
} else {
|
||||
return -1;
|
||||
}
|
||||
}) || [];
|
||||
return { status: statusArray };
|
||||
})
|
||||
.catch((err) => {
|
||||
this.logger.error(err.message);
|
||||
throw new Error(err);
|
||||
});
|
||||
this.logger.info('Done');
|
||||
|
||||
return statusPipelineResponse;
|
||||
}
|
||||
|
||||
async runPipeline({ id, info }: IIdRequest) {
|
||||
this.logger.info('PipelinesClientService - RunPipeline');
|
||||
const statusPipelineResponse = await lastValueFrom(
|
||||
this.pipelineService.triggerPipeline({ id, info }),
|
||||
).catch((err) => {
|
||||
this.logger.error(err.message);
|
||||
throw new Error(err);
|
||||
});
|
||||
|
||||
if (statusPipelineResponse.status == false) {
|
||||
throw new ConflictException(
|
||||
'This pipeline is not ready yet to execute, Try again later!',
|
||||
);
|
||||
}
|
||||
|
||||
this.logger.info('Done');
|
||||
return statusPipelineResponse;
|
||||
}
|
||||
}
|
||||
-36
@@ -1,36 +0,0 @@
|
||||
import { Info } from '@dadosfera/protospack/dist/lib/interfaces';
|
||||
|
||||
export interface ICreatePipelineDto {
|
||||
input: IdRequest;
|
||||
transformations: IdRequest[];
|
||||
output: IdRequest;
|
||||
tags: string[];
|
||||
name: string;
|
||||
description: string;
|
||||
info: Info;
|
||||
}
|
||||
|
||||
export interface IdRequest {
|
||||
id: string;
|
||||
}
|
||||
|
||||
export interface IIdRequest {
|
||||
id: string;
|
||||
info: Info;
|
||||
}
|
||||
|
||||
export interface IUpdatePipelineRequest {
|
||||
input: IdRequest;
|
||||
transformations: IdRequest[];
|
||||
output: IdRequest;
|
||||
tags: string[];
|
||||
name: string;
|
||||
description: string;
|
||||
id: string;
|
||||
info: Info;
|
||||
}
|
||||
|
||||
export interface IGetPipelineLogsRequest {
|
||||
id: string;
|
||||
details: string;
|
||||
}
|
||||
@@ -1,33 +0,0 @@
|
||||
import {
|
||||
ClientsProviderAsyncOptions,
|
||||
GrpcOptions,
|
||||
Transport,
|
||||
} from '@nestjs/microservices';
|
||||
import { PipelinePackages, PipelineProtoFilePath } from '@dadosfera/protospack';
|
||||
import { credentials } from '@grpc/grpc-js';
|
||||
|
||||
const isLocalConnection =
|
||||
process.env.PIFACTORY_URL.startsWith('pi-factory:') ||
|
||||
process.env.PIFACTORY_URL.includes('0.0.0.0');
|
||||
|
||||
export class PipelinesClientConfiguration {
|
||||
public name = 'PipelinesClientConfiguration';
|
||||
private config: GrpcOptions = {
|
||||
transport: Transport.GRPC,
|
||||
options: {
|
||||
url: process.env.PIFACTORY_URL,
|
||||
package: PipelinePackages,
|
||||
credentials: isLocalConnection ? undefined : credentials.createSsl(),
|
||||
protoPath: PipelineProtoFilePath,
|
||||
loader: {
|
||||
keepCase: true,
|
||||
enums: String,
|
||||
defaults: false,
|
||||
},
|
||||
},
|
||||
};
|
||||
providerOptions: ClientsProviderAsyncOptions = {
|
||||
name: this.name,
|
||||
...this.config,
|
||||
};
|
||||
}
|
||||
@@ -1,89 +0,0 @@
|
||||
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;
|
||||
}
|
||||
}
|
||||
@@ -1,19 +0,0 @@
|
||||
import { Module } from '@nestjs/common';
|
||||
import { ClientsModule } from '@nestjs/microservices';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
|
||||
import { PipelinesController } from './pipelines.controller';
|
||||
import { PipelinesService } from './pipelines.service';
|
||||
|
||||
import { PipelinesClientConfiguration } from './pipelines-client';
|
||||
import { PipelinesClientService } from './client.service';
|
||||
|
||||
const client = new PipelinesClientConfiguration();
|
||||
|
||||
@Module({
|
||||
imports: [ClientsModule.register([client.providerOptions])],
|
||||
controllers: [PipelinesController],
|
||||
providers: [PipelinesService, PipelinesClientService, DadosferaLogger],
|
||||
exports: [PipelinesService],
|
||||
})
|
||||
export class PipelinesModule {}
|
||||
@@ -1,33 +0,0 @@
|
||||
import { HttpException, HttpStatus, Injectable } from '@nestjs/common';
|
||||
import { PipelinesClientService } from './client.service';
|
||||
import { IIdRequest } from './interfaces';
|
||||
import { objectCamelToSnake } from 'src/utils/CaseConverter';
|
||||
|
||||
@Injectable()
|
||||
export class PipelinesService {
|
||||
constructor(private pipelineClient: PipelinesClientService) {}
|
||||
|
||||
async getPipelineStatus(data: IIdRequest) {
|
||||
try {
|
||||
const pipelineStatusResponse =
|
||||
await this.pipelineClient.getPipelineStatus(data);
|
||||
|
||||
return objectCamelToSnake(pipelineStatusResponse);
|
||||
} catch (err) {
|
||||
throw new HttpException(err.message, HttpStatus.NOT_FOUND);
|
||||
}
|
||||
}
|
||||
|
||||
async runPipeline({ id, info }: IIdRequest) {
|
||||
try {
|
||||
const triggerPipelineResponse = await this.pipelineClient.runPipeline({
|
||||
id,
|
||||
info,
|
||||
});
|
||||
|
||||
return objectCamelToSnake(triggerPipelineResponse);
|
||||
} catch (err) {
|
||||
throw new HttpException(err.message, HttpStatus.NOT_FOUND);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,13 @@
|
||||
import { ApiProperty, ApiPropertyOptional, OmitType } from '@nestjs/swagger';
|
||||
import { Info } from '@dadosfera/protospack/dist/lib/interfaces';
|
||||
|
||||
export class PipelineInputsDTO {
|
||||
@ApiProperty()
|
||||
tables: Array<{
|
||||
name: string,
|
||||
type: string,
|
||||
|
||||
}>
|
||||
}
|
||||
|
||||
export class IPipelineV2 {
|
||||
@ApiProperty()
|
||||
@@ -53,6 +61,12 @@ export interface IIdRequest {
|
||||
info: Info;
|
||||
}
|
||||
|
||||
export interface Info {
|
||||
user_id: string;
|
||||
customer_id: string;
|
||||
customer: string;
|
||||
}
|
||||
|
||||
export interface IUpdatePipelineRequest {
|
||||
input: IdRequest;
|
||||
transformations: IdRequest[];
|
||||
@@ -123,3 +137,30 @@ 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,7 +14,7 @@ import {
|
||||
Patch,
|
||||
HttpException,
|
||||
BadRequestException,
|
||||
CacheTTL,
|
||||
UseGuards,
|
||||
} from '@nestjs/common';
|
||||
import {
|
||||
ApiCreatedResponse,
|
||||
@@ -24,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';
|
||||
@@ -34,7 +34,6 @@ 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,52 +41,34 @@ 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 });
|
||||
|
||||
@@ -100,6 +81,7 @@ export class PipelinesController {
|
||||
}
|
||||
|
||||
@Post()
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.CREATE)
|
||||
@ApiCreatedResponse({ type: IPipelineV2 })
|
||||
async create(
|
||||
@Language() language: LanguageEnum,
|
||||
@@ -126,6 +108,7 @@ export class PipelinesController {
|
||||
}
|
||||
|
||||
@Get()
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.IMPORT_FILES.permissions.VIEW, PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async findAll(
|
||||
@User() user: RequestUser,
|
||||
@Language() language: LanguageEnum,
|
||||
@@ -150,6 +133,7 @@ export class PipelinesController {
|
||||
}
|
||||
|
||||
@Get('/download-logs')
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async downloadLogs(
|
||||
@User() user: RequestUser,
|
||||
@Language() language: LanguageEnum,
|
||||
@@ -180,6 +164,7 @@ 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,
|
||||
@@ -191,6 +176,7 @@ export class PipelinesController {
|
||||
}
|
||||
|
||||
@Get(':id/objects')
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getPipelineObjects(
|
||||
@Language() language: LanguageEnum,
|
||||
@User() user: RequestUser,
|
||||
@@ -202,23 +188,23 @@ 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(
|
||||
process.env.DEV_URL + `/pipeline/${id} - ON GET PIPELINE STATUS ROUTE`,
|
||||
{
|
||||
user: body.info.user_id,
|
||||
customer: body.info.customer,
|
||||
},
|
||||
);
|
||||
this.logger.info(`/pipeline/${id} - ON GET PIPELINE STATUS ROUTE`, {
|
||||
user: body.info.user_id,
|
||||
customer: body.info.customer,
|
||||
});
|
||||
|
||||
const response = await this.oldPipelinesService.getPipelineStatus(body);
|
||||
const response = await this.pipelinesClientService.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,
|
||||
@@ -236,30 +222,56 @@ export class PipelinesController {
|
||||
language,
|
||||
});
|
||||
|
||||
const result = await this.pipelinesClientService
|
||||
.findOne({ id }, metadata)
|
||||
.then((res) => {
|
||||
//{pipeline:{tables: {tables: [], input_id: ''}}}
|
||||
let tables = JSON.parse(res.pipeline.config.tables);
|
||||
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;
|
||||
});
|
||||
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);
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
@Patch('/:id')
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||
async update(
|
||||
@Language() language: LanguageEnum,
|
||||
@Body() updatePipelineDto,
|
||||
@@ -293,12 +305,54 @@ 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,
|
||||
@@ -311,9 +365,25 @@ 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({
|
||||
@@ -327,7 +397,7 @@ export class PipelinesController {
|
||||
|
||||
@ApiInternalOnlyEndpoint()
|
||||
@Post('/init-upload')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.CREATE)
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.IMPORT_FILES.permissions.VIEW)
|
||||
async initUploadFile(
|
||||
@User() user: RequestUser,
|
||||
@Body() body: IInitUploadCSVFile,
|
||||
@@ -359,7 +429,7 @@ export class PipelinesController {
|
||||
|
||||
@ApiInternalOnlyEndpoint()
|
||||
@Post('/complete-upload')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.CREATE)
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.IMPORT_FILES.permissions.VIEW)
|
||||
async completeUploadFile(
|
||||
@User() user: RequestUser,
|
||||
@Body() body: ICompleteUploadCSVFile,
|
||||
@@ -377,7 +447,7 @@ export class PipelinesController {
|
||||
|
||||
@ApiInternalOnlyEndpoint()
|
||||
@Post('/file')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.CREATE)
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.IMPORT_FILES.permissions.VIEW)
|
||||
async uploadedFile(
|
||||
@User() user: RequestUser,
|
||||
@Body() body: ICreatePipelineCSVFile,
|
||||
@@ -427,6 +497,7 @@ export class PipelinesController {
|
||||
|
||||
@ApiInternalOnlyEndpoint()
|
||||
@Post('start/:id')
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.CREATE)
|
||||
async activate(@Param('id') id: string, @Body() body) {
|
||||
const { info } = body;
|
||||
|
||||
@@ -438,7 +509,7 @@ export class PipelinesController {
|
||||
},
|
||||
);
|
||||
|
||||
const response = await this.oldPipelinesService.runPipeline({ id, info });
|
||||
const response = await this.pipelinesClientService.runPipeline({ id, info });
|
||||
|
||||
return response;
|
||||
}
|
||||
|
||||
@@ -7,23 +7,28 @@ 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],
|
||||
providers: [PipelinesService, DadosferaLogger, NimbusService],
|
||||
exports: [PipelinesService],
|
||||
})
|
||||
export class PipelinesV2Module {}
|
||||
|
||||
@@ -1,5 +1,7 @@
|
||||
/* eslint-disable no-async-promise-executor */
|
||||
import {
|
||||
BadRequestException,
|
||||
ConflictException,
|
||||
HttpException,
|
||||
HttpStatus,
|
||||
Inject,
|
||||
@@ -16,7 +18,7 @@ import { lastValueFrom } from 'rxjs';
|
||||
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import { PipelinesClientConfiguration } from './pipelines-client';
|
||||
import { ICreatePipelineV2Req } from './interfaces';
|
||||
import { ICreatePipelineV2Req, IIdRequest, UpdatePlatformInputRequest, UpdateTableDTO } 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';
|
||||
@@ -26,6 +28,17 @@ 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;
|
||||
@@ -39,6 +52,9 @@ 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;
|
||||
}
|
||||
@@ -138,6 +154,7 @@ 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;
|
||||
@@ -160,6 +177,15 @@ 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 = {
|
||||
@@ -339,4 +365,317 @@ 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;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -0,0 +1,11 @@
|
||||
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
@@ -0,0 +1,9 @@
|
||||
import { ApiProperty } from "@nestjs/swagger";
|
||||
|
||||
export class ValidationTableDTO {
|
||||
@ApiProperty()
|
||||
tables: Array<{
|
||||
table_name: string;
|
||||
table_schema: string;
|
||||
}>
|
||||
}
|
||||
@@ -0,0 +1,19 @@
|
||||
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 {}
|
||||
@@ -0,0 +1,131 @@
|
||||
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);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,14 @@
|
||||
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;
|
||||
};
|
||||
|
||||
@@ -0,0 +1,21 @@
|
||||
import { Test, TestingModule } from '@nestjs/testing';
|
||||
import { ReleaseNoteController } from './release_note.controller';
|
||||
import { ReleaseNoteService } from './release_note.service';
|
||||
import DadosferaLogger from '@dadosfera/dadosfera-logs';
|
||||
|
||||
describe('ReleaseNoteController', () => {
|
||||
let controller: ReleaseNoteController;
|
||||
|
||||
beforeEach(async () => {
|
||||
const module: TestingModule = await Test.createTestingModule({
|
||||
controllers: [ReleaseNoteController],
|
||||
providers: [ReleaseNoteService, DadosferaLogger],
|
||||
}).compile();
|
||||
|
||||
controller = module.get<ReleaseNoteController>(ReleaseNoteController);
|
||||
});
|
||||
|
||||
it('should be defined', () => {
|
||||
expect(controller).toBeDefined();
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,26 @@
|
||||
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);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,10 @@
|
||||
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 {}
|
||||
@@ -0,0 +1,19 @@
|
||||
import { Test, TestingModule } from '@nestjs/testing';
|
||||
import { ReleaseNoteService } from './release_note.service';
|
||||
import DadosferaLogger from '@dadosfera/dadosfera-logs';
|
||||
|
||||
describe('ReleaseNoteService', () => {
|
||||
let service: ReleaseNoteService;
|
||||
|
||||
beforeEach(async () => {
|
||||
const module: TestingModule = await Test.createTestingModule({
|
||||
providers: [ReleaseNoteService, DadosferaLogger],
|
||||
}).compile();
|
||||
|
||||
service = module.get<ReleaseNoteService>(ReleaseNoteService);
|
||||
});
|
||||
|
||||
it('should be defined', () => {
|
||||
expect(service).toBeDefined();
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,46 @@
|
||||
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,27 +229,33 @@ 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 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,
|
||||
};
|
||||
const newRole: RoleDto = this.formatRole(role);
|
||||
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,
|
||||
};;
|
||||
}
|
||||
}
|
||||
|
||||
@@ -0,0 +1,13 @@
|
||||
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),
|
||||
};
|
||||
@@ -0,0 +1,383 @@
|
||||
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 },
|
||||
);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,13 @@
|
||||
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 {}
|
||||
@@ -0,0 +1,178 @@
|
||||
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,4 +124,24 @@ 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,6 +44,18 @@ 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);
|
||||
|
||||
+5
-2
@@ -1,5 +1,8 @@
|
||||
import { Info } from '@dadosfera/protospack/dist/lib/interfaces';
|
||||
|
||||
export interface Info {
|
||||
user_id: string;
|
||||
customer_id: string;
|
||||
customer: string;
|
||||
}
|
||||
export interface ICreateTransformationsRequest {
|
||||
transformations: Transformation[];
|
||||
info: Info;
|
||||
|
||||
@@ -38,6 +38,12 @@ export class User {
|
||||
department?: string;
|
||||
@ApiProperty()
|
||||
hierarchy?: string;
|
||||
@ApiProperty()
|
||||
bio?: string;
|
||||
@ApiProperty()
|
||||
companyName?: string;
|
||||
@ApiProperty()
|
||||
personalSite?: string;
|
||||
@ApiPropertyOptional()
|
||||
customer?: Customer;
|
||||
@ApiProperty()
|
||||
@@ -58,6 +64,8 @@ export class UserNoRolesAndCustomer extends OmitType(UserNoRoles, [
|
||||
export class IUserByCustomer extends OmitType(User, ['customer']) {
|
||||
@ApiPropertyOptional()
|
||||
permissions?: string[];
|
||||
@ApiPropertyOptional()
|
||||
authProvider?: string;
|
||||
}
|
||||
|
||||
export class CreateUserReq {
|
||||
@@ -109,6 +117,12 @@ export class UpdateUserReq {
|
||||
@ApiPropertyOptional()
|
||||
hierarchy?: string;
|
||||
@ApiPropertyOptional()
|
||||
bio?: string;
|
||||
@ApiPropertyOptional()
|
||||
personalSite?: string;
|
||||
@ApiPropertyOptional()
|
||||
companyName?: string;
|
||||
@ApiPropertyOptional()
|
||||
roleNames?: string[];
|
||||
}
|
||||
|
||||
|
||||
@@ -266,7 +266,6 @@ export class UsersController {
|
||||
}
|
||||
|
||||
@Patch(':id')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.USERS.permissions.ADMIN)
|
||||
@ApiOkResponse({ type: UpdateUserRes })
|
||||
async updateUser(
|
||||
@User() user: RequestUser,
|
||||
@@ -274,6 +273,18 @@ 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,6 +125,7 @@ export class UsersService implements OnModuleInit {
|
||||
return { permissions };
|
||||
});
|
||||
res.user.permissions = permissions;
|
||||
res.user.authProvider = process.env.AUTH_PROVIDER || 'cognito';
|
||||
return res;
|
||||
}
|
||||
|
||||
@@ -148,20 +149,23 @@ export class UsersService implements OnModuleInit {
|
||||
}
|
||||
|
||||
async updateUser(req: UpdateUserReq, id: string, customerId: string) {
|
||||
const { department, hierarchy, jobTitle, name, roleNames, email } = req;
|
||||
if (roleNames) {
|
||||
const { roleNames, ...updateUserDTO } = req;
|
||||
if (roleNames && roleNames.length > 0) {
|
||||
await this.setRoles({ roleNames, userId: id }, customerId);
|
||||
}
|
||||
|
||||
const { user } = await lastValueFrom(
|
||||
this.usersClientService.UserUpdate({
|
||||
name,
|
||||
department: updateUserDTO.department,
|
||||
email: updateUserDTO.email,
|
||||
hierarchy: updateUserDTO.hierarchy,
|
||||
jobTitle: updateUserDTO.jobTitle,
|
||||
name: updateUserDTO.name,
|
||||
bio: updateUserDTO.bio,
|
||||
companyName: updateUserDTO.companyName,
|
||||
personalSite: updateUserDTO.personalSite,
|
||||
customerId,
|
||||
id,
|
||||
department,
|
||||
hierarchy,
|
||||
jobTitle,
|
||||
email,
|
||||
metabaseUserId: undefined,
|
||||
}),
|
||||
);
|
||||
|
||||
@@ -14,7 +14,10 @@ export class ValidationPipe implements PipeTransform<any> {
|
||||
return value;
|
||||
}
|
||||
const object = plainToInstance(metatype, value);
|
||||
const errors = await validate(object);
|
||||
const errors = await validate(object, {
|
||||
forbidUnknownValues: false,
|
||||
whitelist: true,
|
||||
});
|
||||
if (errors.length > 0) {
|
||||
const errorMessages = errors.map((err) => err.constraints);
|
||||
throw new BadRequestException(errorMessages);
|
||||
|
||||
@@ -0,0 +1,4 @@
|
||||
export const DYNAMODB_CONFIG = {
|
||||
region: () => process.env.AWS_REGION || 'us-east-1',
|
||||
inputsTable: () => process.env.INPUTS_DB || 'dadosfera-inputs-prd',
|
||||
};
|
||||
@@ -0,0 +1,9 @@
|
||||
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 {}
|
||||
@@ -0,0 +1,238 @@
|
||||
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;
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
@@ -0,0 +1,3 @@
|
||||
export * from './dynamodb.service';
|
||||
export * from './dynamodb.module';
|
||||
export * from './dynamodb.config';
|
||||
@@ -0,0 +1,5 @@
|
||||
export const ELASTICSEARCH_CONFIG = {
|
||||
getUrl: () => process.env.ELASTICSEARCH_URL || 'http://localhost:9200',
|
||||
getApiKey: () => process.env.ELASTICSEARCH_API_KEY || '',
|
||||
timeout: 10000, // 10 seconds
|
||||
};
|
||||
@@ -0,0 +1,9 @@
|
||||
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 {}
|
||||
@@ -0,0 +1,457 @@
|
||||
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,
|
||||
});
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,3 @@
|
||||
export * from './elasticsearch.module';
|
||||
export * from './elasticsearch.service';
|
||||
export * from './elasticsearch.config';
|
||||
@@ -0,0 +1,9 @@
|
||||
import { Module } from "@nestjs/common";
|
||||
import { NimbusService } from "./nimbus.service";
|
||||
import DadosferaLogger from "@dadosfera/dadosfera-logs";
|
||||
|
||||
@Module({
|
||||
providers: [NimbusService, DadosferaLogger],
|
||||
exports: [NimbusService],
|
||||
})
|
||||
export class NimbusServicesModule {}
|
||||
@@ -0,0 +1,56 @@
|
||||
import DadosferaLogger from "@dadosfera/dadosfera-logs";
|
||||
import { Inject, Injectable } from "@nestjs/common";
|
||||
import axios from "axios";
|
||||
|
||||
type TableUpdate = {
|
||||
table_schema: string;
|
||||
table_name: string;
|
||||
}
|
||||
|
||||
@Injectable()
|
||||
export class NimbusService {
|
||||
private logger: DadosferaLogger;
|
||||
|
||||
constructor(
|
||||
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
}
|
||||
|
||||
private buildUrl(customerName: string) {
|
||||
if (process.env.ENV === 'prd') {
|
||||
return `https://nimbus-${customerName}.dadosfera.ai`;
|
||||
}
|
||||
|
||||
return `https://nimbus-${customerName}.${process.env.ENV.replace(
|
||||
'local',
|
||||
'stg',
|
||||
)}.dadosfera.ai`;
|
||||
}
|
||||
|
||||
async renameTable(customerName: string, database: string, old: TableUpdate, update: TableUpdate) {
|
||||
const nimbusUrl = this.buildUrl(customerName);
|
||||
|
||||
const path = `/api/catalog/rename-tables/?database_name=${encodeURIComponent(database)}&table_name=${encodeURIComponent(old.table_name)}&table_schema=${encodeURIComponent(old.table_schema)}`;
|
||||
|
||||
try {
|
||||
this.logger.info("Request for PATCH " + nimbusUrl + path);
|
||||
this.logger.info("Payload: " + JSON.stringify(update));
|
||||
const { data } = await axios.patch(nimbusUrl + path, {
|
||||
table_name: update.table_name,
|
||||
table_schema: update.table_schema
|
||||
})
|
||||
|
||||
return data;
|
||||
} catch (error) {
|
||||
this.logger.error(error);
|
||||
return {
|
||||
message: error.message,
|
||||
database,
|
||||
old,
|
||||
update
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
@@ -5,15 +5,25 @@ import { redisStore } from 'cache-manager-ioredis-yet';
|
||||
@Module({
|
||||
imports: [
|
||||
CacheModule.registerAsync({
|
||||
useFactory: async () => ({
|
||||
store: await redisStore({
|
||||
ttl: 1000 * 60, //1 minute
|
||||
useFactory: async () => {
|
||||
const baseRedisConfig = {
|
||||
ttl: 5 * 1000 * 60, // 5 minute
|
||||
host: process.env.REDIS_HOST,
|
||||
port: process.env.REDIS_PORT && Number(process.env.REDIS_PORT),
|
||||
db: process.env.REDIS_DATABASE && Number(process.env.REDIS_DATABASE),
|
||||
keyPrefix: 'maestro:sso',
|
||||
}),
|
||||
}),
|
||||
}
|
||||
|
||||
if (process.env.REDIS_TLS === 'true') {
|
||||
baseRedisConfig['tls'] = {
|
||||
servername: process.env.REDIS_HOST,
|
||||
}
|
||||
}
|
||||
|
||||
return {
|
||||
store: await redisStore(baseRedisConfig),
|
||||
}
|
||||
},
|
||||
}),
|
||||
],
|
||||
providers: [CacheService],
|
||||
|
||||
@@ -0,0 +1,90 @@
|
||||
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