mirror of
https://github.com/dadosfera/maestro.git
synced 2026-09-01 20:28:17 +00:00
Compare commits
173
Commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
e06e7498db | ||
|
|
922a34a393 | ||
|
|
a795f543f6 | ||
|
|
15b5c49fb2 | ||
|
|
1730ec2753 | ||
|
|
dc4609cbc5 | ||
|
|
e1b0e88bd8 | ||
|
|
cef1184908 | ||
|
|
31dda867d1 | ||
|
|
b89909ad66 | ||
|
|
2bb280e8de | ||
|
|
00cbadbb45 | ||
|
|
47ad527d38 | ||
|
|
51044a23b3 | ||
|
|
bb29d126c1 | ||
|
|
9d0f449eeb | ||
|
|
0eafa67e6f | ||
|
|
c3937472ec | ||
|
|
84d64424ca | ||
|
|
f0bfc5c94b | ||
|
|
ad86a6a698 | ||
|
|
af3b11ad54 | ||
|
|
d54f381998 | ||
|
|
3fd586753e | ||
|
|
9b3894f9e4 | ||
|
|
4a5f8f679a | ||
|
|
a444f5e5ec | ||
|
|
0d30c1cf83 | ||
|
|
15860519e1 | ||
|
|
cffda86eab | ||
|
|
a93fbfbbd9 | ||
|
|
e1cfc6e1a8 | ||
|
|
1e63df6536 | ||
|
|
8eb7fd0169 | ||
|
|
f8c6a8b747 | ||
|
|
8a9d6c2f9c | ||
|
|
bcb2a0b7cb | ||
|
|
2e181e70af | ||
|
|
ef615adb8e | ||
|
|
e03900811b | ||
|
|
5bc5fb0977 | ||
|
|
011032e3d4 | ||
|
|
17363e74f4 | ||
|
|
63efff6adf | ||
|
|
b5f569e522 | ||
|
|
5c77577992 | ||
|
|
4b9e113185 | ||
|
|
06d505c50a | ||
|
|
c70abfc826 | ||
|
|
db2c7d6c02 | ||
|
|
0a99ce1aa4 | ||
|
|
39c66f030a | ||
|
|
1223e21ac4 | ||
|
|
2177f6725c | ||
|
|
65ab16236f | ||
|
|
b4cc8151d7 | ||
|
|
e8b982998f | ||
|
|
ccd4159c59 | ||
|
|
9e49abb40d | ||
|
|
21a82b64f3 | ||
|
|
2414fcf21e | ||
|
|
94fbdb2226 | ||
|
|
985170d7ae | ||
|
|
b38b8f51e2 | ||
|
|
ea16e62d6b | ||
|
|
6182705410 | ||
|
|
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 | ||
|
|
3bcbba9581 | ||
|
|
6f9c967c96 | ||
|
|
b8bdc5beea | ||
|
|
269f70b309 | ||
|
|
93f452ae05 | ||
|
|
e39378229f | ||
|
|
61a4f724ef | ||
|
|
38a9e21f5f | ||
|
|
d25bfd147c | ||
|
|
6c57bac235 | ||
|
|
cf8eed35a3 | ||
|
|
55fc85c544 | ||
|
|
72ed637640 | ||
|
|
890364f597 | ||
|
|
30eec733b2 | ||
|
|
30a41ba144 | ||
|
|
2dc032e7e7 | ||
|
|
7f5981731f | ||
|
|
66309c7bbe | ||
|
|
15048eaf8a | ||
|
|
0e169a3cbc | ||
|
|
b5d933eaf3 | ||
|
|
6fa9bf861a | ||
|
|
e08734c97f | ||
|
|
54b75ce11b | ||
|
|
85234fe0dd |
@@ -71,6 +71,11 @@ jobs:
|
||||
sudo mv helmfile /usr/local/bin/
|
||||
helmfile --version
|
||||
|
||||
- name: Install Helm Diff plugin
|
||||
run: |
|
||||
helm plugin install https://github.com/databus23/helm-diff --version v3.9.3
|
||||
helm diff version
|
||||
|
||||
- name: Debug Helm env
|
||||
run: |
|
||||
helm env
|
||||
@@ -102,4 +107,5 @@ jobs:
|
||||
- name: Run Helmfile Diff
|
||||
env:
|
||||
ENV: ${{ needs.extract_environment.outputs.environment }}
|
||||
HELM_PLUGINS: /home/runner/.local/share/helm/plugins
|
||||
run: helmfile -f deploy/helmfiles/${ENV}.yaml diff
|
||||
|
||||
+3
-3
@@ -1,5 +1,5 @@
|
||||
FROM node:20-alpine AS base_image
|
||||
RUN npm install -g npm@latest
|
||||
RUN npm install -g npm@10.8.2
|
||||
|
||||
FROM base_image AS build_base
|
||||
WORKDIR /app
|
||||
@@ -22,7 +22,7 @@ ENV PUPPETEER_SKIP_CHROMIUM_DOWNLOAD=true \
|
||||
# run aws cli without mounting secret, because CI already has AWS credentials
|
||||
FROM build_base AS ci_image
|
||||
RUN aws codeartifact login --tool npm --namespace @dadosfera --repository dadosfera-npm --domain dadosfera --domain-owner 611330257153 --region us-east-1
|
||||
RUN npm ci
|
||||
RUN npm ci --ignore-scripts
|
||||
COPY . .
|
||||
|
||||
|
||||
@@ -37,7 +37,7 @@ FROM build_base AS dev
|
||||
RUN --mount=type=secret,id=aws,target=/root/.aws/credentials \
|
||||
aws codeartifact login --tool npm --namespace @dadosfera --repository dadosfera-npm --domain dadosfera --domain-owner 611330257153 --region us-east-1
|
||||
# flag --build-from-source is required to force-build sqlite3
|
||||
RUN npm ci
|
||||
RUN npm ci --ignore-scripts
|
||||
COPY . .
|
||||
ENTRYPOINT npm run start:dev
|
||||
|
||||
|
||||
+1
-1
@@ -22,7 +22,7 @@ ENV PUPPETEER_SKIP_CHROMIUM_DOWNLOAD=true \
|
||||
FROM build_base AS build
|
||||
RUN --mount=type=secret,id=aws,target=/root/.aws/credentials \
|
||||
aws codeartifact login --tool npm --namespace @dadosfera --repository dadosfera-npm --domain dadosfera --domain-owner 611330257153 --region us-east-1
|
||||
RUN npm ci
|
||||
RUN npm ci --ignore-scripts
|
||||
COPY . .
|
||||
RUN npm run build
|
||||
|
||||
|
||||
@@ -4,6 +4,7 @@
|
||||
|
||||
# Maestro
|
||||
|
||||
|
||||
Maestro é a API principal da Dadosfera. É responsável pela comunicação do Frontend com nossos microsserviços.
|
||||
|
||||
```mermaid
|
||||
|
||||
@@ -111,8 +111,12 @@ spec:
|
||||
value: "{{ .Values.maestro.redis_tls }}"
|
||||
- name: PLATFORM_API_URL
|
||||
value: {{ .Values.maestro.platform_api_url }}
|
||||
- name: CONNECTIONS_API_URL
|
||||
value: {{ .Values.maestro.connections_api_url | default "" | quote }}
|
||||
- name: STORAGE_EXPLORER_API_URL
|
||||
value: {{ .Values.maestro.storage_explorer_api_url | quote }}
|
||||
- name: FIREBASE_BASE_URL
|
||||
value: {{ .Values.maestro.firebase_base_url }}
|
||||
- name: JWT_PRIVATE_KEY
|
||||
valueFrom:
|
||||
secretKeyRef:
|
||||
|
||||
@@ -9,7 +9,9 @@ maestro:
|
||||
cookie_secret: "ff7bc13823edb2ae50d248e5780bddc9d4b31c36"
|
||||
redis_database: "1"
|
||||
platform_api_url: https://xs2hkhq07k.execute-api.us-east-1.amazonaws.com
|
||||
connections_api_url: https://iy40eans64.execute-api.us-east-1.amazonaws.com
|
||||
storage_explorer_api_url: "http://storage-explorer-{customer}.data-apps.svc.cluster.local:8000/api"
|
||||
firebase_base_url: https://feature-flag-25bf6-default-rtdb.firebaseio.com/stg
|
||||
|
||||
hostname: maestro.stg.dadosfera.ai
|
||||
|
||||
|
||||
@@ -55,6 +55,7 @@ maestro:
|
||||
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
|
||||
|
||||
@@ -0,0 +1,759 @@
|
||||
# Orchest Module Identity Implementation Plan
|
||||
|
||||
> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking.
|
||||
|
||||
**Goal:** Add a Maestro endpoint `GET /auth/module-identity` that translates a Dadosfera session into the `X-Auth-*` identity headers Orchest's RBAC trusts, gating on a module permission and (for service ingresses) authorizing per-project against orchest-api.
|
||||
|
||||
**Architecture:** One route serves two nginx `auth_request` callers, keyed on whether the ingress annotation carries authz query params. Phase 1: authenticate the `ddf-auth` cookie via Maestro's existing JWKS verify, gate on the module permission (seqid 31), emit identity headers, set `X-Auth-Roles: admin` for Super Admin (seqid 34). Phase 2: when the annotation carries `permission`+`project_uuid`, additionally call the calling tenant's orchest-api `/api/authz/check` (host derived from the JWT's `customer_name` + the `orchest-{module}-{customer_name}` namespace convention) and relay allow/deny, fail-closed. A small orchest-side change flips `service_access_auth_url` to build a scoped URL in Maestro mode.
|
||||
|
||||
**Tech Stack:** NestJS (controllers/providers, Jest via `Test.createTestingModule`), Express `Request`/`Response`, `jsonwebtoken`, plain `process.env` config. Orchest side: Python (`lib/python/orchest-internals`).
|
||||
|
||||
**Spec:** `dbt-to-orchest` repo — `docs/superpowers/specs/2026-08-20-maestro-module-identity-design.md` (the Maestro repo does not hold the spec; executors read it there).
|
||||
|
||||
## Global Constraints
|
||||
|
||||
- **Config via `process.env`** — Maestro reads env directly (no ConfigService). Pattern: `const X = process.env.X || '<default>'` (see `authentication.guard.ts:140`).
|
||||
- **Config values (verbatim):** `ORCHEST_MODULE_PERMISSION_SEQID` default `31` (Intelligence/"Orchest Module", `intelligence:open`); `ORCHEST_ADMIN_PERMISSION_SEQIDS` default `34` (Super Admin, `users:admin`), a comma-separated list parsed to numbers; `ORCHEST_NAMESPACE_MODULE` default `intelli` (∈ `intelli`|`process`).
|
||||
- **Token source is the `ddf-auth` COOKIE**, not the `Authorization` header (the nginx `auth_request` subrequest carries the browser cookie). Never read `Authorization` in this route.
|
||||
- **`permission`/`project_uuid` come from `request.query`** (the orchest-api-authored annotation), never from the end-user URL. The end-user URL is `X-Original-URI` and MUST NOT be read for authz.
|
||||
- **Fail-closed:** any error reaching orchest-api, or an unexpected status, → deny (never allow).
|
||||
- **`/auth/me` is not modified.** Other services depend on it.
|
||||
- **No orchest-api code change.** `/api/authz/check`, the `AuthCustom/AuthUrl` ingress path, and RBAC already exist.
|
||||
- **Namespace pattern:** `orchest-{module}-{customer_name}`. `dadosferademo2` is the one exception (being removed) — explicitly unsupported, no special-casing.
|
||||
- Branch: `feat/orchest-module-identity` (off `origin/beta`), already created.
|
||||
|
||||
---
|
||||
|
||||
## File Structure
|
||||
|
||||
**Maestro (`feat/orchest-module-identity` off beta):**
|
||||
- Modify `src/modules/auth/auth.service.ts` — make `validateJwtToken` public (or add a public `verifyAccessToken` wrapper); add `authorizeOrchestServiceAccess(...)` helper (Phase 2).
|
||||
- Modify `src/modules/auth/auth.controller.ts` — add the `GET /auth/module-identity` route.
|
||||
- Create `src/modules/auth/orchest-identity.ts` — pure, testable helpers: `parseAdminSeqids(env)`, `isModuleAllowed(perms, gateSeqid)`, `isAdmin(perms, adminSeqids)`, `tenantOrchestApiHost(customerName, module)`, `isAllowedOrchestApiHost(host)`. Keeps set-membership/string logic out of the controller so it unit-tests without HTTP.
|
||||
- Create `src/modules/auth/orchest-identity.spec.ts` — unit tests for the helpers.
|
||||
- Create `src/modules/auth/module-identity.controller.spec.ts` — controller tests (mock `AuthClientService`, fake `Request`/`Response`).
|
||||
|
||||
**dbt-to-orchest repo (Phase 2 orchest-side, separate branch there):**
|
||||
- Modify `lib/python/orchest-internals/_orchest/internals/utils.py:19-51` — `service_access_auth_url` Maestro branch.
|
||||
- Modify `lib/python/orchest-internals/tests/…` (or wherever `utils` is tested) — add the Maestro-mode case.
|
||||
|
||||
---
|
||||
|
||||
## PHASE 1 — Identity (webserver ingress)
|
||||
|
||||
### Task 1: Pure helpers for permission mapping
|
||||
|
||||
**Files:**
|
||||
- Create: `src/modules/auth/orchest-identity.ts`
|
||||
- Test: `src/modules/auth/orchest-identity.spec.ts`
|
||||
|
||||
**Interfaces:**
|
||||
- Consumes: nothing.
|
||||
- Produces:
|
||||
- `parseAdminSeqids(raw: string | undefined): number[]` — parse `"34"` / `"34,40"` → `[34]` / `[34,40]`; empty/undefined → `[34]`.
|
||||
- `moduleGateSeqid(raw: string | undefined): number` — parse `ORCHEST_MODULE_PERMISSION_SEQID` → number; default `31`.
|
||||
- `isModuleAllowed(perms: number[], gateSeqid: number): boolean`
|
||||
- `isAdmin(perms: number[], adminSeqids: number[]): boolean`
|
||||
|
||||
- [ ] **Step 1: Write the failing test**
|
||||
|
||||
```typescript
|
||||
import {
|
||||
parseAdminSeqids,
|
||||
moduleGateSeqid,
|
||||
isModuleAllowed,
|
||||
isAdmin,
|
||||
} from './orchest-identity';
|
||||
|
||||
describe('orchest-identity mapping', () => {
|
||||
it('parseAdminSeqids: default, single, list, whitespace', () => {
|
||||
expect(parseAdminSeqids(undefined)).toEqual([34]);
|
||||
expect(parseAdminSeqids('')).toEqual([34]);
|
||||
expect(parseAdminSeqids('34')).toEqual([34]);
|
||||
expect(parseAdminSeqids('34,40')).toEqual([34, 40]);
|
||||
expect(parseAdminSeqids(' 34 , 40 ')).toEqual([34, 40]);
|
||||
});
|
||||
|
||||
it('moduleGateSeqid: default and override', () => {
|
||||
expect(moduleGateSeqid(undefined)).toBe(31);
|
||||
expect(moduleGateSeqid('43')).toBe(43);
|
||||
});
|
||||
|
||||
it('isModuleAllowed', () => {
|
||||
expect(isModuleAllowed([31, 5], 31)).toBe(true);
|
||||
expect(isModuleAllowed([5, 7], 31)).toBe(false);
|
||||
expect(isModuleAllowed([], 31)).toBe(false);
|
||||
});
|
||||
|
||||
it('isAdmin: intersection', () => {
|
||||
expect(isAdmin([31, 34], [34])).toBe(true);
|
||||
expect(isAdmin([31], [34])).toBe(false);
|
||||
expect(isAdmin([99], [34, 99])).toBe(true);
|
||||
});
|
||||
});
|
||||
```
|
||||
|
||||
- [ ] **Step 2: Run test to verify it fails**
|
||||
|
||||
Run: `npx jest src/modules/auth/orchest-identity.spec.ts -t 'orchest-identity mapping'`
|
||||
Expected: FAIL — `Cannot find module './orchest-identity'`.
|
||||
|
||||
- [ ] **Step 3: Write minimal implementation**
|
||||
|
||||
```typescript
|
||||
// src/modules/auth/orchest-identity.ts
|
||||
export function parseAdminSeqids(raw: string | undefined): number[] {
|
||||
if (!raw || !raw.trim()) return [34];
|
||||
return raw
|
||||
.split(',')
|
||||
.map((s) => Number(s.trim()))
|
||||
.filter((n) => Number.isInteger(n));
|
||||
}
|
||||
|
||||
export function moduleGateSeqid(raw: string | undefined): number {
|
||||
const n = Number(raw);
|
||||
return Number.isInteger(n) && n > 0 ? n : 31;
|
||||
}
|
||||
|
||||
export function isModuleAllowed(perms: number[], gateSeqid: number): boolean {
|
||||
return Array.isArray(perms) && perms.includes(gateSeqid);
|
||||
}
|
||||
|
||||
export function isAdmin(perms: number[], adminSeqids: number[]): boolean {
|
||||
return (
|
||||
Array.isArray(perms) && perms.some((p) => adminSeqids.includes(p))
|
||||
);
|
||||
}
|
||||
```
|
||||
|
||||
- [ ] **Step 4: Run test to verify it passes**
|
||||
|
||||
Run: `npx jest src/modules/auth/orchest-identity.spec.ts -t 'orchest-identity mapping'`
|
||||
Expected: PASS.
|
||||
|
||||
- [ ] **Step 5: Commit**
|
||||
|
||||
```bash
|
||||
git add src/modules/auth/orchest-identity.ts src/modules/auth/orchest-identity.spec.ts
|
||||
git commit -m "feat(orchest-identity): permission-mapping helpers"
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
### Task 2: Expose JWKS verify on AuthClientService
|
||||
|
||||
**Files:**
|
||||
- Modify: `src/modules/auth/auth.service.ts:413-425`
|
||||
|
||||
**Interfaces:**
|
||||
- Consumes: existing `getPublicKeys()`.
|
||||
- Produces: `public async verifyAccessToken(token: string): Promise<any>` — verifies the JWT against JWKS and returns its payload; throws on missing/invalid token or unknown `kid`. (Rename of the existing private `validateJwtToken`, kept callable by `validateUserSession`.)
|
||||
|
||||
- [ ] **Step 1: Make the method public and rename**
|
||||
|
||||
The method already does exactly the needed decode+verify. Rename `validateJwtToken` → `verifyAccessToken`, change `private` → `public`, and update its one caller.
|
||||
|
||||
In `src/modules/auth/auth.service.ts`, change line 413:
|
||||
|
||||
```typescript
|
||||
public async verifyAccessToken(token: string) {
|
||||
const decoded: any = token && jwt.decode(token, { complete: true });
|
||||
if (!decoded) throw new Error('Invalid token');
|
||||
|
||||
const { kid } = decoded.header;
|
||||
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;
|
||||
}
|
||||
```
|
||||
|
||||
And update the caller in `validateUserSession` (was line 310):
|
||||
|
||||
```typescript
|
||||
const payload = await this.verifyAccessToken(accessToken);
|
||||
```
|
||||
|
||||
- [ ] **Step 2: Verify existing suite still compiles/passes for auth.service**
|
||||
|
||||
Run: `npx jest src/modules/auth`
|
||||
Expected: PASS (no behavior change; `/auth/me` path unaffected). If there is no existing auth.service spec, run `npx tsc --noEmit` to confirm the rename compiles.
|
||||
|
||||
- [ ] **Step 3: Commit**
|
||||
|
||||
```bash
|
||||
git add src/modules/auth/auth.service.ts
|
||||
git commit -m "refactor(auth): expose verifyAccessToken (was private validateJwtToken)"
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
### Task 3: `GET /auth/module-identity` — identity + module gate
|
||||
|
||||
**Files:**
|
||||
- Modify: `src/modules/auth/auth.controller.ts` (add route beside `getMe`, ~after line 533)
|
||||
- Test: `src/modules/auth/module-identity.controller.spec.ts`
|
||||
|
||||
**Interfaces:**
|
||||
- Consumes: `AuthClientService.verifyAccessToken` (Task 2); `parseAdminSeqids`/`moduleGateSeqid`/`isModuleAllowed`/`isAdmin` (Task 1).
|
||||
- Produces: route `GET /auth/module-identity`. On 200 sets response headers `X-Auth-User`, `X-Auth-Username`, and (admin only) `X-Auth-Roles: admin`; empty body. 401 (no/invalid cookie), 403 (lacks module gate).
|
||||
|
||||
- [ ] **Step 1: Write the failing test**
|
||||
|
||||
```typescript
|
||||
import { Test } from '@nestjs/testing';
|
||||
import { AuthController } from './auth.controller';
|
||||
import { AuthClientService } from './auth.service';
|
||||
|
||||
function res() {
|
||||
const headers: Record<string, string> = {};
|
||||
const r: any = {
|
||||
_status: 0,
|
||||
_sent: undefined,
|
||||
set: (k: string, v: string) => { headers[k] = v; return r; },
|
||||
status: (c: number) => { r._status = c; return r; },
|
||||
send: (b?: any) => { r._sent = b ?? ''; return r; },
|
||||
json: (b?: any) => { r._sent = b; return r; },
|
||||
_headers: headers,
|
||||
};
|
||||
return r;
|
||||
}
|
||||
function req(cookie?: string, query: Record<string, string> = {}) {
|
||||
return { cookies: cookie ? { 'ddf-auth': cookie } : {}, query } as any;
|
||||
}
|
||||
|
||||
describe('GET /auth/module-identity — identity', () => {
|
||||
let controller: AuthController;
|
||||
const auth = { verifyAccessToken: jest.fn() } as unknown as AuthClientService;
|
||||
|
||||
beforeEach(async () => {
|
||||
jest.resetAllMocks();
|
||||
process.env.ORCHEST_MODULE_PERMISSION_SEQID = '31';
|
||||
process.env.ORCHEST_ADMIN_PERMISSION_SEQIDS = '34';
|
||||
const mod = await Test.createTestingModule({
|
||||
controllers: [AuthController],
|
||||
providers: [{ provide: AuthClientService, useValue: auth }],
|
||||
})
|
||||
// Any other providers AuthController injects must be stubbed here the
|
||||
// same way (ApiKeyService, DadosferaLogger, etc.). Add them as the
|
||||
// compile step reports missing providers.
|
||||
.compile();
|
||||
controller = mod.get(AuthController);
|
||||
});
|
||||
|
||||
it('no cookie → 401', async () => {
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req(undefined), r);
|
||||
expect(r._status).toBe(401);
|
||||
});
|
||||
|
||||
it('valid + module + admin → 200 with X-Auth-Roles: admin', async () => {
|
||||
(auth.verifyAccessToken as jest.Mock).mockResolvedValue({
|
||||
user_id: 'u-1', username: 'alice', permissions: [31, 34],
|
||||
});
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok'), r);
|
||||
expect(r._status).toBe(200);
|
||||
expect(r._headers['X-Auth-User']).toBe('u-1');
|
||||
expect(r._headers['X-Auth-Username']).toBe('alice');
|
||||
expect(r._headers['X-Auth-Roles']).toBe('admin');
|
||||
});
|
||||
|
||||
it('valid + module, not admin → 200, no X-Auth-Roles', async () => {
|
||||
(auth.verifyAccessToken as jest.Mock).mockResolvedValue({
|
||||
user_id: 'u-2', username: 'bob', permissions: [31],
|
||||
});
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok'), r);
|
||||
expect(r._status).toBe(200);
|
||||
expect(r._headers['X-Auth-Roles']).toBeUndefined();
|
||||
});
|
||||
|
||||
it('valid, lacks module → 403', async () => {
|
||||
(auth.verifyAccessToken as jest.Mock).mockResolvedValue({
|
||||
user_id: 'u-3', username: 'carol', permissions: [5],
|
||||
});
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok'), r);
|
||||
expect(r._status).toBe(403);
|
||||
});
|
||||
|
||||
it('verify throws (expired/bad) → 401', async () => {
|
||||
(auth.verifyAccessToken as jest.Mock).mockRejectedValue(new Error('bad'));
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok'), r);
|
||||
expect(r._status).toBe(401);
|
||||
});
|
||||
|
||||
it('admin seqids extended by config → 200 admin', async () => {
|
||||
process.env.ORCHEST_ADMIN_PERMISSION_SEQIDS = '34,99';
|
||||
(auth.verifyAccessToken as jest.Mock).mockResolvedValue({
|
||||
user_id: 'u-4', username: 'dana', permissions: [31, 99],
|
||||
});
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok'), r);
|
||||
expect(r._headers['X-Auth-Roles']).toBe('admin');
|
||||
});
|
||||
});
|
||||
```
|
||||
|
||||
- [ ] **Step 2: Run test to verify it fails**
|
||||
|
||||
Run: `npx jest src/modules/auth/module-identity.controller.spec.ts`
|
||||
Expected: FAIL — `controller.moduleIdentity is not a function` (and possibly missing-provider errors, which tell you which providers to stub — add them to the `providers` array per the comment).
|
||||
|
||||
- [ ] **Step 3: Write minimal implementation**
|
||||
|
||||
Add to `auth.controller.ts` (import the helpers at top; `AuthClientService` is already injected as `this.authClient`):
|
||||
|
||||
```typescript
|
||||
import {
|
||||
parseAdminSeqids,
|
||||
moduleGateSeqid,
|
||||
isModuleAllowed,
|
||||
isAdmin,
|
||||
} from './orchest-identity';
|
||||
```
|
||||
|
||||
```typescript
|
||||
@Get('module-identity')
|
||||
async moduleIdentity(@Req() req: Request, @Res() res: Response) {
|
||||
const token = req.cookies?.['ddf-auth'];
|
||||
if (!token) {
|
||||
return res.status(401).send();
|
||||
}
|
||||
|
||||
let payload: any;
|
||||
try {
|
||||
payload = await this.authClient.verifyAccessToken(token);
|
||||
} catch (e) {
|
||||
return res.status(401).send();
|
||||
}
|
||||
|
||||
const perms: number[] = payload?.permissions ?? [];
|
||||
const gate = moduleGateSeqid(process.env.ORCHEST_MODULE_PERMISSION_SEQID);
|
||||
if (!isModuleAllowed(perms, gate)) {
|
||||
return res.status(403).send();
|
||||
}
|
||||
|
||||
res.set('X-Auth-User', String(payload.user_id));
|
||||
res.set('X-Auth-Username', String(payload.username ?? ''));
|
||||
if (isAdmin(perms, parseAdminSeqids(process.env.ORCHEST_ADMIN_PERMISSION_SEQIDS))) {
|
||||
res.set('X-Auth-Roles', 'admin');
|
||||
}
|
||||
|
||||
// Phase 2 authz branch is inserted here (Task 5) before the 200.
|
||||
return res.status(200).send();
|
||||
}
|
||||
```
|
||||
|
||||
- [ ] **Step 4: Run test to verify it passes**
|
||||
|
||||
Run: `npx jest src/modules/auth/module-identity.controller.spec.ts`
|
||||
Expected: PASS (all identity cases).
|
||||
|
||||
- [ ] **Step 5: Commit**
|
||||
|
||||
```bash
|
||||
git add src/modules/auth/auth.controller.ts src/modules/auth/module-identity.controller.spec.ts
|
||||
git commit -m "feat(auth): GET /auth/module-identity — identity + module gate"
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## PHASE 2 — Per-service authz
|
||||
|
||||
### Task 4: Tenant orchest-api host derivation + allowlist
|
||||
|
||||
**Files:**
|
||||
- Modify: `src/modules/auth/orchest-identity.ts`
|
||||
- Modify: `src/modules/auth/orchest-identity.spec.ts`
|
||||
|
||||
**Interfaces:**
|
||||
- Consumes: nothing.
|
||||
- Produces:
|
||||
- `tenantOrchestApiHost(customerName: string, module: string): string` — returns `orchest-api.orchest-{module}-{customerName}.svc.cluster.local`.
|
||||
- `isAllowedOrchestApiHost(host: string): boolean` — matches `^orchest-api\.orchest-(intelli|process)-[a-z0-9-]+\.svc\.cluster\.local$`.
|
||||
- `namespaceModule(raw: string | undefined): 'intelli' | 'process'` — parse `ORCHEST_NAMESPACE_MODULE`; default `intelli`; anything not `process` → `intelli`.
|
||||
|
||||
- [ ] **Step 1: Write the failing test (append to orchest-identity.spec.ts)**
|
||||
|
||||
```typescript
|
||||
import {
|
||||
tenantOrchestApiHost,
|
||||
isAllowedOrchestApiHost,
|
||||
namespaceModule,
|
||||
} from './orchest-identity';
|
||||
|
||||
describe('orchest-identity tenant routing', () => {
|
||||
it('namespaceModule default and values', () => {
|
||||
expect(namespaceModule(undefined)).toBe('intelli');
|
||||
expect(namespaceModule('process')).toBe('process');
|
||||
expect(namespaceModule('garbage')).toBe('intelli');
|
||||
});
|
||||
|
||||
it('tenantOrchestApiHost builds the namespace pattern', () => {
|
||||
expect(tenantOrchestApiHost('acme', 'intelli')).toBe(
|
||||
'orchest-api.orchest-intelli-acme.svc.cluster.local',
|
||||
);
|
||||
expect(tenantOrchestApiHost('acme', 'process')).toBe(
|
||||
'orchest-api.orchest-process-acme.svc.cluster.local',
|
||||
);
|
||||
});
|
||||
|
||||
it('isAllowedOrchestApiHost guards against malformed values', () => {
|
||||
expect(
|
||||
isAllowedOrchestApiHost('orchest-api.orchest-intelli-acme.svc.cluster.local'),
|
||||
).toBe(true);
|
||||
expect(isAllowedOrchestApiHost('evil.example.com')).toBe(false);
|
||||
expect(
|
||||
isAllowedOrchestApiHost('orchest-api.orchest-intelli-.svc.cluster.local'),
|
||||
).toBe(false);
|
||||
expect(
|
||||
isAllowedOrchestApiHost('orchest-api.orchest-other-acme.svc.cluster.local'),
|
||||
).toBe(false);
|
||||
});
|
||||
});
|
||||
```
|
||||
|
||||
- [ ] **Step 2: Run test to verify it fails**
|
||||
|
||||
Run: `npx jest src/modules/auth/orchest-identity.spec.ts -t 'tenant routing'`
|
||||
Expected: FAIL — the three functions are not exported.
|
||||
|
||||
- [ ] **Step 3: Write minimal implementation (append to orchest-identity.ts)**
|
||||
|
||||
```typescript
|
||||
export function namespaceModule(raw: string | undefined): 'intelli' | 'process' {
|
||||
return raw === 'process' ? 'process' : 'intelli';
|
||||
}
|
||||
|
||||
export function tenantOrchestApiHost(
|
||||
customerName: string,
|
||||
module: string,
|
||||
): string {
|
||||
return `orchest-api.orchest-${module}-${customerName}.svc.cluster.local`;
|
||||
}
|
||||
|
||||
const ORCHEST_API_HOST_RE =
|
||||
/^orchest-api\.orchest-(intelli|process)-[a-z0-9-]+\.svc\.cluster\.local$/;
|
||||
|
||||
export function isAllowedOrchestApiHost(host: string): boolean {
|
||||
return ORCHEST_API_HOST_RE.test(host);
|
||||
}
|
||||
```
|
||||
|
||||
**Slug-normalization note (verify during implementation):** confirm the
|
||||
JWT's `customer_name` is *exactly* the namespace slug (lowercase, kebab, no
|
||||
spaces). If it is not, normalize deterministically inside
|
||||
`tenantOrchestApiHost` (e.g. `customerName.toLowerCase().replace(/[^a-z0-9-]/g, '-')`)
|
||||
and extend the test with the raw→normalized case. Do NOT guess the rule —
|
||||
inspect a real token or ask the team.
|
||||
|
||||
- [ ] **Step 4: Run test to verify it passes**
|
||||
|
||||
Run: `npx jest src/modules/auth/orchest-identity.spec.ts -t 'tenant routing'`
|
||||
Expected: PASS.
|
||||
|
||||
- [ ] **Step 5: Commit**
|
||||
|
||||
```bash
|
||||
git add src/modules/auth/orchest-identity.ts src/modules/auth/orchest-identity.spec.ts
|
||||
git commit -m "feat(orchest-identity): tenant orchest-api host derivation + allowlist"
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
### Task 5: authz branch — call `/api/authz/check`, relay, fail-closed
|
||||
|
||||
**Files:**
|
||||
- Modify: `src/modules/auth/auth.service.ts` (add `authorizeOrchestServiceAccess`)
|
||||
- Modify: `src/modules/auth/auth.controller.ts` (insert the authz branch in `moduleIdentity`)
|
||||
- Modify: `src/modules/auth/module-identity.controller.spec.ts`
|
||||
|
||||
**Interfaces:**
|
||||
- Consumes: `tenantOrchestApiHost`, `isAllowedOrchestApiHost`, `namespaceModule` (Task 4); an HTTP client. Maestro uses gRPC for its own services but plain HTTP for this cross-service call — use `axios` if already a dependency, else the `@nestjs/axios` `HttpService`; confirm which is present before writing (grep `import axios` / `HttpService`).
|
||||
- Produces: `AuthClientService.authorizeOrchestServiceAccess(args: { host: string; permission: string; projectUuid?: string; headers: Record<string,string> }): Promise<'allow' | 'deny' | 'error'>` — GET `http://{host}/api/authz/check?permission=…[&project_uuid=…]` with the identity headers; 200→`allow`, 403→`deny`, anything else/throw→`error`.
|
||||
|
||||
- [ ] **Step 1: Write the failing test (append to module-identity.controller.spec.ts)**
|
||||
|
||||
```typescript
|
||||
describe('GET /auth/module-identity — per-service authz', () => {
|
||||
let controller: AuthController;
|
||||
const auth = {
|
||||
verifyAccessToken: jest.fn(),
|
||||
authorizeOrchestServiceAccess: jest.fn(),
|
||||
} as unknown as AuthClientService;
|
||||
|
||||
beforeEach(async () => {
|
||||
jest.resetAllMocks();
|
||||
process.env.ORCHEST_MODULE_PERMISSION_SEQID = '31';
|
||||
process.env.ORCHEST_ADMIN_PERMISSION_SEQIDS = '34';
|
||||
process.env.ORCHEST_NAMESPACE_MODULE = 'intelli';
|
||||
const mod = await Test.createTestingModule({
|
||||
controllers: [AuthController],
|
||||
providers: [{ provide: AuthClientService, useValue: auth }],
|
||||
}).compile(); // add the same stubbed providers as Task 3
|
||||
controller = mod.get(AuthController);
|
||||
(auth.verifyAccessToken as jest.Mock).mockResolvedValue({
|
||||
user_id: 'u-1', username: 'alice', permissions: [31], customer_name: 'acme',
|
||||
});
|
||||
});
|
||||
|
||||
const q = { permission: 'session.open', project_uuid: 'p-1' };
|
||||
|
||||
it('has grant → 200 and calls the tenant host', async () => {
|
||||
(auth.authorizeOrchestServiceAccess as jest.Mock).mockResolvedValue('allow');
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok', q), r);
|
||||
expect(r._status).toBe(200);
|
||||
expect(auth.authorizeOrchestServiceAccess).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
host: 'orchest-api.orchest-intelli-acme.svc.cluster.local',
|
||||
permission: 'session.open',
|
||||
projectUuid: 'p-1',
|
||||
}),
|
||||
);
|
||||
});
|
||||
|
||||
it('lacks grant → 403', async () => {
|
||||
(auth.authorizeOrchestServiceAccess as jest.Mock).mockResolvedValue('deny');
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok', q), r);
|
||||
expect(r._status).toBe(403);
|
||||
});
|
||||
|
||||
it('orchest-api error → 502 (fail-closed)', async () => {
|
||||
(auth.authorizeOrchestServiceAccess as jest.Mock).mockResolvedValue('error');
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok', q), r);
|
||||
expect(r._status).toBe(502);
|
||||
});
|
||||
|
||||
it('permission without project_uuid → 403 (all-or-nothing)', async () => {
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok', { permission: 'session.open' }), r);
|
||||
expect(r._status).toBe(403);
|
||||
expect(auth.authorizeOrchestServiceAccess).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('no authz params → 200 identity-only (webserver case)', async () => {
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok', {}), r);
|
||||
expect(r._status).toBe(200);
|
||||
expect(auth.authorizeOrchestServiceAccess).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('malformed customer_name → host fails allowlist → 403, no call', async () => {
|
||||
(auth.verifyAccessToken as jest.Mock).mockResolvedValue({
|
||||
user_id: 'u-1', username: 'alice', permissions: [31], customer_name: 'Bad Name!',
|
||||
});
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok', q), r);
|
||||
expect(r._status).toBe(403);
|
||||
expect(auth.authorizeOrchestServiceAccess).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
```
|
||||
|
||||
- [ ] **Step 2: Run test to verify it fails**
|
||||
|
||||
Run: `npx jest src/modules/auth/module-identity.controller.spec.ts -t 'per-service authz'`
|
||||
Expected: FAIL — `authorizeOrchestServiceAccess` undefined / authz branch absent.
|
||||
|
||||
- [ ] **Step 3a: Implement the service helper**
|
||||
|
||||
In `auth.service.ts` (use the HTTP client confirmed in the Interfaces note; `axios` shown):
|
||||
|
||||
```typescript
|
||||
public async authorizeOrchestServiceAccess(args: {
|
||||
host: string;
|
||||
permission: string;
|
||||
projectUuid?: string;
|
||||
headers: Record<string, string>;
|
||||
}): Promise<'allow' | 'deny' | 'error'> {
|
||||
const params: Record<string, string> = { permission: args.permission };
|
||||
if (args.projectUuid) params.project_uuid = args.projectUuid;
|
||||
try {
|
||||
const resp = await axios.get(`http://${args.host}/api/authz/check`, {
|
||||
params,
|
||||
headers: args.headers,
|
||||
timeout: 5000,
|
||||
validateStatus: () => true, // never throw on 4xx/5xx; we branch below
|
||||
});
|
||||
if (resp.status === 200) return 'allow';
|
||||
if (resp.status === 403) return 'deny';
|
||||
return 'error';
|
||||
} catch (e) {
|
||||
return 'error';
|
||||
}
|
||||
}
|
||||
```
|
||||
|
||||
- [ ] **Step 3b: Insert the authz branch in the controller**
|
||||
|
||||
Replace the `// Phase 2 authz branch is inserted here` marker (Task 3) with:
|
||||
|
||||
```typescript
|
||||
const permission = req.query?.permission as string | undefined;
|
||||
const projectUuid = req.query?.project_uuid as string | undefined;
|
||||
if (permission) {
|
||||
// Service-ingress caller: authorize per-project. All-or-nothing —
|
||||
// an incomplete annotation must not silently skip the check.
|
||||
if (!projectUuid) {
|
||||
return res.status(403).send();
|
||||
}
|
||||
const module = namespaceModule(process.env.ORCHEST_NAMESPACE_MODULE);
|
||||
const host = tenantOrchestApiHost(String(payload.customer_name ?? ''), module);
|
||||
if (!isAllowedOrchestApiHost(host)) {
|
||||
return res.status(403).send();
|
||||
}
|
||||
const identityHeaders: Record<string, string> = {
|
||||
'X-Auth-User': String(payload.user_id),
|
||||
'X-Auth-Username': String(payload.username ?? ''),
|
||||
};
|
||||
if (isAdmin(perms, parseAdminSeqids(process.env.ORCHEST_ADMIN_PERMISSION_SEQIDS))) {
|
||||
identityHeaders['X-Auth-Roles'] = 'admin';
|
||||
}
|
||||
const decision = await this.authClient.authorizeOrchestServiceAccess({
|
||||
host,
|
||||
permission,
|
||||
projectUuid,
|
||||
headers: identityHeaders,
|
||||
});
|
||||
if (decision === 'deny') return res.status(403).send();
|
||||
if (decision === 'error') return res.status(502).send();
|
||||
// 'allow' falls through to the 200 below (identity headers already set).
|
||||
}
|
||||
|
||||
return res.status(200).send();
|
||||
```
|
||||
|
||||
Add the imports `tenantOrchestApiHost`, `isAllowedOrchestApiHost`, `namespaceModule` to the existing `./orchest-identity` import line.
|
||||
|
||||
- [ ] **Step 4: Run test to verify it passes**
|
||||
|
||||
Run: `npx jest src/modules/auth/module-identity.controller.spec.ts`
|
||||
Expected: PASS (identity + per-service authz suites).
|
||||
|
||||
- [ ] **Step 5: Commit**
|
||||
|
||||
```bash
|
||||
git add src/modules/auth/auth.service.ts src/modules/auth/auth.controller.ts src/modules/auth/module-identity.controller.spec.ts
|
||||
git commit -m "feat(auth): per-service authz branch on /auth/module-identity"
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
### Task 6: orchest-side — scope the Maestro service auth-url (dbt-to-orchest repo)
|
||||
|
||||
**Files:**
|
||||
- Modify: `lib/python/orchest-internals/_orchest/internals/utils.py:19-51` (`service_access_auth_url`)
|
||||
- Test: the module's existing test (grep `service_access_auth_url` under `lib/python/…/tests`; if none, create `lib/python/orchest-internals/tests/test_service_access_auth_url.py`)
|
||||
|
||||
**This task is in the `dbt-to-orchest` repo, not Maestro.** Do it on a branch there (e.g. off the current RBAC branch). It has no dependency on Tasks 1-5 compiling, but the annotation it produces is what Task 5 consumes at runtime.
|
||||
|
||||
**Interfaces:**
|
||||
- Consumes: nothing new.
|
||||
- Produces: `service_access_auth_url(base_auth_url, permission, project_uuid, auth_custom)` — in Maestro mode (`auth_custom=True`), returns `f"{base_auth_url}?permission={permission}"` (plus `&project_uuid=…` when set) instead of returning `base_auth_url` unchanged.
|
||||
|
||||
- [ ] **Step 1: Write the failing test**
|
||||
|
||||
```python
|
||||
from _orchest.internals.utils import service_access_auth_url
|
||||
|
||||
def test_maestro_mode_appends_scoped_params():
|
||||
url = service_access_auth_url(
|
||||
"http://maestro/auth/module-identity", "session.open", "p-1",
|
||||
auth_custom=True,
|
||||
)
|
||||
assert url == (
|
||||
"http://maestro/auth/module-identity?permission=session.open&project_uuid=p-1"
|
||||
)
|
||||
|
||||
def test_maestro_mode_without_project_uuid():
|
||||
url = service_access_auth_url(
|
||||
"http://maestro/auth/module-identity", "project.view", None,
|
||||
auth_custom=True,
|
||||
)
|
||||
assert url == "http://maestro/auth/module-identity?permission=project.view"
|
||||
|
||||
def test_local_mode_unchanged():
|
||||
url = service_access_auth_url(
|
||||
"http://auth-server/auth", "session.open", "p-1", auth_custom=False,
|
||||
)
|
||||
assert url == (
|
||||
"http://auth-server/auth/service-access?permission=session.open&project_uuid=p-1"
|
||||
)
|
||||
```
|
||||
|
||||
- [ ] **Step 2: Run test to verify it fails**
|
||||
|
||||
Run (in the pod or a venv with the lib on path):
|
||||
`python -m pytest lib/python/orchest-internals/tests/test_service_access_auth_url.py -q`
|
||||
Expected: FAIL on the two Maestro cases (current code returns the base URL unchanged).
|
||||
|
||||
- [ ] **Step 3: Write minimal implementation**
|
||||
|
||||
Replace the `if auth_custom:` short-circuit in `utils.py`:
|
||||
|
||||
```python
|
||||
if auth_custom:
|
||||
# Maestro mode: no /auth/service-access sibling route — the unified
|
||||
# /auth/module-identity route authorizes when the annotation carries
|
||||
# the scope. Append the same permission/project_uuid params. (Maestro
|
||||
# derives the tenant orchest-api host from the JWT, not from the URL.)
|
||||
query = f"permission={permission}"
|
||||
if project_uuid:
|
||||
query += f"&project_uuid={project_uuid}"
|
||||
return f"{base_auth_url}?{query}"
|
||||
```
|
||||
|
||||
Update the docstring's "falls back to the plain base_auth_url" paragraph to describe the new scoped behavior.
|
||||
|
||||
- [ ] **Step 4: Run test to verify it passes**
|
||||
|
||||
Run: `python -m pytest lib/python/orchest-internals/tests/test_service_access_auth_url.py -q`
|
||||
Expected: PASS (all three).
|
||||
|
||||
- [ ] **Step 5: Commit (dbt-to-orchest repo)**
|
||||
|
||||
```bash
|
||||
git add lib/python/orchest-internals/_orchest/internals/utils.py lib/python/orchest-internals/tests/test_service_access_auth_url.py
|
||||
git commit -m "feat(rbac): scope Maestro-mode service auth-url (close direct-URL bypass)"
|
||||
```
|
||||
|
||||
---
|
||||
|
||||
## Deployment & manual verification (after Tasks 1-6)
|
||||
|
||||
Not code steps — run after merging, per the spec §8/§10.
|
||||
|
||||
- [ ] **Maestro env** on the Orchest-serving deployment: `ORCHEST_MODULE_PERMISSION_SEQID=31`, `ORCHEST_ADMIN_PERMISSION_SEQIDS=34`, `ORCHEST_NAMESPACE_MODULE=intelli`.
|
||||
- [ ] **OrchestCluster spec** (webserver ingress): `AuthCustom: true`, `AuthUrl: http://maestro.<maestro-ns>.svc.cluster.local/auth/module-identity`, `AuthSignin: <platform login URL>`.
|
||||
- [ ] **⚠️ celery-worker rebuild (Phase 2 / Task 6):** `service_access_auth_url` runs in the **celery-worker baked image** — rebuild + roll it, and verify the scoped annotation on a **freshly launched** service ingress (`kubectl get ingress …`), never an existing one. (This gotcha cost a debugging cycle last week; see `docs/source/development/building_images_minikube.md`.)
|
||||
- [ ] **Verify path A (spec §10):** curl-drive orchest-api with the `X-Auth-*` headers Maestro would set and watch the RBAC ladder + `external_admin` resolve, before standing up real Maestro.
|
||||
- [ ] **Forgery checks (spec §9.11, §9.12):** a non-admin's forged `X-Auth-Roles: admin` through the ingress must not reach orchest-api as admin; `?permission=…` appended to the app URL must not trigger an authz check.
|
||||
|
||||
---
|
||||
|
||||
## Self-Review
|
||||
|
||||
**Spec coverage:**
|
||||
- §4.1 cookie auth + JWKS reuse → Task 2 (expose verify) + Task 3 (read cookie).
|
||||
- §4.2 ladder (401/403/200 + authz + all-or-nothing + no refresh) → Task 3 (401/403/200) + Task 5 (authz, 502, all-or-nothing).
|
||||
- §4.3 three headers, no email → Task 3 (sets exactly the three; email never set).
|
||||
- §5 module gate + admin set (config, list) → Task 1 + Task 3.
|
||||
- §5.1 tenant routing (namespace pattern, allowlist, slug note, dadosferademo2) → Task 4 + Task 5.
|
||||
- §6 forgery boundary → deployment verification checklist.
|
||||
- §7 service_access_auth_url scoping → Task 6.
|
||||
- §8 deployment config → deployment checklist.
|
||||
- §9 tests 1-12 → Tasks 1/3 (1-5), Task 5 (6-10), deployment checklist (11-12).
|
||||
- §10 minikube verification → deployment checklist.
|
||||
- §11 phase split → Phase 1 (Tasks 1-3) / Phase 2 (Tasks 4-6).
|
||||
|
||||
**Placeholder scan:** none — every code step has real content; the only deliberately-open items are flagged verify-steps (HTTP client choice in Task 5; `customer_name` slug normalization in Task 4), each with an explicit instruction to inspect rather than guess.
|
||||
|
||||
**Type consistency:** `verifyAccessToken` (Task 2) consumed in Tasks 3/5; `moduleIdentity(req,res)` signature identical across Tasks 3/5; helper names (`parseAdminSeqids`, `moduleGateSeqid`, `isModuleAllowed`, `isAdmin`, `tenantOrchestApiHost`, `isAllowedOrchestApiHost`, `namespaceModule`) defined in Tasks 1/4 and used verbatim in Tasks 3/5; `authorizeOrchestServiceAccess` return union `'allow'|'deny'|'error'` consistent between service (Task 5 3a) and controller (Task 5 3b).
|
||||
+1878
-1377
File diff suppressed because it is too large
Load Diff
@@ -5,6 +5,9 @@ const config: Config.InitialOptions = {
|
||||
roots: ['<rootDir>/src/', '<rootDir>/test/'],
|
||||
testRegex: '.*\\.(test|spec)\\.[jt]s$',
|
||||
transform: { '\\.[jt]s$': 'ts-jest' },
|
||||
// Dummy values for env vars read at module-import time (see the setup file),
|
||||
// so specs importing those modules don't crash on load.
|
||||
setupFiles: ['<rootDir>/test/jest.setup-env.ts'],
|
||||
collectCoverageFrom: ['**/*.[jt]s'],
|
||||
coverageDirectory: 'coverage',
|
||||
coveragePathIgnorePatterns: [
|
||||
|
||||
Generated
+83
-20
@@ -16,8 +16,7 @@
|
||||
"@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.26",
|
||||
"@dadosfera/protospack-v2": "^3.40.0-beta.14",
|
||||
"@grpc/grpc-js": "^1.9.3",
|
||||
"@grpc/proto-loader": "^0.7.9",
|
||||
"@nestjs/cli": "^9.5.0",
|
||||
@@ -31,7 +30,7 @@
|
||||
"@nestjs/schematics": "^9.2.0",
|
||||
"@nestjs/swagger": "^6.3.0",
|
||||
"@nestjs/testing": "^9.4.3",
|
||||
"axios": "^0.30.2",
|
||||
"axios": "0.30.3",
|
||||
"cache-manager": "^5.1.4",
|
||||
"cache-manager-ioredis-yet": "^1.1.0",
|
||||
"class-transformer": "^0.5.1",
|
||||
@@ -47,6 +46,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",
|
||||
@@ -1734,19 +1734,10 @@
|
||||
"winston-log2gelf": "^2.4.0"
|
||||
}
|
||||
},
|
||||
"node_modules/@dadosfera/protospack": {
|
||||
"version": "2.5.3",
|
||||
"resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack/-/protospack-2.5.3.tgz",
|
||||
"integrity": "sha512-yOLnd+s6n9VkPpZXO8HnUY27CQPHj/qs+ecddviA4Ldn0Gx4KGRgbVdsSSP45nPm0GHhCd2bHg4ap+la7xtRmA==",
|
||||
"license": "ISC",
|
||||
"dependencies": {
|
||||
"rxjs": "^7.5.5"
|
||||
}
|
||||
},
|
||||
"node_modules/@dadosfera/protospack-v2": {
|
||||
"version": "3.38.0-beta.26",
|
||||
"resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.38.0-beta.26.tgz",
|
||||
"integrity": "sha512-N8NS7+djLGy0wJXk00+4oupqd/wBIQ1f+YBBhK2y9x4guFXYK1KWIrPZhPj+gaU8g2KNkqKoT7SnE9PXNxlLSQ==",
|
||||
"version": "3.40.0-beta.14",
|
||||
"resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.40.0-beta.14.tgz",
|
||||
"integrity": "sha512-pv3pxq0x1XcBgf3ajD6QOFRLOduh8iEozKFA3AKlIW4gid+gT4iL0GcU2M+O7h0QFeO4JIzRZe/nEMN82nqk7A==",
|
||||
"dependencies": {
|
||||
"@grpc/grpc-js": "^1.9.3",
|
||||
"rxjs": "^7.5.5"
|
||||
@@ -2927,6 +2918,20 @@
|
||||
"node": ">= 0.6"
|
||||
}
|
||||
},
|
||||
"node_modules/@nestjs/platform-express/node_modules/concat-stream": {
|
||||
"version": "1.6.2",
|
||||
"resolved": "https://registry.npmjs.org/concat-stream/-/concat-stream-1.6.2.tgz",
|
||||
"integrity": "sha512-27HBghJxjiZtIk3Ycvn/4kbJk/1uZuJFfuPEns6LaEvpvG1f0hTea8lilrouyo9mVc2GWdcEZ8OLoGmSADlrCw==",
|
||||
"engines": [
|
||||
"node >= 0.8"
|
||||
],
|
||||
"dependencies": {
|
||||
"buffer-from": "^1.0.0",
|
||||
"inherits": "^2.0.3",
|
||||
"readable-stream": "^2.2.2",
|
||||
"typedarray": "^0.0.6"
|
||||
}
|
||||
},
|
||||
"node_modules/@nestjs/platform-express/node_modules/content-disposition": {
|
||||
"version": "0.5.4",
|
||||
"resolved": "https://registry.npmjs.org/content-disposition/-/content-disposition-0.5.4.tgz",
|
||||
@@ -3096,6 +3101,24 @@
|
||||
"integrity": "sha512-Tpp60P6IUJDTuOq/5Z8cdskzJujfwqfOTkrwIwj7IRISpnkJnT6SyJ4PCPnGMoFjC9ddhal5KVIYtAt97ix05A==",
|
||||
"license": "MIT"
|
||||
},
|
||||
"node_modules/@nestjs/platform-express/node_modules/multer": {
|
||||
"version": "1.4.4-lts.1",
|
||||
"resolved": "https://registry.npmjs.org/multer/-/multer-1.4.4-lts.1.tgz",
|
||||
"integrity": "sha512-WeSGziVj6+Z2/MwQo3GvqzgR+9Uc+qt8SwHKh3gvNPiISKfsMfG4SvCOFYlxxgkXt7yIV2i1yczehm0EOKIxIg==",
|
||||
"deprecated": "Multer 1.x is impacted by a number of vulnerabilities, which have been patched in 2.x. You should upgrade to the latest 2.x version.",
|
||||
"dependencies": {
|
||||
"append-field": "^1.0.0",
|
||||
"busboy": "^1.0.0",
|
||||
"concat-stream": "^1.5.2",
|
||||
"mkdirp": "^0.5.4",
|
||||
"object-assign": "^4.1.1",
|
||||
"type-is": "^1.6.4",
|
||||
"xtend": "^4.0.0"
|
||||
},
|
||||
"engines": {
|
||||
"node": ">= 6.0.0"
|
||||
}
|
||||
},
|
||||
"node_modules/@nestjs/platform-express/node_modules/negotiator": {
|
||||
"version": "0.6.3",
|
||||
"resolved": "https://registry.npmjs.org/negotiator/-/negotiator-0.6.3.tgz",
|
||||
@@ -3120,6 +3143,25 @@
|
||||
"url": "https://github.com/sponsors/ljharb"
|
||||
}
|
||||
},
|
||||
"node_modules/@nestjs/platform-express/node_modules/readable-stream": {
|
||||
"version": "2.3.8",
|
||||
"resolved": "https://registry.npmjs.org/readable-stream/-/readable-stream-2.3.8.tgz",
|
||||
"integrity": "sha512-8p0AUk4XODgIewSi0l8Epjs+EVnWiK7NoDIEGU0HhE7+ZyY8D1IMY7odu5lRrFXGg71L15KG8QrPmum45RTtdA==",
|
||||
"dependencies": {
|
||||
"core-util-is": "~1.0.0",
|
||||
"inherits": "~2.0.3",
|
||||
"isarray": "~1.0.0",
|
||||
"process-nextick-args": "~2.0.0",
|
||||
"safe-buffer": "~5.1.1",
|
||||
"string_decoder": "~1.1.1",
|
||||
"util-deprecate": "~1.0.1"
|
||||
}
|
||||
},
|
||||
"node_modules/@nestjs/platform-express/node_modules/readable-stream/node_modules/safe-buffer": {
|
||||
"version": "5.1.2",
|
||||
"resolved": "https://registry.npmjs.org/safe-buffer/-/safe-buffer-5.1.2.tgz",
|
||||
"integrity": "sha512-Gd2UZBJDkXlY7GbJxfsE8/nvKkUEU1G38c1siN6QP6a9PT9MmHB8GnpscSmMJSoF8LOIrt8ud/wPtojys4G6+g=="
|
||||
},
|
||||
"node_modules/@nestjs/platform-express/node_modules/safe-buffer": {
|
||||
"version": "5.2.1",
|
||||
"resolved": "https://registry.npmjs.org/safe-buffer/-/safe-buffer-5.2.1.tgz",
|
||||
@@ -3194,6 +3236,19 @@
|
||||
"node": ">= 0.8"
|
||||
}
|
||||
},
|
||||
"node_modules/@nestjs/platform-express/node_modules/string_decoder": {
|
||||
"version": "1.1.1",
|
||||
"resolved": "https://registry.npmjs.org/string_decoder/-/string_decoder-1.1.1.tgz",
|
||||
"integrity": "sha512-n/ShnvDi6FHbbVfviro+WojiFzv+s8MPMHBczVePfUpDJLwoLT0ht1l4YwBCbi8pJAveEEdnkHyPyTP/mzRfwg==",
|
||||
"dependencies": {
|
||||
"safe-buffer": "~5.1.0"
|
||||
}
|
||||
},
|
||||
"node_modules/@nestjs/platform-express/node_modules/string_decoder/node_modules/safe-buffer": {
|
||||
"version": "5.1.2",
|
||||
"resolved": "https://registry.npmjs.org/safe-buffer/-/safe-buffer-5.1.2.tgz",
|
||||
"integrity": "sha512-Gd2UZBJDkXlY7GbJxfsE8/nvKkUEU1G38c1siN6QP6a9PT9MmHB8GnpscSmMJSoF8LOIrt8ud/wPtojys4G6+g=="
|
||||
},
|
||||
"node_modules/@nestjs/platform-express/node_modules/tslib": {
|
||||
"version": "2.5.3",
|
||||
"resolved": "https://registry.npmjs.org/tslib/-/tslib-2.5.3.tgz",
|
||||
@@ -5533,9 +5588,9 @@
|
||||
}
|
||||
},
|
||||
"node_modules/axios": {
|
||||
"version": "0.30.2",
|
||||
"resolved": "https://registry.npmjs.org/axios/-/axios-0.30.2.tgz",
|
||||
"integrity": "sha512-0pE4RQ4UQi1jKY6p7u6i1Tkzqmu+d+/tHS7Q7rKunWLB9WyilBTpHHpXzPNMDj5hTbK0B0PTLSz07yqMBiF6xg==",
|
||||
"version": "0.30.3",
|
||||
"resolved": "https://registry.npmjs.org/axios/-/axios-0.30.3.tgz",
|
||||
"integrity": "sha512-5/tmEb6TmE/ax3mdXBc/Mi6YdPGxQsv+0p5YlciXWt3PHIn0VamqCXhRMtScnwY3lbgSXLneOuXAKUhgmSRpwg==",
|
||||
"license": "MIT",
|
||||
"dependencies": {
|
||||
"follow-redirects": "^1.15.4",
|
||||
@@ -6508,7 +6563,6 @@
|
||||
"engines": [
|
||||
"node >= 6.0"
|
||||
],
|
||||
"license": "MIT",
|
||||
"dependencies": {
|
||||
"buffer-from": "^1.0.0",
|
||||
"inherits": "^2.0.3",
|
||||
@@ -9127,6 +9181,11 @@
|
||||
"url": "https://github.com/sponsors/sindresorhus"
|
||||
}
|
||||
},
|
||||
"node_modules/isarray": {
|
||||
"version": "1.0.0",
|
||||
"resolved": "https://registry.npmjs.org/isarray/-/isarray-1.0.0.tgz",
|
||||
"integrity": "sha512-VLghIWNM6ELQzo7zwmcg0NmTVyWKYjvIeM83yjp0wRDTmUnrM678fQbcKBo6n2CJEF0szoG//ytg+TKla89ALQ=="
|
||||
},
|
||||
"node_modules/isexe": {
|
||||
"version": "2.0.0",
|
||||
"resolved": "https://registry.npmjs.org/isexe/-/isexe-2.0.0.tgz",
|
||||
@@ -10724,7 +10783,6 @@
|
||||
"version": "2.0.2",
|
||||
"resolved": "https://registry.npmjs.org/multer/-/multer-2.0.2.tgz",
|
||||
"integrity": "sha512-u7f2xaZ/UG8oLXHvtF/oWTRvT44p9ecwBBqTwgJVq0+4BW1g8OW01TyMEGWBHbyMOYVHXslaut7qEQ1meATXgw==",
|
||||
"license": "MIT",
|
||||
"dependencies": {
|
||||
"append-field": "^1.0.0",
|
||||
"busboy": "^1.6.0",
|
||||
@@ -11727,6 +11785,11 @@
|
||||
"url": "https://github.com/chalk/ansi-styles?sponsor=1"
|
||||
}
|
||||
},
|
||||
"node_modules/process-nextick-args": {
|
||||
"version": "2.0.1",
|
||||
"resolved": "https://registry.npmjs.org/process-nextick-args/-/process-nextick-args-2.0.1.tgz",
|
||||
"integrity": "sha512-3ouUOpQhtgrbOa17J7+uxOTpITYWaGP7/AhoR3+A+/1e9skrzelGi/dXzEYyvbxubEF6Wn2ypscTKiKJFFn1ag=="
|
||||
},
|
||||
"node_modules/process-warning": {
|
||||
"version": "1.0.0",
|
||||
"resolved": "https://registry.npmjs.org/process-warning/-/process-warning-1.0.0.tgz",
|
||||
|
||||
+8
-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\"",
|
||||
@@ -34,8 +34,7 @@
|
||||
"@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.26",
|
||||
"@dadosfera/protospack-v2": "^3.40.0-beta.14",
|
||||
"@grpc/grpc-js": "^1.9.3",
|
||||
"@grpc/proto-loader": "^0.7.9",
|
||||
"@nestjs/cli": "^9.5.0",
|
||||
@@ -49,7 +48,7 @@
|
||||
"@nestjs/schematics": "^9.2.0",
|
||||
"@nestjs/swagger": "^6.3.0",
|
||||
"@nestjs/testing": "^9.4.3",
|
||||
"axios": "^0.30.2",
|
||||
"axios": "0.30.3",
|
||||
"cache-manager": "^5.1.4",
|
||||
"cache-manager-ioredis-yet": "^1.1.0",
|
||||
"class-transformer": "^0.5.1",
|
||||
@@ -65,6 +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",
|
||||
@@ -80,7 +80,7 @@
|
||||
"swagger-ui-express": "^4.6.3"
|
||||
},
|
||||
"overrides": {
|
||||
"multer": "2.0.2",
|
||||
"axios": "0.30.3",
|
||||
"form-data": "^4.0.4",
|
||||
"body-parser": "^1.20.3",
|
||||
"cross-spawn": "^7.0.5",
|
||||
@@ -117,5 +117,8 @@
|
||||
"ts-node": "^10.9.1",
|
||||
"tsconfig-paths": "^3.14.2",
|
||||
"typescript": "^4.9.5"
|
||||
},
|
||||
"resolutions": {
|
||||
"axios": "0.30.3"
|
||||
}
|
||||
}
|
||||
|
||||
+3
-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';
|
||||
@@ -35,6 +34,8 @@ import { ShareMetadataModule } from './modules/share-metadata/share-metadata.mod
|
||||
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: [
|
||||
@@ -58,7 +59,6 @@ import { StorageExplorerModule } from './modules/storage-explorer/storage-explor
|
||||
PermissionsModule,
|
||||
TermsOfUseModule,
|
||||
ConnectionTestModule,
|
||||
PipelinesModule,
|
||||
TransformationsModule,
|
||||
UsersModule,
|
||||
RolesModule,
|
||||
@@ -79,6 +79,7 @@ import { StorageExplorerModule } from './modules/storage-explorer/storage-explor
|
||||
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,
|
||||
|
||||
@@ -11,6 +11,7 @@ export function extractUserFrom(aRawJwt: string) {
|
||||
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,
|
||||
|
||||
@@ -357,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',
|
||||
@@ -712,6 +722,8 @@ export const DADOSFERA_MODULES_KEYS = {
|
||||
PII: 'pii',
|
||||
EMBED: 'embedded-analytics',
|
||||
EMBED_ASSIGNED: 'embed-assigned',
|
||||
CATALOG: 'catalog',
|
||||
COLLECT: 'collect',
|
||||
}
|
||||
|
||||
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);
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
@@ -111,3 +111,4 @@ function configureSwagger(app: INestApplication) {
|
||||
);
|
||||
}
|
||||
bootstrap();
|
||||
|
||||
|
||||
@@ -32,6 +32,15 @@ import {
|
||||
} from '@dadosfera/protospack-v2/dist/lib/Duc/interfaces/messages';
|
||||
|
||||
import { PERMISSIONS_GROUPS } from 'src/authentication/permissions.enum';
|
||||
import {
|
||||
parseAdminSeqids,
|
||||
moduleGateSeqid,
|
||||
isModuleAllowed,
|
||||
isAdmin,
|
||||
tenantOrchestApiHost,
|
||||
isAllowedOrchestApiHost,
|
||||
namespaceModule,
|
||||
} from './orchest-identity';
|
||||
import {
|
||||
Authenticated,
|
||||
RequireAllPermissions,
|
||||
@@ -478,6 +487,7 @@ export class AuthController {
|
||||
@Get('me')
|
||||
async getMe(@Req() req: Request, @Res() res: Response) {
|
||||
this.logger.info('GET /auth/me ')
|
||||
this.logger.info(JSON.stringify(req.headers));
|
||||
|
||||
// Check for API key header first
|
||||
const apiKey = req.get('X-Api-key');
|
||||
@@ -488,6 +498,7 @@ export class AuthController {
|
||||
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,
|
||||
@@ -529,4 +540,84 @@ export class AuthController {
|
||||
return res.status(200).json(user);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Ingress auth subrequest for Orchest (see dbt-to-orchest
|
||||
* docs/superpowers/specs/2026-08-20-maestro-module-identity-design.md).
|
||||
*
|
||||
* Authenticates the ddf-auth cookie, gates on the module permission, and
|
||||
* returns the X-Auth-* identity headers orchest-api trusts. Never writes
|
||||
* cookies (an auth_request response cannot). No token refresh: an expired
|
||||
* token is a 401 (the signin flow re-auths).
|
||||
*/
|
||||
@Get('module-identity')
|
||||
async moduleIdentity(@Req() req: Request, @Res() res: Response) {
|
||||
const token = req.cookies?.['ddf-auth'];
|
||||
if (!token) {
|
||||
return res.status(401).send();
|
||||
}
|
||||
|
||||
let payload: any;
|
||||
try {
|
||||
payload = await this.authClient.verifyAccessToken(token);
|
||||
} catch (e) {
|
||||
return res.status(401).send();
|
||||
}
|
||||
|
||||
const perms: number[] = payload?.permissions ?? [];
|
||||
const gate = moduleGateSeqid(process.env.ORCHEST_MODULE_PERMISSION_SEQID);
|
||||
if (!isModuleAllowed(perms, gate)) {
|
||||
return res.status(403).send();
|
||||
}
|
||||
|
||||
const adminUser = isAdmin(
|
||||
perms,
|
||||
parseAdminSeqids(process.env.ORCHEST_ADMIN_PERMISSION_SEQIDS),
|
||||
);
|
||||
|
||||
res.set('X-Auth-User', String(payload.user_id));
|
||||
res.set('X-Auth-Username', String(payload.username ?? ''));
|
||||
if (adminUser) {
|
||||
res.set('X-Auth-Roles', 'admin');
|
||||
}
|
||||
|
||||
// Service-ingress caller: the orchest-api-authored annotation carries the
|
||||
// scope. Authorize per-project against the tenant's orchest-api. A
|
||||
// param-less request is the webserver-ingress case → identity only.
|
||||
const permission = req.query?.permission as string | undefined;
|
||||
const projectUuid = req.query?.project_uuid as string | undefined;
|
||||
if (permission) {
|
||||
// All-or-nothing: an incomplete annotation must not silently skip the
|
||||
// per-project check.
|
||||
if (!projectUuid) {
|
||||
return res.status(403).send();
|
||||
}
|
||||
const module = namespaceModule(process.env.ORCHEST_NAMESPACE_MODULE);
|
||||
const host = tenantOrchestApiHost(
|
||||
String(payload.customer_name ?? ''),
|
||||
module,
|
||||
);
|
||||
if (!isAllowedOrchestApiHost(host)) {
|
||||
return res.status(403).send();
|
||||
}
|
||||
const identityHeaders: Record<string, string> = {
|
||||
'X-Auth-User': String(payload.user_id),
|
||||
'X-Auth-Username': String(payload.username ?? ''),
|
||||
};
|
||||
if (adminUser) {
|
||||
identityHeaders['X-Auth-Roles'] = 'admin';
|
||||
}
|
||||
const decision = await this.authClient.authorizeOrchestServiceAccess({
|
||||
host,
|
||||
permission,
|
||||
projectUuid,
|
||||
headers: identityHeaders,
|
||||
});
|
||||
if (decision === 'deny') return res.status(403).send();
|
||||
if (decision === 'error') return res.status(502).send();
|
||||
// 'allow' falls through to the 200 (identity headers already set).
|
||||
}
|
||||
|
||||
return res.status(200).send();
|
||||
}
|
||||
}
|
||||
|
||||
@@ -9,6 +9,7 @@ import {
|
||||
import { ClientGrpc } from '@nestjs/microservices';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import { lastValueFrom } from 'rxjs';
|
||||
import axios from 'axios';
|
||||
|
||||
import { ProtoServices } from '@dadosfera/protospack-v2/dist/lib/Duc';
|
||||
import {
|
||||
@@ -307,7 +308,7 @@ export class AuthClientService implements OnModuleInit {
|
||||
}
|
||||
|
||||
public async validateUserSession(accessToken: any, resourceHost: string) {
|
||||
const payload = await this.validateJwtToken(accessToken);
|
||||
const payload = await this.verifyAccessToken(accessToken);
|
||||
|
||||
const userDto = await this.getUserfromPayload(payload);
|
||||
|
||||
@@ -410,7 +411,12 @@ export class AuthClientService implements OnModuleInit {
|
||||
this.logger.info('Clean cookie sessions');
|
||||
}
|
||||
|
||||
private async validateJwtToken(token: string) {
|
||||
/**
|
||||
* Verify a DUC access token against the JWKS and return its payload.
|
||||
* Public so the /auth/module-identity route can authenticate the ddf-auth
|
||||
* cookie the same way (was the private validateJwtToken).
|
||||
*/
|
||||
public async verifyAccessToken(token: string) {
|
||||
const decoded: any = token && jwt.decode(token, { complete: true });
|
||||
if (!decoded) throw new Error('Invalid token');
|
||||
|
||||
@@ -437,7 +443,8 @@ export class AuthClientService implements OnModuleInit {
|
||||
|
||||
const userDto: UserDTO = {
|
||||
id: user.id,
|
||||
name: user.username,
|
||||
name: user.name,
|
||||
email: user.email,
|
||||
jobTitle: user?.jobTitle || null,
|
||||
department: user?.department || null,
|
||||
hierarchy: user?.hierarchy || null,
|
||||
@@ -491,4 +498,33 @@ export class AuthClientService implements OnModuleInit {
|
||||
|
||||
return;
|
||||
}
|
||||
|
||||
/**
|
||||
* Ask a tenant's orchest-api whether the acting user holds `permission` on
|
||||
* `projectUuid` (the per-service authz half of /auth/module-identity).
|
||||
* Fail-closed: any error / unexpected status → 'error' (the route denies).
|
||||
*/
|
||||
public async authorizeOrchestServiceAccess(args: {
|
||||
host: string;
|
||||
permission: string;
|
||||
projectUuid?: string;
|
||||
headers: Record<string, string>;
|
||||
}): Promise<'allow' | 'deny' | 'error'> {
|
||||
const params: Record<string, string> = { permission: args.permission };
|
||||
if (args.projectUuid) params.project_uuid = args.projectUuid;
|
||||
try {
|
||||
const resp = await axios.get(`http://${args.host}/api/authz/check`, {
|
||||
params,
|
||||
headers: args.headers,
|
||||
timeout: 5000,
|
||||
// Never throw on 4xx/5xx; branch on the status ourselves.
|
||||
validateStatus: () => true,
|
||||
});
|
||||
if (resp.status === 200) return 'allow';
|
||||
if (resp.status === 403) return 'deny';
|
||||
return 'error';
|
||||
} catch (e) {
|
||||
return 'error';
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -144,6 +144,7 @@ export interface BulkEditResponse {
|
||||
export type UserDTO = {
|
||||
id: string,
|
||||
name: string,
|
||||
email: string,
|
||||
jobTitle?: string,
|
||||
department?: string,
|
||||
hierarchy?: string,
|
||||
|
||||
@@ -0,0 +1,206 @@
|
||||
import { Test } from '@nestjs/testing';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import { AuthController } from './auth.controller';
|
||||
import { AuthClientService } from './auth.service';
|
||||
import { ApiKeyService } from 'src/modules/api-key/api-key.service';
|
||||
|
||||
function res() {
|
||||
const headers: Record<string, string> = {};
|
||||
const r: any = {
|
||||
_status: 0,
|
||||
_sent: undefined,
|
||||
set: (k: string, v: string) => {
|
||||
headers[k] = v;
|
||||
return r;
|
||||
},
|
||||
status: (c: number) => {
|
||||
r._status = c;
|
||||
return r;
|
||||
},
|
||||
send: (b?: any) => {
|
||||
r._sent = b ?? '';
|
||||
return r;
|
||||
},
|
||||
json: (b?: any) => {
|
||||
r._sent = b;
|
||||
return r;
|
||||
},
|
||||
_headers: headers,
|
||||
};
|
||||
return r;
|
||||
}
|
||||
|
||||
function req(cookie?: string, query: Record<string, string> = {}) {
|
||||
return { cookies: cookie ? { 'ddf-auth': cookie } : {}, query } as any;
|
||||
}
|
||||
|
||||
const loggerStub = {
|
||||
logger: { info: jest.fn(), error: jest.fn(), warn: jest.fn() },
|
||||
} as unknown as DadosferaLogger;
|
||||
|
||||
async function makeController(auth: AuthClientService): Promise<AuthController> {
|
||||
const mod = await Test.createTestingModule({
|
||||
controllers: [AuthController],
|
||||
providers: [
|
||||
{ provide: DadosferaLogger, useValue: loggerStub },
|
||||
{ provide: AuthClientService, useValue: auth },
|
||||
{ provide: ApiKeyService, useValue: {} },
|
||||
],
|
||||
}).compile();
|
||||
return mod.get(AuthController);
|
||||
}
|
||||
|
||||
describe('GET /auth/module-identity — identity', () => {
|
||||
let controller: AuthController;
|
||||
const auth = {
|
||||
verifyAccessToken: jest.fn(),
|
||||
} as unknown as AuthClientService;
|
||||
|
||||
beforeEach(async () => {
|
||||
jest.resetAllMocks();
|
||||
process.env.ORCHEST_MODULE_PERMISSION_SEQID = '31';
|
||||
process.env.ORCHEST_ADMIN_PERMISSION_SEQIDS = '34';
|
||||
controller = await makeController(auth);
|
||||
});
|
||||
|
||||
it('no cookie → 401', async () => {
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req(undefined), r);
|
||||
expect(r._status).toBe(401);
|
||||
});
|
||||
|
||||
it('valid + module + admin → 200 with X-Auth-Roles: admin', async () => {
|
||||
(auth.verifyAccessToken as jest.Mock).mockResolvedValue({
|
||||
user_id: 'u-1',
|
||||
username: 'alice',
|
||||
permissions: [31, 34],
|
||||
});
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok'), r);
|
||||
expect(r._status).toBe(200);
|
||||
expect(r._headers['X-Auth-User']).toBe('u-1');
|
||||
expect(r._headers['X-Auth-Username']).toBe('alice');
|
||||
expect(r._headers['X-Auth-Roles']).toBe('admin');
|
||||
});
|
||||
|
||||
it('valid + module, not admin → 200, no X-Auth-Roles', async () => {
|
||||
(auth.verifyAccessToken as jest.Mock).mockResolvedValue({
|
||||
user_id: 'u-2',
|
||||
username: 'bob',
|
||||
permissions: [31],
|
||||
});
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok'), r);
|
||||
expect(r._status).toBe(200);
|
||||
expect(r._headers['X-Auth-Roles']).toBeUndefined();
|
||||
});
|
||||
|
||||
it('valid, lacks module → 403', async () => {
|
||||
(auth.verifyAccessToken as jest.Mock).mockResolvedValue({
|
||||
user_id: 'u-3',
|
||||
username: 'carol',
|
||||
permissions: [5],
|
||||
});
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok'), r);
|
||||
expect(r._status).toBe(403);
|
||||
});
|
||||
|
||||
it('verify throws (expired/bad) → 401', async () => {
|
||||
(auth.verifyAccessToken as jest.Mock).mockRejectedValue(new Error('bad'));
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok'), r);
|
||||
expect(r._status).toBe(401);
|
||||
});
|
||||
|
||||
it('admin seqids extended by config → 200 admin', async () => {
|
||||
process.env.ORCHEST_ADMIN_PERMISSION_SEQIDS = '34,99';
|
||||
(auth.verifyAccessToken as jest.Mock).mockResolvedValue({
|
||||
user_id: 'u-4',
|
||||
username: 'dana',
|
||||
permissions: [31, 99],
|
||||
});
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok'), r);
|
||||
expect(r._headers['X-Auth-Roles']).toBe('admin');
|
||||
});
|
||||
});
|
||||
|
||||
describe('GET /auth/module-identity — per-service authz', () => {
|
||||
let controller: AuthController;
|
||||
const auth = {
|
||||
verifyAccessToken: jest.fn(),
|
||||
authorizeOrchestServiceAccess: jest.fn(),
|
||||
} as unknown as AuthClientService;
|
||||
|
||||
const q = { permission: 'session.open', project_uuid: 'p-1' };
|
||||
|
||||
beforeEach(async () => {
|
||||
jest.resetAllMocks();
|
||||
process.env.ORCHEST_MODULE_PERMISSION_SEQID = '31';
|
||||
process.env.ORCHEST_ADMIN_PERMISSION_SEQIDS = '34';
|
||||
process.env.ORCHEST_NAMESPACE_MODULE = 'intelli';
|
||||
controller = await makeController(auth);
|
||||
(auth.verifyAccessToken as jest.Mock).mockResolvedValue({
|
||||
user_id: 'u-1',
|
||||
username: 'alice',
|
||||
permissions: [31],
|
||||
customer_name: 'acme',
|
||||
});
|
||||
});
|
||||
|
||||
it('has grant → 200 and calls the tenant host', async () => {
|
||||
(auth.authorizeOrchestServiceAccess as jest.Mock).mockResolvedValue('allow');
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok', q), r);
|
||||
expect(r._status).toBe(200);
|
||||
expect(auth.authorizeOrchestServiceAccess).toHaveBeenCalledWith(
|
||||
expect.objectContaining({
|
||||
host: 'orchest-api.orchest-intelli-acme.svc.cluster.local',
|
||||
permission: 'session.open',
|
||||
projectUuid: 'p-1',
|
||||
}),
|
||||
);
|
||||
});
|
||||
|
||||
it('lacks grant → 403', async () => {
|
||||
(auth.authorizeOrchestServiceAccess as jest.Mock).mockResolvedValue('deny');
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok', q), r);
|
||||
expect(r._status).toBe(403);
|
||||
});
|
||||
|
||||
it('orchest-api error → 502 (fail-closed)', async () => {
|
||||
(auth.authorizeOrchestServiceAccess as jest.Mock).mockResolvedValue('error');
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok', q), r);
|
||||
expect(r._status).toBe(502);
|
||||
});
|
||||
|
||||
it('permission without project_uuid → 403 (all-or-nothing)', async () => {
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok', { permission: 'session.open' }), r);
|
||||
expect(r._status).toBe(403);
|
||||
expect(auth.authorizeOrchestServiceAccess).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('no authz params → 200 identity-only (webserver case)', async () => {
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok', {}), r);
|
||||
expect(r._status).toBe(200);
|
||||
expect(auth.authorizeOrchestServiceAccess).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('malformed customer_name → host fails allowlist → 403, no call', async () => {
|
||||
(auth.verifyAccessToken as jest.Mock).mockResolvedValue({
|
||||
user_id: 'u-1',
|
||||
username: 'alice',
|
||||
permissions: [31],
|
||||
customer_name: '',
|
||||
});
|
||||
const r = res();
|
||||
await controller.moduleIdentity(req('tok', q), r);
|
||||
expect(r._status).toBe(403);
|
||||
expect(auth.authorizeOrchestServiceAccess).not.toHaveBeenCalled();
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,76 @@
|
||||
import {
|
||||
parseAdminSeqids,
|
||||
moduleGateSeqid,
|
||||
isModuleAllowed,
|
||||
isAdmin,
|
||||
tenantOrchestApiHost,
|
||||
isAllowedOrchestApiHost,
|
||||
namespaceModule,
|
||||
} from './orchest-identity';
|
||||
|
||||
describe('orchest-identity mapping', () => {
|
||||
it('parseAdminSeqids: default, single, list, whitespace', () => {
|
||||
expect(parseAdminSeqids(undefined)).toEqual([34]);
|
||||
expect(parseAdminSeqids('')).toEqual([34]);
|
||||
expect(parseAdminSeqids('34')).toEqual([34]);
|
||||
expect(parseAdminSeqids('34,40')).toEqual([34, 40]);
|
||||
expect(parseAdminSeqids(' 34 , 40 ')).toEqual([34, 40]);
|
||||
});
|
||||
|
||||
it('moduleGateSeqid: default and override', () => {
|
||||
expect(moduleGateSeqid(undefined)).toBe(31);
|
||||
expect(moduleGateSeqid('43')).toBe(43);
|
||||
});
|
||||
|
||||
it('isModuleAllowed', () => {
|
||||
expect(isModuleAllowed([31, 5], 31)).toBe(true);
|
||||
expect(isModuleAllowed([5, 7], 31)).toBe(false);
|
||||
expect(isModuleAllowed([], 31)).toBe(false);
|
||||
});
|
||||
|
||||
it('isAdmin: intersection', () => {
|
||||
expect(isAdmin([31, 34], [34])).toBe(true);
|
||||
expect(isAdmin([31], [34])).toBe(false);
|
||||
expect(isAdmin([99], [34, 99])).toBe(true);
|
||||
});
|
||||
});
|
||||
|
||||
describe('orchest-identity tenant routing', () => {
|
||||
it('namespaceModule default and values', () => {
|
||||
expect(namespaceModule(undefined)).toBe('intelli');
|
||||
expect(namespaceModule('process')).toBe('process');
|
||||
expect(namespaceModule('garbage')).toBe('intelli');
|
||||
});
|
||||
|
||||
it('tenantOrchestApiHost builds the namespace pattern', () => {
|
||||
expect(tenantOrchestApiHost('acme', 'intelli')).toBe(
|
||||
'orchest-api.orchest-intelli-acme.svc.cluster.local',
|
||||
);
|
||||
expect(tenantOrchestApiHost('acme', 'process')).toBe(
|
||||
'orchest-api.orchest-process-acme.svc.cluster.local',
|
||||
);
|
||||
});
|
||||
|
||||
it('tenantOrchestApiHost normalizes a non-slug customer_name', () => {
|
||||
expect(tenantOrchestApiHost('Acme Corp', 'intelli')).toBe(
|
||||
'orchest-api.orchest-intelli-acme-corp.svc.cluster.local',
|
||||
);
|
||||
});
|
||||
|
||||
it('isAllowedOrchestApiHost guards against malformed values', () => {
|
||||
expect(
|
||||
isAllowedOrchestApiHost(
|
||||
'orchest-api.orchest-intelli-acme.svc.cluster.local',
|
||||
),
|
||||
).toBe(true);
|
||||
expect(isAllowedOrchestApiHost('evil.example.com')).toBe(false);
|
||||
expect(
|
||||
isAllowedOrchestApiHost('orchest-api.orchest-intelli-.svc.cluster.local'),
|
||||
).toBe(false);
|
||||
expect(
|
||||
isAllowedOrchestApiHost(
|
||||
'orchest-api.orchest-other-acme.svc.cluster.local',
|
||||
),
|
||||
).toBe(false);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,70 @@
|
||||
/**
|
||||
* Pure helpers backing `GET /auth/module-identity` (see
|
||||
* dbt-to-orchest docs/superpowers/specs/2026-08-20-maestro-module-identity-design.md).
|
||||
*
|
||||
* Kept free of HTTP/Nest so the permission mapping and tenant-routing logic
|
||||
* unit-test without a request. Authorization decisions that reach orchest-api
|
||||
* live on AuthClientService; this module only maps Maestro permissions to the
|
||||
* Orchest header contract and derives the tenant orchest-api host.
|
||||
*/
|
||||
|
||||
/** Parse ORCHEST_ADMIN_PERMISSION_SEQIDS ("34" / "34,40"); default [34]. */
|
||||
export function parseAdminSeqids(raw: string | undefined): number[] {
|
||||
if (!raw || !raw.trim()) return [34];
|
||||
return raw
|
||||
.split(',')
|
||||
.map((s) => Number(s.trim()))
|
||||
.filter((n) => Number.isInteger(n));
|
||||
}
|
||||
|
||||
/** Parse ORCHEST_MODULE_PERMISSION_SEQID; default 31 (Intelligence/Orchest). */
|
||||
export function moduleGateSeqid(raw: string | undefined): number {
|
||||
const n = Number(raw);
|
||||
return Number.isInteger(n) && n > 0 ? n : 31;
|
||||
}
|
||||
|
||||
/** Whether the user's permissions include the module-access gate seqid. */
|
||||
export function isModuleAllowed(perms: number[], gateSeqid: number): boolean {
|
||||
return Array.isArray(perms) && perms.includes(gateSeqid);
|
||||
}
|
||||
|
||||
/** Whether the user's permissions intersect the admin seqid set. */
|
||||
export function isAdmin(perms: number[], adminSeqids: number[]): boolean {
|
||||
return Array.isArray(perms) && perms.some((p) => adminSeqids.includes(p));
|
||||
}
|
||||
|
||||
/** The `{module}` slug in the tenant namespace pattern; default 'intelli'. */
|
||||
export function namespaceModule(
|
||||
raw: string | undefined,
|
||||
): 'intelli' | 'process' {
|
||||
return raw === 'process' ? 'process' : 'intelli';
|
||||
}
|
||||
|
||||
/**
|
||||
* The in-cluster DNS of the calling tenant's orchest-api, from the tenant
|
||||
* namespace convention `orchest-{module}-{customer_name}`.
|
||||
*
|
||||
* `customer_name` is normalized to the namespace slug shape (lowercase,
|
||||
* non-[a-z0-9-] → '-') so a display-name value still yields a valid host; a
|
||||
* value that is already a clean slug is unchanged. NOTE: confirm the exact
|
||||
* prod `customer_name` → namespace mapping against a real token before relying
|
||||
* on this in production (the `dadosferademo2` tenant is a known exception and
|
||||
* is unsupported — it is being removed).
|
||||
*/
|
||||
export function tenantOrchestApiHost(
|
||||
customerName: string,
|
||||
module: string,
|
||||
): string {
|
||||
const slug = String(customerName)
|
||||
.toLowerCase()
|
||||
.replace(/[^a-z0-9-]/g, '-');
|
||||
return `orchest-api.orchest-${module}-${slug}.svc.cluster.local`;
|
||||
}
|
||||
|
||||
const ORCHEST_API_HOST_RE =
|
||||
/^orchest-api\.orchest-(intelli|process)-[a-z0-9-]+\.svc\.cluster\.local$/;
|
||||
|
||||
/** Defense in depth: only call an orchest-api host matching the convention. */
|
||||
export function isAllowedOrchestApiHost(host: string): boolean {
|
||||
return ORCHEST_API_HOST_RE.test(host);
|
||||
}
|
||||
@@ -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,
|
||||
@@ -83,6 +85,9 @@ export class CatalogController {
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async searchCatalog(
|
||||
@User() user: RequestUser,
|
||||
@Query() query: ICatalogAllRequest,
|
||||
@@ -122,6 +127,9 @@ export class CatalogController {
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async dowloadAsserts(
|
||||
@User() user: RequestUser,
|
||||
@Query() query: ICatalogAllRequest,
|
||||
@@ -165,6 +173,9 @@ export class CatalogController {
|
||||
|
||||
@ApiInternalOnlyEndpoint()
|
||||
@Get('data-asset')
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async findByPipelineAndObject(@User() user: RequestUser, @Query() query) {
|
||||
const { username, user_id, customer_id, customer_name, permissions } = user;
|
||||
const { pipeline, object } = query;
|
||||
@@ -223,6 +234,9 @@ export class CatalogController {
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async findAllTags(@Body() body) {
|
||||
this.logger.info(`/catalog - ON FIND ALL TAGS ROUTE`, {
|
||||
user: body.info.user_id,
|
||||
@@ -241,11 +255,59 @@ export class CatalogController {
|
||||
return res;
|
||||
}
|
||||
|
||||
@Get('schemas')
|
||||
@RequireSomePermission(
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
async findSchemas(@User() user: RequestUser) {
|
||||
const { username, user_id, customer_id, customer_name } = user;
|
||||
this.logger.info(`/catalog - ON FIND SCHEMAS ROUTE`, {
|
||||
username,
|
||||
customer_name,
|
||||
});
|
||||
|
||||
const metadata = PackTheMetadata({
|
||||
username,
|
||||
user_id,
|
||||
customer_id,
|
||||
customer_name,
|
||||
});
|
||||
|
||||
try {
|
||||
const res = await this.catalogService.findSchemas(metadata);
|
||||
return res;
|
||||
} catch (error) {
|
||||
throw new HttpException(error.message, HttpStatus.NOT_FOUND);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@Get('custom-properties')
|
||||
@RequireSomePermission(
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
async getCustomPropertyDefinitions(@User() user: RequestUser) {
|
||||
const { customer_id, customer_name, user_id, username } = user;
|
||||
const metadata = PackTheMetadata({
|
||||
customer_id,
|
||||
customer_name,
|
||||
user_id,
|
||||
username,
|
||||
});
|
||||
|
||||
return this.catalogService.getCustomPropertyDefinitions(metadata);
|
||||
}
|
||||
|
||||
@Get('data-asset/:id')
|
||||
@RequireSomePermission(
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async getDataAsset(
|
||||
@User() user: RequestUser,
|
||||
@Param('id') id: string,
|
||||
@@ -357,6 +419,9 @@ export class CatalogController {
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async getDataAssetColumnsMetadata(
|
||||
@User() user: RequestUser,
|
||||
@Language() language: LanguageEnum,
|
||||
@@ -388,6 +453,9 @@ export class CatalogController {
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async getDataAssetPreview(
|
||||
@User() user: RequestUser,
|
||||
@Language() language: LanguageEnum,
|
||||
@@ -419,6 +487,9 @@ export class CatalogController {
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async getDataAssetDocs(
|
||||
@User() user: RequestUser,
|
||||
@Language() language: LanguageEnum,
|
||||
@@ -450,6 +521,9 @@ export class CatalogController {
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.UPDATE,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async updateDataAsset(
|
||||
@User() user: RequestUser,
|
||||
@Language() language: LanguageEnum,
|
||||
@@ -465,6 +539,8 @@ export class CatalogController {
|
||||
language,
|
||||
});
|
||||
|
||||
delete (body as any).certification_status;
|
||||
|
||||
const result = await this.catalogService.updateOneDataAsset({
|
||||
body,
|
||||
data_asset_id,
|
||||
@@ -478,11 +554,44 @@ export class CatalogController {
|
||||
return result;
|
||||
}
|
||||
|
||||
@Put('data-asset/:id/certification-status')
|
||||
@RequireSomePermission(
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.CERTIFY,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
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,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async manageDataAssetDocs(
|
||||
@User() user: RequestUser,
|
||||
@Headers() headers,
|
||||
@@ -520,6 +629,9 @@ export class CatalogController {
|
||||
|
||||
@ApiInternalOnlyEndpoint()
|
||||
@Put('data-asset/:id/manage-permissions')
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async manageDataAssetPermissions(
|
||||
@Param('id') id: string,
|
||||
@User() user: RequestUser,
|
||||
@@ -542,6 +654,9 @@ export class CatalogController {
|
||||
|
||||
@ApiInternalOnlyEndpoint()
|
||||
@Put('data-asset/:id/revoke-permissions')
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async revokeDataAssetPermissions(
|
||||
@Param('id') id: string,
|
||||
@User() user: RequestUser,
|
||||
@@ -567,6 +682,9 @@ export class CatalogController {
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.CREATE,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async createDataAsset(
|
||||
@User() user: RequestUser,
|
||||
@Body() body: ICreateDataAsset,
|
||||
@@ -591,6 +709,9 @@ export class CatalogController {
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.UPDATE,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async commentOnDataAsset(
|
||||
@Param('id') id: string,
|
||||
@User() user: RequestUser,
|
||||
@@ -617,6 +738,9 @@ export class CatalogController {
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DELETE,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async deleteDataAsset(@Param('id') id: string, @User() user: RequestUser) {
|
||||
const { customer_id, customer_name, user_id, username } = user;
|
||||
const metadata = PackTheMetadata({
|
||||
@@ -638,6 +762,9 @@ export class CatalogController {
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.UPDATE,
|
||||
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async deleteComment(
|
||||
@Param('id') id: string,
|
||||
@User() user: RequestUser,
|
||||
@@ -791,6 +918,9 @@ export class CatalogController {
|
||||
|
||||
@Get('nimbus-dashboards')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.CATALOG
|
||||
)
|
||||
async getNimbusDashboards(
|
||||
@User() user: RequestUser,
|
||||
@Body() body: GetNimbusDashboardsRequest,
|
||||
@@ -943,4 +1073,4 @@ export class CatalogController {
|
||||
this.logger.error(error.message);
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -29,6 +29,7 @@ import {
|
||||
AssetReporter,
|
||||
BatchRemoveRlsRulesRequest,
|
||||
CreateDataDocsDTO,
|
||||
IUpdateCertificationStatusRequest,
|
||||
IUpdateDataRequest,
|
||||
TriggerCatalogReq,
|
||||
} from './dtos';
|
||||
@@ -119,6 +120,10 @@ class CatalogService implements OnModuleInit {
|
||||
}
|
||||
}
|
||||
|
||||
async getCustomPropertyDefinitions(metadata: Metadata) {
|
||||
return lastValueFrom(this.catalogReadService.GetCustomPropertyDefinitions({}, metadata));
|
||||
}
|
||||
|
||||
async createDataAsset(data: Messages.CreateDataAssetRequest, metadata) {
|
||||
this.logger.info('CatalogService - Manage Data assets permissions');
|
||||
if (!data.embed) data.embed = undefined;
|
||||
@@ -214,15 +219,6 @@ class CatalogService implements OnModuleInit {
|
||||
|
||||
this.logger.debug('Extracted filters:', { filters });
|
||||
|
||||
console.log('MAESTRO VAI CHAMAR PI-FACTORY COM (ANTES AJUSTE):', {
|
||||
search,
|
||||
page,
|
||||
size,
|
||||
sort_by,
|
||||
order,
|
||||
filters,
|
||||
});
|
||||
|
||||
if (
|
||||
filters.manually !== undefined &&
|
||||
filters.manually !== null &&
|
||||
@@ -233,15 +229,6 @@ class CatalogService implements OnModuleInit {
|
||||
delete filters.manually;
|
||||
}
|
||||
|
||||
console.log('MAESTRO VAI CHAMAR PI-FACTORY COM (DEPOIS AJUSTE):', {
|
||||
search,
|
||||
page,
|
||||
size,
|
||||
sort_by,
|
||||
order,
|
||||
filters,
|
||||
});
|
||||
|
||||
if (filters.owner) {
|
||||
const { users: customer_users } =
|
||||
await this.userService.findAllUsersByCustomerId(customer_id);
|
||||
@@ -402,6 +389,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;
|
||||
@@ -517,6 +526,22 @@ 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);
|
||||
@@ -734,6 +759,64 @@ class CatalogService implements OnModuleInit {
|
||||
}
|
||||
}
|
||||
|
||||
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');
|
||||
|
||||
|
||||
@@ -1,4 +1,10 @@
|
||||
import { ApiProperty, ApiPropertyOptional, PickType } from '@nestjs/swagger';
|
||||
import {
|
||||
IsEnum,
|
||||
IsNotEmpty,
|
||||
IsOptional,
|
||||
IsString,
|
||||
} from 'class-validator';
|
||||
import { CreateDataAssetRequest } from '@dadosfera/protospack-v2/dist/lib/Catalog/interfaces/messages';
|
||||
|
||||
export enum DataAssetShareType {
|
||||
@@ -6,6 +12,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',
|
||||
@@ -191,6 +203,27 @@ export class IData {
|
||||
day_opening: number;
|
||||
}
|
||||
|
||||
|
||||
export enum CustomPropertyType {
|
||||
TEXT = 'text',
|
||||
NUMBER = 'number',
|
||||
DATE = 'date',
|
||||
BOOLEAN = 'boolean',
|
||||
}
|
||||
|
||||
export class CustomPropertyDto {
|
||||
@ApiProperty()
|
||||
key: string;
|
||||
@ApiProperty()
|
||||
value: string;
|
||||
@ApiProperty({ enum: CustomPropertyType })
|
||||
type: CustomPropertyType;
|
||||
@ApiPropertyOptional()
|
||||
color?: string;
|
||||
@ApiPropertyOptional()
|
||||
emoji?: string;
|
||||
}
|
||||
|
||||
export class IUpdateDataRequest {
|
||||
@ApiProperty()
|
||||
name: string;
|
||||
@@ -204,7 +237,16 @@ export class IUpdateDataRequest {
|
||||
share_type?: DataAssetShareType;
|
||||
@ApiPropertyOptional()
|
||||
docs?: string;
|
||||
@ApiPropertyOptional({ type: [CustomPropertyDto] })
|
||||
custom_properties?: CustomPropertyDto[];
|
||||
}
|
||||
|
||||
export class IUpdateCertificationStatusRequest {
|
||||
@ApiProperty({ enum: CertificationStatus })
|
||||
@IsEnum(CertificationStatus)
|
||||
certification_status: CertificationStatus;
|
||||
}
|
||||
|
||||
export class ICreateDataAsset implements CreateDataAssetRequest {
|
||||
@ApiProperty()
|
||||
display_name: string;
|
||||
@@ -358,4 +400,4 @@ export type CreateDataDocsDTO = {
|
||||
docs: string;
|
||||
asset_type: string;
|
||||
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -22,17 +22,24 @@ import {
|
||||
ConnectionTestListTablesRes,
|
||||
GetTableMetadataRes,
|
||||
GetTableMetadataReq,
|
||||
RefreshCatalogReq,
|
||||
RefreshCatalogRes,
|
||||
RefreshCatalogStatusReq,
|
||||
} from './dto/connection-test';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import { Authenticated } from 'src/decorators/authentication.decorator';
|
||||
import { Authenticated, RequireModule } from 'src/decorators/authentication.decorator';
|
||||
import { GrpcToHttpExceptionFilter } from 'src/error/grpc-to-http-exception.filter';
|
||||
import { ApiInternalOnlyController } from 'src/decorators/swagger.decorator';
|
||||
import { DADOSFERA_MODULES_KEYS } from 'src/authentication/permissions.enum';
|
||||
|
||||
@ApiInternalOnlyController()
|
||||
@ApiTags('Connection Test')
|
||||
@Controller('connection-test')
|
||||
@UseFilters(new GrpcToHttpExceptionFilter())
|
||||
@Authenticated()
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.COLLECT
|
||||
)
|
||||
export class ConnectionTestController {
|
||||
logger: any;
|
||||
constructor(
|
||||
@@ -84,7 +91,7 @@ export class ConnectionTestController {
|
||||
});
|
||||
return this.connectionTestService.connectionTestListSchemas(
|
||||
body,
|
||||
user.customer_name,
|
||||
user,
|
||||
);
|
||||
}
|
||||
|
||||
@@ -101,7 +108,7 @@ export class ConnectionTestController {
|
||||
});
|
||||
return this.connectionTestService.connectionTestListTables(
|
||||
body,
|
||||
user.customer_name,
|
||||
user,
|
||||
);
|
||||
}
|
||||
|
||||
@@ -118,7 +125,38 @@ export class ConnectionTestController {
|
||||
});
|
||||
return this.connectionTestService.getTableMetadata(
|
||||
body,
|
||||
user.customer_name,
|
||||
user,
|
||||
);
|
||||
}
|
||||
|
||||
@Post('refresh-catalog')
|
||||
@ApiOkResponse({ type: RefreshCatalogRes })
|
||||
@HttpCode(HttpStatus.ACCEPTED)
|
||||
async refreshCatalog(
|
||||
@User() user: RequestUser,
|
||||
@Body(new ValidationPipe()) body: RefreshCatalogReq,
|
||||
) {
|
||||
this.logger.info('/connection-test/refresh-catalog', {
|
||||
user: user.user_id,
|
||||
customer: user.customer_name,
|
||||
connection: body.connection_id,
|
||||
});
|
||||
return this.connectionTestService.refreshCatalog(body, user);
|
||||
}
|
||||
|
||||
@Post('refresh-catalog/status')
|
||||
@ApiOkResponse({ type: RefreshCatalogRes })
|
||||
@HttpCode(HttpStatus.OK)
|
||||
async refreshCatalogStatus(
|
||||
@User() user: RequestUser,
|
||||
@Body(new ValidationPipe()) body: RefreshCatalogStatusReq,
|
||||
) {
|
||||
this.logger.info('/connection-test/refresh-catalog/status', {
|
||||
user: user.user_id,
|
||||
customer: user.customer_name,
|
||||
connection: body.connection_id,
|
||||
session: body.session_id,
|
||||
});
|
||||
return this.connectionTestService.refreshCatalogStatus(body, user);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -5,10 +5,17 @@ import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import { ClientsModule } from '@nestjs/microservices';
|
||||
import { ConnectionTestClientConfiguration } from './connection-test-client.config';
|
||||
import { ConnectionModule } from '../connection/connection.module';
|
||||
import { ConnectionsApiModule } from '../connections-api/connections-api.module';
|
||||
import { PlatformApiModule } from '../platform-api/platform-api.module';
|
||||
const client = new ConnectionTestClientConfiguration();
|
||||
@Module({
|
||||
controllers: [ConnectionTestController],
|
||||
providers: [ConnectionTestService, DadosferaLogger],
|
||||
imports: [ClientsModule.register([client.providerOptions]), ConnectionModule],
|
||||
imports: [
|
||||
ClientsModule.register([client.providerOptions]),
|
||||
ConnectionModule,
|
||||
ConnectionsApiModule,
|
||||
PlatformApiModule,
|
||||
],
|
||||
})
|
||||
export class ConnectionTestModule {}
|
||||
|
||||
@@ -0,0 +1,205 @@
|
||||
import { ConnectionTestService } from './connection-test.service';
|
||||
import { RequestUser } from 'src/decorators/user.decorator';
|
||||
|
||||
describe('ConnectionTestService catalog cache', () => {
|
||||
const user: RequestUser = {
|
||||
user_id: 'user-id',
|
||||
username: 'user@example.com',
|
||||
permissions: [],
|
||||
customer_id: 'customer-id',
|
||||
customer_name: 'customer-name',
|
||||
customer_tier: 'standard',
|
||||
access_token: 'token',
|
||||
customer_modules: [],
|
||||
roles: [],
|
||||
};
|
||||
const grpcClient = { getService: jest.fn().mockReturnValue({}) };
|
||||
const connectionsService = {};
|
||||
const connectionsApiService = { proxy: jest.fn() };
|
||||
const platformApiService = { proxy: jest.fn() };
|
||||
let service: ConnectionTestService;
|
||||
|
||||
beforeEach(() => {
|
||||
jest.clearAllMocks();
|
||||
service = new ConnectionTestService(
|
||||
grpcClient as any,
|
||||
connectionsService as any,
|
||||
connectionsApiService as any,
|
||||
platformApiService as any,
|
||||
);
|
||||
});
|
||||
|
||||
it('keeps the existing schemas response contract', async () => {
|
||||
connectionsApiService.proxy.mockResolvedValue({
|
||||
schemas: [{ schema_name: 'analytics' }, { schema_name: 'public' }],
|
||||
});
|
||||
|
||||
await expect(
|
||||
service.connectionTestListSchemas(
|
||||
{ connection_id: 'config-id', plugin: 'postgresql' },
|
||||
user,
|
||||
),
|
||||
).resolves.toEqual({
|
||||
operation_result: true,
|
||||
schema_list: ['analytics', 'public'],
|
||||
});
|
||||
});
|
||||
|
||||
it('keeps the existing tables response contract', async () => {
|
||||
connectionsApiService.proxy.mockResolvedValue({
|
||||
tables: [{ table_name: 'customers' }, { table_name: 'orders' }],
|
||||
});
|
||||
|
||||
await expect(
|
||||
service.connectionTestListTables(
|
||||
{
|
||||
connection_id: 'config-id',
|
||||
plugin: 'postgresql',
|
||||
schema: 'public',
|
||||
},
|
||||
user,
|
||||
),
|
||||
).resolves.toEqual({
|
||||
operation_result: true,
|
||||
table_list: ['customers', 'orders'],
|
||||
});
|
||||
});
|
||||
|
||||
it('maps cached columns to the existing table metadata contract', async () => {
|
||||
connectionsApiService.proxy.mockResolvedValue({
|
||||
columns: [
|
||||
{
|
||||
column_name: 'id',
|
||||
data_type: 'bigint',
|
||||
is_primary_key: true,
|
||||
},
|
||||
],
|
||||
});
|
||||
|
||||
await expect(
|
||||
service.getTableMetadata(
|
||||
{
|
||||
connection_id: 'config-id',
|
||||
plugin: 'postgresql',
|
||||
schema: 'public',
|
||||
table_list: ['customers'],
|
||||
},
|
||||
user,
|
||||
),
|
||||
).resolves.toEqual({
|
||||
operation_result: true,
|
||||
tables_metadata: [
|
||||
{
|
||||
table_name: 'customers',
|
||||
columns: [
|
||||
{
|
||||
name: 'id',
|
||||
type: 'bigint',
|
||||
is_primary_key: true,
|
||||
},
|
||||
],
|
||||
references: [],
|
||||
},
|
||||
],
|
||||
});
|
||||
expect(connectionsApiService.proxy).toHaveBeenCalledWith(
|
||||
'GET',
|
||||
'/connection_catalog/config-id/schemas/public/tables/customers/columns',
|
||||
user,
|
||||
);
|
||||
});
|
||||
|
||||
it('submits a catalog refresh without holding the request open', async () => {
|
||||
platformApiService.proxy.mockResolvedValue({
|
||||
session_id: 'session-id',
|
||||
date: '20260731',
|
||||
});
|
||||
|
||||
await expect(
|
||||
service.refreshCatalog(
|
||||
{ connection_id: 'config-id', plugin: 'postgresql' },
|
||||
user,
|
||||
),
|
||||
).resolves.toEqual({
|
||||
operation_result: true,
|
||||
status: 'PENDING',
|
||||
session_id: 'session-id',
|
||||
date: '20260731',
|
||||
});
|
||||
|
||||
expect(platformApiService.proxy).toHaveBeenCalledWith(
|
||||
'POST',
|
||||
'/connection_test',
|
||||
user,
|
||||
{
|
||||
customer_id: user.customer_name,
|
||||
plugin: 'postgresql',
|
||||
task: {
|
||||
task_type: 'refresh_catalog',
|
||||
connection: {
|
||||
provider: 'connection_manager',
|
||||
config_id: 'config-id',
|
||||
},
|
||||
},
|
||||
},
|
||||
);
|
||||
});
|
||||
|
||||
it('keeps polling without changing the catalog pointer while pending', async () => {
|
||||
platformApiService.proxy.mockResolvedValue({ status: 'PENDING' });
|
||||
|
||||
await expect(
|
||||
service.refreshCatalogStatus(
|
||||
{
|
||||
connection_id: 'config-id',
|
||||
plugin: 'postgresql',
|
||||
session_id: 'session-id',
|
||||
date: '20260731',
|
||||
},
|
||||
user,
|
||||
),
|
||||
).resolves.toEqual({
|
||||
operation_result: false,
|
||||
status: 'PENDING',
|
||||
session_id: 'session-id',
|
||||
date: '20260731',
|
||||
});
|
||||
|
||||
expect(connectionsApiService.proxy).not.toHaveBeenCalled();
|
||||
});
|
||||
|
||||
it('publishes the catalog pointer after the refresh finishes', async () => {
|
||||
platformApiService.proxy.mockResolvedValue({ status: 'DONE' });
|
||||
connectionsApiService.proxy.mockResolvedValue({
|
||||
last_catalog_refresh_status: 'SUCCESS',
|
||||
});
|
||||
|
||||
await expect(
|
||||
service.refreshCatalogStatus(
|
||||
{
|
||||
connection_id: 'config/id',
|
||||
plugin: 'postgresql',
|
||||
session_id: 'session-id',
|
||||
date: '20260731',
|
||||
},
|
||||
user,
|
||||
),
|
||||
).resolves.toEqual({
|
||||
operation_result: true,
|
||||
status: 'DONE',
|
||||
session_id: 'session-id',
|
||||
date: '20260731',
|
||||
});
|
||||
|
||||
expect(connectionsApiService.proxy).toHaveBeenCalledWith(
|
||||
'PUT',
|
||||
'/connection_config/config%2Fid/catalog_metadata',
|
||||
user,
|
||||
{
|
||||
last_catalog_refresh_status: 'SUCCESS',
|
||||
last_catalog_connection_test_date: '20260731',
|
||||
last_catalog_connection_test_session_id: 'session-id',
|
||||
},
|
||||
);
|
||||
});
|
||||
});
|
||||
@@ -1,4 +1,4 @@
|
||||
import { Inject, Injectable } from '@nestjs/common';
|
||||
import { HttpException, HttpStatus, Inject, Injectable } from '@nestjs/common';
|
||||
import { ClientGrpc } from '@nestjs/microservices';
|
||||
import { ConnectionTest } from '@dadosfera/protospack-v2';
|
||||
import { lastValueFrom } from 'rxjs';
|
||||
@@ -13,6 +13,9 @@ import {
|
||||
ConnectionTestPingRes,
|
||||
GetTableMetadataReq,
|
||||
GetTableMetadataRes,
|
||||
RefreshCatalogReq,
|
||||
RefreshCatalogRes,
|
||||
RefreshCatalogStatusReq,
|
||||
} from './dto/connection-test';
|
||||
import { ConnectionClientService } from '../connection/client.service';
|
||||
import {
|
||||
@@ -21,6 +24,8 @@ import {
|
||||
} from '../connection/dtos/connection';
|
||||
import { RequestUser } from 'src/decorators/user.decorator';
|
||||
import { PackTheMetadata } from 'src/utils/PackTheMetadata';
|
||||
import { ConnectionsApiService } from '../connections-api/connections-api.service';
|
||||
import { PlatformApiService } from '../platform-api/platform-api.service';
|
||||
|
||||
@Injectable()
|
||||
export class ConnectionTestService {
|
||||
@@ -28,6 +33,8 @@ export class ConnectionTestService {
|
||||
constructor(
|
||||
@Inject('ConnectionTestGrpcClient') private readonly grpcClient: ClientGrpc,
|
||||
private connectionsService: ConnectionClientService,
|
||||
private connectionsApiService: ConnectionsApiService,
|
||||
private platformApiService: PlatformApiService,
|
||||
) {
|
||||
this.connectionTestReadClient =
|
||||
grpcClient.getService<ConnectionTest.ReadService.ConnectionTestReadServices>(
|
||||
@@ -147,45 +154,137 @@ export class ConnectionTestService {
|
||||
}
|
||||
async connectionTestListSchemas(
|
||||
body: ConnectionTestListSchemasReq,
|
||||
customer_name: string,
|
||||
user: RequestUser,
|
||||
): Promise<ConnectionTestListSchemasRes> {
|
||||
const { connection_id, plugin } = body;
|
||||
return lastValueFrom(
|
||||
this.connectionTestReadClient.ListSchemas({
|
||||
connection_id,
|
||||
customer_name,
|
||||
plugin,
|
||||
}),
|
||||
const result = await this.connectionsApiService.proxy(
|
||||
'GET',
|
||||
`/connection_catalog/${encodeURIComponent(body.connection_id)}/schemas`,
|
||||
user,
|
||||
);
|
||||
return {
|
||||
operation_result: true,
|
||||
schema_list: result.schemas.map((schema) => schema.schema_name),
|
||||
};
|
||||
}
|
||||
|
||||
async connectionTestListTables(
|
||||
body: ConnectionTestListTablesReq,
|
||||
customer_name: string,
|
||||
user: RequestUser,
|
||||
): Promise<ConnectionTestListTablesRes> {
|
||||
const { connection_id, plugin, schema } = body;
|
||||
return lastValueFrom(
|
||||
this.connectionTestReadClient.ListTables({
|
||||
connection_id,
|
||||
customer_name,
|
||||
plugin,
|
||||
schema,
|
||||
}),
|
||||
const result = await this.connectionsApiService.proxy(
|
||||
'GET',
|
||||
`/connection_catalog/${encodeURIComponent(body.connection_id)}` +
|
||||
`/schemas/${encodeURIComponent(body.schema)}/tables`,
|
||||
user,
|
||||
);
|
||||
return {
|
||||
operation_result: true,
|
||||
table_list: result.tables.map((table) => table.table_name),
|
||||
};
|
||||
}
|
||||
|
||||
async getTableMetadata(
|
||||
body: GetTableMetadataReq,
|
||||
customer_name: string,
|
||||
user: RequestUser,
|
||||
): Promise<GetTableMetadataRes> {
|
||||
const { schema, plugin, table_list, connection_id } = body;
|
||||
return lastValueFrom(
|
||||
this.connectionTestReadClient.GetTableMetadata({
|
||||
connection_id,
|
||||
customer_name,
|
||||
plugin,
|
||||
schema,
|
||||
table_list,
|
||||
const tables_metadata = await Promise.all(
|
||||
body.table_list.map(async (table_name) => {
|
||||
const result = await this.connectionsApiService.proxy(
|
||||
'GET',
|
||||
`/connection_catalog/${encodeURIComponent(body.connection_id)}` +
|
||||
`/schemas/${encodeURIComponent(body.schema)}` +
|
||||
`/tables/${encodeURIComponent(table_name)}/columns`,
|
||||
user,
|
||||
);
|
||||
return {
|
||||
table_name,
|
||||
columns: result.columns.map((column) => ({
|
||||
name: column.column_name,
|
||||
type: column.data_type,
|
||||
is_primary_key: column.is_primary_key,
|
||||
})),
|
||||
references: [],
|
||||
};
|
||||
}),
|
||||
);
|
||||
return { operation_result: true, tables_metadata };
|
||||
}
|
||||
|
||||
async refreshCatalog(
|
||||
body: RefreshCatalogReq,
|
||||
user: RequestUser,
|
||||
): Promise<RefreshCatalogRes> {
|
||||
const task = await this.platformApiService.proxy(
|
||||
'POST',
|
||||
'/connection_test',
|
||||
user,
|
||||
{
|
||||
customer_id: user.customer_name,
|
||||
plugin: body.plugin,
|
||||
task: {
|
||||
task_type: 'refresh_catalog',
|
||||
connection: {
|
||||
provider: 'connection_manager',
|
||||
config_id: body.connection_id,
|
||||
},
|
||||
},
|
||||
},
|
||||
);
|
||||
|
||||
if (!task.session_id || !task.date) {
|
||||
throw new HttpException(
|
||||
'Platform API did not return a catalog refresh task identifier',
|
||||
HttpStatus.BAD_GATEWAY,
|
||||
);
|
||||
}
|
||||
|
||||
return {
|
||||
operation_result: true,
|
||||
status: 'PENDING',
|
||||
session_id: task.session_id,
|
||||
date: task.date,
|
||||
};
|
||||
}
|
||||
|
||||
async refreshCatalogStatus(
|
||||
body: RefreshCatalogStatusReq,
|
||||
user: RequestUser,
|
||||
): Promise<RefreshCatalogRes> {
|
||||
const result = await this.platformApiService.proxy(
|
||||
'POST',
|
||||
'/connection_test/status',
|
||||
user,
|
||||
{
|
||||
session_id: body.session_id,
|
||||
date: body.date,
|
||||
},
|
||||
);
|
||||
|
||||
if (result.status === 'DONE') {
|
||||
await this.connectionsApiService.proxy(
|
||||
'PUT',
|
||||
`/connection_config/${encodeURIComponent(
|
||||
body.connection_id,
|
||||
)}/catalog_metadata`,
|
||||
user,
|
||||
{
|
||||
last_catalog_refresh_status: 'SUCCESS',
|
||||
last_catalog_connection_test_date: body.date,
|
||||
last_catalog_connection_test_session_id: body.session_id,
|
||||
},
|
||||
);
|
||||
} else if (result.status === 'ERROR' || result.status === 'EXPIRED') {
|
||||
throw new HttpException(
|
||||
`Catalog refresh finished with status ${result.status}`,
|
||||
HttpStatus.BAD_GATEWAY,
|
||||
);
|
||||
}
|
||||
|
||||
return {
|
||||
operation_result: result.status === 'DONE',
|
||||
status: result.status,
|
||||
session_id: body.session_id,
|
||||
date: body.date,
|
||||
};
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,5 +1,5 @@
|
||||
import { ApiProperty, ApiPropertyOptional, OmitType } from '@nestjs/swagger';
|
||||
import { IsString, IsOptional } from 'class-validator';
|
||||
import { IsIn, IsString, IsOptional } from 'class-validator';
|
||||
import { DatabaseConnectionPropertiesDto } from 'src/modules/connection/dtos/connection';
|
||||
import { CreateConnectionDto } from 'src/modules/connection/dtos/connection';
|
||||
export class ColumnDto {
|
||||
@@ -7,6 +7,8 @@ export class ColumnDto {
|
||||
name: string;
|
||||
@ApiProperty()
|
||||
type: string;
|
||||
@ApiProperty()
|
||||
is_primary_key: boolean;
|
||||
}
|
||||
export class TableMetadataDto {
|
||||
@ApiProperty()
|
||||
@@ -131,3 +133,37 @@ export class GetTableMetadataRes {
|
||||
@ApiProperty({ type: [TableMetadataDto] })
|
||||
tables_metadata: TableMetadataDto[];
|
||||
}
|
||||
|
||||
export class RefreshCatalogReq {
|
||||
@ApiProperty()
|
||||
@IsString()
|
||||
connection_id: string;
|
||||
|
||||
@ApiProperty({ enum: ['oracle', 'mysql', 'postgresql', 'sqlserver'] })
|
||||
@IsIn(['oracle', 'mysql', 'postgresql', 'sqlserver'])
|
||||
plugin: string;
|
||||
}
|
||||
|
||||
export class RefreshCatalogStatusReq extends RefreshCatalogReq {
|
||||
@ApiProperty()
|
||||
@IsString()
|
||||
session_id: string;
|
||||
|
||||
@ApiProperty()
|
||||
@IsString()
|
||||
date: string;
|
||||
}
|
||||
|
||||
export class RefreshCatalogRes {
|
||||
@ApiProperty()
|
||||
operation_result: boolean;
|
||||
|
||||
@ApiProperty()
|
||||
status: string;
|
||||
|
||||
@ApiProperty()
|
||||
session_id: string;
|
||||
|
||||
@ApiProperty()
|
||||
date: string;
|
||||
}
|
||||
|
||||
@@ -16,8 +16,9 @@ import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import {
|
||||
Authenticated,
|
||||
RequireAllPermissions,
|
||||
RequireModule,
|
||||
} from 'src/decorators/authentication.decorator';
|
||||
import { PERMISSIONS_GROUPS } from 'src/authentication/permissions.enum';
|
||||
import { DADOSFERA_MODULES_KEYS, PERMISSIONS_GROUPS } from 'src/authentication/permissions.enum';
|
||||
import { RequestUser, User } from 'src/decorators/user.decorator';
|
||||
import { ValidationPipe } from '../../pipes/object-validation.pipe';
|
||||
import {
|
||||
@@ -39,6 +40,9 @@ const connectionPermissions = PERMISSIONS_GROUPS.CONNECTION.permissions;
|
||||
@ApiTags('connections')
|
||||
@Authenticated()
|
||||
@Controller('connections')
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.COLLECT
|
||||
)
|
||||
export class ConnectionController {
|
||||
logger: any;
|
||||
constructor(
|
||||
|
||||
@@ -0,0 +1,11 @@
|
||||
export const CONNECTIONS_API_CONFIG = {
|
||||
getUrl: (): string => {
|
||||
const url = process.env.CONNECTIONS_API_URL;
|
||||
if (!url) {
|
||||
throw new Error('CONNECTIONS_API_URL environment variable is not set');
|
||||
}
|
||||
return url;
|
||||
},
|
||||
region: process.env.AWS_REGION || 'us-east-1',
|
||||
timeout: parseInt(process.env.CONNECTIONS_API_TIMEOUT || '30000', 10),
|
||||
};
|
||||
@@ -0,0 +1,10 @@
|
||||
import { Module } from '@nestjs/common';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
|
||||
import { ConnectionsApiService } from './connections-api.service';
|
||||
|
||||
@Module({
|
||||
providers: [ConnectionsApiService, DadosferaLogger],
|
||||
exports: [ConnectionsApiService],
|
||||
})
|
||||
export class ConnectionsApiModule {}
|
||||
@@ -0,0 +1,99 @@
|
||||
import { Injectable, Inject, HttpException } from '@nestjs/common';
|
||||
import { SignatureV4 } from '@aws-sdk/signature-v4';
|
||||
import { Sha256 } from '@aws-crypto/sha256-js';
|
||||
import { defaultProvider } from '@aws-sdk/credential-provider-node';
|
||||
import axios, { AxiosResponse, Method } from 'axios';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
|
||||
import { RequestUser } from '../../decorators/user.decorator';
|
||||
import { CONNECTIONS_API_CONFIG } from './connections-api.config';
|
||||
|
||||
@Injectable()
|
||||
export class ConnectionsApiService {
|
||||
private signer: SignatureV4;
|
||||
private logger: any;
|
||||
|
||||
constructor(@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
this.signer = new SignatureV4({
|
||||
service: 'execute-api',
|
||||
region: CONNECTIONS_API_CONFIG.region,
|
||||
credentials: defaultProvider(),
|
||||
sha256: Sha256,
|
||||
});
|
||||
}
|
||||
|
||||
async proxy(
|
||||
method: string,
|
||||
path: string,
|
||||
user: RequestUser,
|
||||
body?: any,
|
||||
query?: Record<string, string>,
|
||||
): Promise<any> {
|
||||
const baseUrl = CONNECTIONS_API_CONFIG.getUrl();
|
||||
const url = new URL(`${baseUrl}${path}`);
|
||||
|
||||
if (query) {
|
||||
Object.entries(query).forEach(([key, value]) => {
|
||||
if (value !== undefined && value !== null) {
|
||||
url.searchParams.set(key, String(value));
|
||||
}
|
||||
});
|
||||
}
|
||||
const headers: Record<string, string> = {
|
||||
host: url.hostname,
|
||||
'content-type': 'application/json',
|
||||
customer_name: user.customer_name || '',
|
||||
customer_id: user.customer_id || '',
|
||||
'x-user-id': user.user_id || '',
|
||||
'x-username': user.username || '',
|
||||
'x-customer-tier': user.customer_tier || '',
|
||||
'x-customer-id': user.customer_id || '',
|
||||
};
|
||||
const requestToSign = {
|
||||
method: method.toUpperCase(),
|
||||
protocol: url.protocol,
|
||||
hostname: url.hostname,
|
||||
port: url.port ? parseInt(url.port, 10) : undefined,
|
||||
path: url.pathname + url.search,
|
||||
headers,
|
||||
body: body ? JSON.stringify(body) : undefined,
|
||||
};
|
||||
|
||||
try {
|
||||
const signedRequest = await this.signer.sign(requestToSign);
|
||||
const response: AxiosResponse = await axios({
|
||||
method: method as Method,
|
||||
url: url.href,
|
||||
headers: signedRequest.headers as Record<string, string>,
|
||||
data: body,
|
||||
timeout: CONNECTIONS_API_CONFIG.timeout,
|
||||
validateStatus: () => true,
|
||||
});
|
||||
|
||||
if (response.status >= 400) {
|
||||
throw new HttpException(response.data, response.status);
|
||||
}
|
||||
return response.data;
|
||||
} catch (error) {
|
||||
this.logger.error('Connections API proxy error', {
|
||||
error: error.message,
|
||||
path,
|
||||
method: method.toUpperCase(),
|
||||
});
|
||||
if (error instanceof HttpException) {
|
||||
throw error;
|
||||
}
|
||||
if (error.response) {
|
||||
throw new HttpException(error.response.data, error.response.status);
|
||||
}
|
||||
if (error.code === 'ECONNREFUSED') {
|
||||
throw new HttpException('Connections API service unavailable', 503);
|
||||
}
|
||||
if (error.code === 'ETIMEDOUT' || error.code === 'ECONNABORTED') {
|
||||
throw new HttpException('Connections API request timeout', 504);
|
||||
}
|
||||
throw new HttpException('Internal server error', 500);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -25,9 +25,10 @@ import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import {
|
||||
Authenticated,
|
||||
RequireAllPermissions,
|
||||
RequireModule,
|
||||
RequireSomePermission,
|
||||
} from 'src/decorators/authentication.decorator';
|
||||
import { PERMISSIONS_GROUPS } from 'src/authentication/permissions.enum';
|
||||
import { DADOSFERA_MODULES_KEYS, PERMISSIONS_GROUPS } from 'src/authentication/permissions.enum';
|
||||
import { Language } from 'src/decorators/language.decorator';
|
||||
import { LanguageEnum } from 'src/utils/languages.enum';
|
||||
import { ApiInternalOnlyController } from 'src/decorators/swagger.decorator';
|
||||
@@ -99,6 +100,9 @@ export class ConnectorController {
|
||||
PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE,
|
||||
PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.COLLECT
|
||||
)
|
||||
async getAllConnectors(
|
||||
@Language() language: LanguageEnum,
|
||||
@Query() queries: GetAllDto,
|
||||
@@ -131,6 +135,9 @@ export class ConnectorController {
|
||||
PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE,
|
||||
PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.COLLECT
|
||||
)
|
||||
async getConnectorsTags() {
|
||||
return await this.connectorClientService.getConnectorsTags();
|
||||
}
|
||||
@@ -143,6 +150,9 @@ export class ConnectorController {
|
||||
PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE,
|
||||
PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.COLLECT
|
||||
)
|
||||
async getConnector(
|
||||
@Language() language: LanguageEnum,
|
||||
@Param('plugin') plugin: string,
|
||||
@@ -171,6 +181,9 @@ export class ConnectorController {
|
||||
PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE,
|
||||
PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE,
|
||||
)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.COLLECT
|
||||
)
|
||||
async getConnectorDetails(
|
||||
@Language() language: LanguageEnum,
|
||||
@Param('plugin') plugin: string,
|
||||
@@ -193,6 +206,9 @@ export class ConnectorController {
|
||||
@Put('/:plugin')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.CONNECTORS.permissions.UPDATE)
|
||||
@ApiConsumes('multipart/form-data')
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.COLLECT
|
||||
)
|
||||
async updateConnector(
|
||||
@Param('plugin') plugin: string,
|
||||
@Body() body: UpdateDto,
|
||||
@@ -214,6 +230,9 @@ export class ConnectorController {
|
||||
|
||||
@Put('/:plugin/add-tag')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.CONNECTORS.permissions.UPDATE)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.COLLECT
|
||||
)
|
||||
async addTagOnConnector(
|
||||
@Param('plugin') plugin: string,
|
||||
@Body() body: AddTagDto,
|
||||
@@ -241,6 +260,9 @@ export class ConnectorController {
|
||||
|
||||
@Put('/:plugin/remove-tag')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.CONNECTORS.permissions.UPDATE)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.COLLECT
|
||||
)
|
||||
async removeTagOnConnector(
|
||||
@Param('plugin') plugin: string,
|
||||
@Body() body: RemoveTagDto,
|
||||
@@ -269,6 +291,9 @@ export class ConnectorController {
|
||||
|
||||
@Delete('/:plugin')
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.CONNECTORS.permissions.DELETE)
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.COLLECT
|
||||
)
|
||||
async deleteConnector(
|
||||
@Param('plugin') plugin: string,
|
||||
@Query('version') version: string,
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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;
|
||||
}
|
||||
@@ -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,10 +17,13 @@ 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';
|
||||
import { CreateInputReq } from './dtos/input.model';
|
||||
import { Metadata } from '@grpc/grpc-js';
|
||||
|
||||
@Injectable()
|
||||
|
||||
@@ -71,10 +74,10 @@ export class InputsService {
|
||||
objectCamelToSnake(createInputResponse);
|
||||
return createInputResponse;
|
||||
},
|
||||
update: async (updateInputDTO: UpdateInputRequest) => {
|
||||
this.logger.info('InputClientService - Update');
|
||||
update: async (updateInputDTO: UpdateInputRequest, metadata: Metadata): Promise<InputUpdateResponse> => {
|
||||
this.logger.info('InputClientService - Update' + JSON.stringify(updateInputDTO));
|
||||
const updateInputResponse = await lastValueFrom(
|
||||
this.inputWriteService.InputUpdate(updateInputDTO),
|
||||
this.inputWriteService.InputUpdate(updateInputDTO, metadata),
|
||||
);
|
||||
|
||||
return updateInputResponse;
|
||||
@@ -164,6 +167,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,
|
||||
};
|
||||
@@ -199,24 +207,47 @@ export class InputsService {
|
||||
return findOneInputResponse;
|
||||
}
|
||||
|
||||
async update(id: string, data, info: Info) {
|
||||
this.validateCron({ ...data, info });
|
||||
async update(id: string, data, info: Info, metadata?: Metadata) {
|
||||
// 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,
|
||||
}, metadata);
|
||||
|
||||
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 +289,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,78 +0,0 @@
|
||||
import { ConflictException, Inject, OnModuleInit } from '@nestjs/common';
|
||||
import { ClientGrpc } from '@nestjs/microservices';
|
||||
import {
|
||||
PipelineServicesNames,
|
||||
PipelinesServiceInterface,
|
||||
} from '@dadosfera/protospack';
|
||||
import { lastValueFrom } from 'rxjs';
|
||||
|
||||
import { IIdRequest } from './interfaces';
|
||||
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import { PipelinesClientConfiguration } from './pipelines-client';
|
||||
|
||||
export class PipelinesClientService implements OnModuleInit {
|
||||
private pipelineService: PipelinesServiceInterface;
|
||||
logger: DadosferaLogger;
|
||||
|
||||
constructor(
|
||||
@Inject(DadosferaLogger)
|
||||
dadosferaLogger: DadosferaLogger,
|
||||
@Inject(PipelinesClientConfiguration.name)
|
||||
private readonly grpcClient: ClientGrpc,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
}
|
||||
|
||||
onModuleInit() {
|
||||
this.pipelineService =
|
||||
this.grpcClient.getService<PipelinesServiceInterface>(
|
||||
PipelineServicesNames.PipelineService,
|
||||
);
|
||||
}
|
||||
|
||||
async getPipelineStatus(data) {
|
||||
this.logger.info('PipelinesClientService - GetPipelineStatus');
|
||||
|
||||
const statusPipelineResponse = await lastValueFrom(
|
||||
this.pipelineService.getPipelineStatus(data),
|
||||
)
|
||||
.then((res) => {
|
||||
const statusArray =
|
||||
res.status?.sort((a, b) => {
|
||||
if (a.id < b.id) {
|
||||
return 1;
|
||||
} else {
|
||||
return -1;
|
||||
}
|
||||
}) || [];
|
||||
return { status: statusArray };
|
||||
})
|
||||
.catch((err) => {
|
||||
this.logger.error(err.message);
|
||||
throw new Error(err);
|
||||
});
|
||||
this.logger.info('Done');
|
||||
|
||||
return statusPipelineResponse;
|
||||
}
|
||||
|
||||
async runPipeline({ id, info }: IIdRequest) {
|
||||
this.logger.info('PipelinesClientService - RunPipeline');
|
||||
const statusPipelineResponse = await lastValueFrom(
|
||||
this.pipelineService.triggerPipeline({ id, info }),
|
||||
).catch((err) => {
|
||||
this.logger.error(err.message);
|
||||
throw new Error(err);
|
||||
});
|
||||
|
||||
if (statusPipelineResponse.status == false) {
|
||||
throw new ConflictException(
|
||||
'This pipeline is not ready yet to execute, Try again later!',
|
||||
);
|
||||
}
|
||||
|
||||
this.logger.info('Done');
|
||||
return statusPipelineResponse;
|
||||
}
|
||||
}
|
||||
-36
@@ -1,36 +0,0 @@
|
||||
import { Info } from '@dadosfera/protospack/dist/lib/interfaces';
|
||||
|
||||
export interface ICreatePipelineDto {
|
||||
input: IdRequest;
|
||||
transformations: IdRequest[];
|
||||
output: IdRequest;
|
||||
tags: string[];
|
||||
name: string;
|
||||
description: string;
|
||||
info: Info;
|
||||
}
|
||||
|
||||
export interface IdRequest {
|
||||
id: string;
|
||||
}
|
||||
|
||||
export interface IIdRequest {
|
||||
id: string;
|
||||
info: Info;
|
||||
}
|
||||
|
||||
export interface IUpdatePipelineRequest {
|
||||
input: IdRequest;
|
||||
transformations: IdRequest[];
|
||||
output: IdRequest;
|
||||
tags: string[];
|
||||
name: string;
|
||||
description: string;
|
||||
id: string;
|
||||
info: Info;
|
||||
}
|
||||
|
||||
export interface IGetPipelineLogsRequest {
|
||||
id: string;
|
||||
details: string;
|
||||
}
|
||||
@@ -1,33 +0,0 @@
|
||||
import {
|
||||
ClientsProviderAsyncOptions,
|
||||
GrpcOptions,
|
||||
Transport,
|
||||
} from '@nestjs/microservices';
|
||||
import { PipelinePackages, PipelineProtoFilePath } from '@dadosfera/protospack';
|
||||
import { credentials } from '@grpc/grpc-js';
|
||||
|
||||
const isLocalConnection =
|
||||
process.env.PIFACTORY_URL.startsWith('pi-factory:') ||
|
||||
process.env.PIFACTORY_URL.includes('0.0.0.0');
|
||||
|
||||
export class PipelinesClientConfiguration {
|
||||
public name = 'PipelinesClientConfiguration';
|
||||
private config: GrpcOptions = {
|
||||
transport: Transport.GRPC,
|
||||
options: {
|
||||
url: process.env.PIFACTORY_URL,
|
||||
package: PipelinePackages,
|
||||
credentials: isLocalConnection ? undefined : credentials.createSsl(),
|
||||
protoPath: PipelineProtoFilePath,
|
||||
loader: {
|
||||
keepCase: true,
|
||||
enums: String,
|
||||
defaults: false,
|
||||
},
|
||||
},
|
||||
};
|
||||
providerOptions: ClientsProviderAsyncOptions = {
|
||||
name: this.name,
|
||||
...this.config,
|
||||
};
|
||||
}
|
||||
@@ -1,72 +0,0 @@
|
||||
import { Body, Controller, Get, Inject, Param, Post } from '@nestjs/common';
|
||||
import { ApiOperation, ApiTags } from '@nestjs/swagger';
|
||||
import {
|
||||
AuthenticateCondition,
|
||||
Authenticated,
|
||||
RequireSomePermission,
|
||||
} from 'src/decorators/authentication.decorator';
|
||||
import { PERMISSIONS_GROUPS } from '../../authentication/permissions.enum';
|
||||
import { PipelinesService } from './pipelines.service';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import { ApiInternalOnlyController } from 'src/decorators/swagger.decorator';
|
||||
|
||||
@ApiInternalOnlyController()
|
||||
@ApiTags('Pipelines')
|
||||
@Controller('pipelines')
|
||||
@Authenticated()
|
||||
export class PipelinesController {
|
||||
logger: DadosferaLogger;
|
||||
constructor(
|
||||
@Inject(DadosferaLogger)
|
||||
dadosferaLogger: DadosferaLogger,
|
||||
private pipelineService: PipelinesService,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
}
|
||||
|
||||
@Post('start/:id')
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.CREATE)
|
||||
@ApiOperation({
|
||||
deprecated: true,
|
||||
description:
|
||||
'This method is deprecated. Please use route /pipelinesV2/start/:id instead',
|
||||
})
|
||||
async activate(@Param('id') id: string, @Body() body) {
|
||||
const { info } = body;
|
||||
|
||||
this.logger.info(
|
||||
process.env.DEV_URL + `/pipeline/start/${id} - ON START PIPELINE ROUTE`,
|
||||
{
|
||||
user: body.info.user_id,
|
||||
customer: body.info.customer,
|
||||
},
|
||||
);
|
||||
|
||||
const response = await this.pipelineService.runPipeline({ id, info });
|
||||
|
||||
return response;
|
||||
}
|
||||
|
||||
@Get(':id/status')
|
||||
@ApiOperation({
|
||||
deprecated: true,
|
||||
description:
|
||||
'This method is deprecated. Please use route /pipelinesV2/:id/status instead',
|
||||
})
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.IMPORT_FILES.permissions.VIEW, PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getPipelineStatus(@Body() body, @Param('id') id: string) {
|
||||
body.id = id;
|
||||
|
||||
this.logger.info(
|
||||
process.env.DEV_URL + `/pipeline/${id} - ON GET PIPELINE STATUS ROUTE`,
|
||||
{
|
||||
user: body.info.user_id,
|
||||
customer: body.info.customer,
|
||||
},
|
||||
);
|
||||
|
||||
const response = await this.pipelineService.getPipelineStatus(body);
|
||||
|
||||
return response;
|
||||
}
|
||||
}
|
||||
@@ -1,19 +0,0 @@
|
||||
import { Module } from '@nestjs/common';
|
||||
import { ClientsModule } from '@nestjs/microservices';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
|
||||
import { PipelinesController } from './pipelines.controller';
|
||||
import { PipelinesService } from './pipelines.service';
|
||||
|
||||
import { PipelinesClientConfiguration } from './pipelines-client';
|
||||
import { PipelinesClientService } from './client.service';
|
||||
|
||||
const client = new PipelinesClientConfiguration();
|
||||
|
||||
@Module({
|
||||
imports: [ClientsModule.register([client.providerOptions])],
|
||||
controllers: [PipelinesController],
|
||||
providers: [PipelinesService, PipelinesClientService, DadosferaLogger],
|
||||
exports: [PipelinesService],
|
||||
})
|
||||
export class PipelinesModule {}
|
||||
@@ -1,33 +0,0 @@
|
||||
import { HttpException, HttpStatus, Injectable } from '@nestjs/common';
|
||||
import { PipelinesClientService } from './client.service';
|
||||
import { IIdRequest } from './interfaces';
|
||||
import { objectCamelToSnake } from 'src/utils/CaseConverter';
|
||||
|
||||
@Injectable()
|
||||
export class PipelinesService {
|
||||
constructor(private pipelineClient: PipelinesClientService) {}
|
||||
|
||||
async getPipelineStatus(data: IIdRequest) {
|
||||
try {
|
||||
const pipelineStatusResponse =
|
||||
await this.pipelineClient.getPipelineStatus(data);
|
||||
|
||||
return objectCamelToSnake(pipelineStatusResponse);
|
||||
} catch (err) {
|
||||
throw new HttpException(err.message, HttpStatus.NOT_FOUND);
|
||||
}
|
||||
}
|
||||
|
||||
async runPipeline({ id, info }: IIdRequest) {
|
||||
try {
|
||||
const triggerPipelineResponse = await this.pipelineClient.runPipeline({
|
||||
id,
|
||||
info,
|
||||
});
|
||||
|
||||
return objectCamelToSnake(triggerPipelineResponse);
|
||||
} catch (err) {
|
||||
throw new HttpException(err.message, HttpStatus.NOT_FOUND);
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -1,5 +1,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,18 +24,17 @@ import {
|
||||
ApiTags,
|
||||
} from '@nestjs/swagger';
|
||||
import {
|
||||
AuthenticateCondition,
|
||||
RequireAllPermissions,
|
||||
RequireModule,
|
||||
RequireSomePermission,
|
||||
} from 'src/decorators/authentication.decorator';
|
||||
import { PERMISSIONS_GROUPS } from '../../authentication/permissions.enum';
|
||||
import { DADOSFERA_MODULES_KEYS, PERMISSIONS_GROUPS } from '../../authentication/permissions.enum';
|
||||
import { PipelinesService } from './pipelines.service';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
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,
|
||||
@@ -43,23 +42,31 @@ 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')
|
||||
@RequireModule(
|
||||
DADOSFERA_MODULES_KEYS.COLLECT
|
||||
)
|
||||
export class PipelinesController {
|
||||
logger: DadosferaLogger;
|
||||
constructor(
|
||||
@Inject(DadosferaLogger)
|
||||
dadosferaLogger: DadosferaLogger,
|
||||
private pipelinesClientService: PipelinesService,
|
||||
private oldPipelinesService: OldPipelineService,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
}
|
||||
@@ -187,6 +194,7 @@ export class PipelinesController {
|
||||
@Get(':id/status')
|
||||
@RequireSomePermission(PERMISSIONS_GROUPS.IMPORT_FILES.permissions.VIEW, PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getPipelineStatus(@Body() body, @Param('id') id: string) {
|
||||
|
||||
body.id = id;
|
||||
|
||||
this.logger.info(`/pipeline/${id} - ON GET PIPELINE STATUS ROUTE`, {
|
||||
@@ -194,7 +202,7 @@ export class PipelinesController {
|
||||
customer: body.info.customer,
|
||||
});
|
||||
|
||||
const response = await this.oldPipelinesService.getPipelineStatus(body);
|
||||
const response = await this.pipelinesClientService.getPipelineStatus(body);
|
||||
|
||||
return response;
|
||||
}
|
||||
@@ -218,26 +226,50 @@ 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;
|
||||
}
|
||||
@@ -277,6 +309,40 @@ 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 info: Info = {
|
||||
user_id: user.user_id,
|
||||
customer: user.customer_name,
|
||||
customer_id: user.customer_id,
|
||||
pipeline_id: pipelineId
|
||||
};
|
||||
|
||||
const metadata = PackTheMetadata(user);
|
||||
|
||||
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({
|
||||
@@ -296,6 +362,21 @@ 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)
|
||||
@@ -425,7 +506,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,320 @@ 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,
|
||||
metadata
|
||||
);
|
||||
|
||||
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 || [];
|
||||
|
||||
if (user.customer_modules.includes('catalog')) {
|
||||
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;
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -10,21 +10,30 @@ import {
|
||||
Query,
|
||||
Inject,
|
||||
BadRequestException,
|
||||
HttpException,
|
||||
NotFoundException,
|
||||
UseGuards,
|
||||
} from '@nestjs/common';
|
||||
import { ApiTags, ApiOperation } from '@nestjs/swagger';
|
||||
import { ApiTags, ApiOperation, ApiOkResponse } from '@nestjs/swagger';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
|
||||
import {
|
||||
Authenticated,
|
||||
RequireAllPermissions,
|
||||
RequireModule,
|
||||
} from '../../decorators/authentication.decorator';
|
||||
import { User, RequestUser } from '../../decorators/user.decorator';
|
||||
import { PlatformApiService } from './platform-api.service';
|
||||
import { PERMISSIONS_GROUPS } from '../../authentication/permissions.enum';
|
||||
import { DADOSFERA_MODULES_KEYS, PERMISSIONS_GROUPS } from '../../authentication/permissions.enum';
|
||||
import { ElasticsearchService } from '../../services/elasticsearch';
|
||||
import { DynamoDBService, ReferenceColumn } from '../../services/dynamodb';
|
||||
import { CustomersService } from '../customers/customers.service';
|
||||
import { validateCronAgainstScheduleLimit } from '../../utils/cron-validation';
|
||||
import { CatalogService } from '../catalog/catalog.service';
|
||||
import { PackTheMetadata } from '../../utils/PackTheMetadata';
|
||||
import { ValidationTableDTO } from './platform-api.dto';
|
||||
import { InputsService } from '../inputs/inputs.service';
|
||||
import { PipelineExecutionGuard } from 'src/guards/pipeline-execution.guard';
|
||||
|
||||
|
||||
type ValidateTablesDTO = {
|
||||
@@ -34,8 +43,14 @@ type ValidateTablesDTO = {
|
||||
}>
|
||||
}
|
||||
|
||||
type RenameTablesBody = {
|
||||
raw?: { table_name: string; table_schema: string };
|
||||
qualify?: { table_name: string; table_schema: string };
|
||||
}
|
||||
|
||||
@ApiTags('Platform API')
|
||||
@Controller('platform')
|
||||
@RequireModule(DADOSFERA_MODULES_KEYS.COLLECT)
|
||||
export class PlatformApiController {
|
||||
private logger: any;
|
||||
|
||||
@@ -44,6 +59,8 @@ export class PlatformApiController {
|
||||
private readonly elasticsearchService: ElasticsearchService,
|
||||
private readonly dynamoDBService: DynamoDBService,
|
||||
private readonly customersService: CustomersService,
|
||||
private readonly catalogService: CatalogService,
|
||||
private readonly inputsService: InputsService,
|
||||
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
@@ -57,6 +74,10 @@ export class PlatformApiController {
|
||||
return id?.replace(/-/g, '_') || '';
|
||||
}
|
||||
|
||||
private decodePathParam(value: string): string {
|
||||
return value ? decodeURIComponent(value) : '';
|
||||
}
|
||||
|
||||
/**
|
||||
* Denormalize ID back to UUID format (replace _ with -).
|
||||
* Used when we receive a normalized ID but need the original UUID.
|
||||
@@ -75,6 +96,23 @@ export class PlatformApiController {
|
||||
return jobId?.replace(/-/g, '_') || '';
|
||||
}
|
||||
|
||||
private async getJobByAnyConnectorType(normalizedJobId: string, user: RequestUser): Promise<any> {
|
||||
const connectorTypes = ['jdbc', 'singer', 's3'];
|
||||
for (const type of connectorTypes) {
|
||||
try {
|
||||
const job = await this.platformApiService.proxy(
|
||||
'GET',
|
||||
`/jobs/${type}/${normalizedJobId}`,
|
||||
user,
|
||||
);
|
||||
return job;
|
||||
} catch (error) {
|
||||
// Continue to next connector type
|
||||
}
|
||||
}
|
||||
throw new HttpException(`Job ${normalizedJobId} not found in any connector type (jdbc, singer, s3)`, 404);
|
||||
}
|
||||
|
||||
/**
|
||||
* Extract the pipeline ID (base UUID) from a job ID.
|
||||
* Job IDs have format "uuid-suffix" where suffix is the job index (e.g., "0", "1").
|
||||
@@ -102,7 +140,7 @@ export class PlatformApiController {
|
||||
}
|
||||
|
||||
private readonly VALID_CONNECTORS = ['jdbc', 'singer', 's3'];
|
||||
private readonly MAX_MEMORY_MB = 12000; // 12GB maximum memory per pipeline/job
|
||||
private readonly MAX_MEMORY_MB = 12000; // 12GB maximum memory per pipelines/job
|
||||
|
||||
/**
|
||||
* Validate that connector is provided and is a valid type.
|
||||
@@ -373,55 +411,9 @@ export class PlatformApiController {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Sync sync-mode changes to DynamoDB for JDBC connectors.
|
||||
* Always passes both target_load_type and incremental_column_name to ensure proper sync.
|
||||
*/
|
||||
private async syncJdbcSyncModeToDynamoDB(
|
||||
jobId: string,
|
||||
body: any,
|
||||
user: RequestUser,
|
||||
): Promise<void> {
|
||||
// JDBC sync mode uses target_load_type field
|
||||
const changes: any = {};
|
||||
|
||||
if ('target_load_type' in body) {
|
||||
changes.target_load_type = body.target_load_type;
|
||||
}
|
||||
|
||||
// Handle incremental_column_name:
|
||||
// - If provided in body, use that value
|
||||
// - If changing to full_load, explicitly clear it
|
||||
if ('incremental_column_name' in body) {
|
||||
changes.incremental_column_name = body.incremental_column_name;
|
||||
changes.incremental_column_type = body.incremental_column_type;
|
||||
} else if (body.target_load_type === 'full_load') {
|
||||
// Changing to full_load without specifying incremental_column - clear it
|
||||
changes.incremental_column_name = null;
|
||||
}
|
||||
|
||||
await this.syncJobInputToDynamoDB(jobId, changes, user, 'jdbc');
|
||||
}
|
||||
|
||||
/**
|
||||
* Sync sync-mode changes to DynamoDB for Singer connectors.
|
||||
*/
|
||||
private async syncSingerSyncModeToDynamoDB(
|
||||
jobId: string,
|
||||
body: any,
|
||||
user: RequestUser,
|
||||
): Promise<void> {
|
||||
// Singer sync mode uses replication_method field
|
||||
// Map to DynamoDB type: FULL_TABLE -> full_load, INCREMENTAL -> incremental
|
||||
if ('replication_method' in body) {
|
||||
const type = body.replication_method === 'INCREMENTAL' ? 'incremental' : 'full_load';
|
||||
await this.syncJobInputToDynamoDB(jobId, { load_type: type }, user, 'singer');
|
||||
}
|
||||
}
|
||||
|
||||
// ==================== PIPELINE ROUTES ====================
|
||||
|
||||
@Post('pipeline')
|
||||
@Post('pipelines')
|
||||
@ApiOperation({ summary: 'Create a new pipeline' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.CREATE)
|
||||
async createPipeline(@Body() body: any, @User() user: RequestUser) {
|
||||
@@ -542,7 +534,7 @@ export class PlatformApiController {
|
||||
);
|
||||
}
|
||||
|
||||
@Get('pipeline/:pipelineId')
|
||||
@Get('pipelines/:pipelineId')
|
||||
@ApiOperation({ summary: 'Get pipeline by ID' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getPipeline(
|
||||
@@ -553,7 +545,7 @@ export class PlatformApiController {
|
||||
return this.platformApiService.proxy('GET', `/pipeline/${normalizedId}`, user);
|
||||
}
|
||||
|
||||
@Patch('pipeline/:pipelineId')
|
||||
@Patch('pipelines/:pipelineId')
|
||||
@ApiOperation({ summary: 'Update pipeline by ID' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||
async updatePipeline(
|
||||
@@ -604,7 +596,7 @@ export class PlatformApiController {
|
||||
return result;
|
||||
}
|
||||
|
||||
@Delete('pipeline/:pipelineId')
|
||||
@Delete('pipelines/:pipelineId')
|
||||
@ApiOperation({ summary: 'Delete pipeline by ID' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE)
|
||||
async deletePipeline(
|
||||
@@ -634,7 +626,7 @@ export class PlatformApiController {
|
||||
return result;
|
||||
}
|
||||
|
||||
@Post('pipeline/execute')
|
||||
@Post('pipelines/execute')
|
||||
@ApiOperation({ summary: 'Execute a pipeline' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||
async executePipeline(@Body() body: any, @User() user: RequestUser) {
|
||||
@@ -644,10 +636,10 @@ export class PlatformApiController {
|
||||
...body,
|
||||
customer_id: user.customer_name,
|
||||
};
|
||||
return this.platformApiService.proxy('POST', '/pipeline/execute', user, enrichedBody);
|
||||
return this.platformApiService.proxy('POST', '/pipelines/execute', user, enrichedBody);
|
||||
}
|
||||
|
||||
@Post('pipeline/pause')
|
||||
@Post('pipelines/pause')
|
||||
@ApiOperation({ summary: 'Pause a pipeline' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||
async pausePipeline(@Body() body: any, @User() user: RequestUser) {
|
||||
@@ -660,7 +652,7 @@ export class PlatformApiController {
|
||||
return this.platformApiService.proxy('POST', '/pipeline/pause', user, enrichedBody);
|
||||
}
|
||||
|
||||
@Post('pipeline/unpause')
|
||||
@Post('pipelines/unpause')
|
||||
@ApiOperation({ summary: 'Unpause a pipeline' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||
async unpausePipeline(@Body() body: any, @User() user: RequestUser) {
|
||||
@@ -673,7 +665,7 @@ export class PlatformApiController {
|
||||
return this.platformApiService.proxy('POST', '/pipeline/unpause', user, enrichedBody);
|
||||
}
|
||||
|
||||
@Put('pipeline/:pipelineId/memory')
|
||||
@Put('pipelines/:pipelineId/memory')
|
||||
@ApiOperation({ summary: 'Update pipeline memory configuration' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||
async updatePipelineMemory(
|
||||
@@ -696,7 +688,7 @@ export class PlatformApiController {
|
||||
|
||||
// ==================== PIPELINE METADATA ROUTES ====================
|
||||
|
||||
@Put('pipeline/:pipelineId/metadata')
|
||||
@Put('pipelines/:pipelineId/metadata')
|
||||
@ApiOperation({ summary: 'Update pipeline metadata' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||
async updatePipelineMetadata(
|
||||
@@ -728,26 +720,10 @@ export class PlatformApiController {
|
||||
);
|
||||
}
|
||||
|
||||
// ==================== Catalog ROUTES ====================
|
||||
|
||||
@Get('pipelines/catalog/tables')
|
||||
@ApiOperation({ summary: 'Get all tables available' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getAvailableTables(
|
||||
@User() user: RequestUser,
|
||||
@Query() query: Record<string, string>,
|
||||
) {
|
||||
return this.platformApiService.proxy(
|
||||
'GET',
|
||||
'/catalog/tables',
|
||||
user,
|
||||
undefined,
|
||||
query,
|
||||
);
|
||||
}
|
||||
// ==================== PIPELINE VALIDATION ====================
|
||||
|
||||
@Get('pipelines/catalog/schemas')
|
||||
@ApiOperation({ summary: 'Get all schemas available' })
|
||||
@ApiOperation({ summary: 'Get available schemas' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getAvailableSchemas(
|
||||
@User() user: RequestUser,
|
||||
@@ -755,7 +731,7 @@ export class PlatformApiController {
|
||||
) {
|
||||
return this.platformApiService.proxy(
|
||||
'GET',
|
||||
'/catalog/schemas',
|
||||
`/catalog/schemas`,
|
||||
user,
|
||||
undefined,
|
||||
query,
|
||||
@@ -763,25 +739,25 @@ export class PlatformApiController {
|
||||
}
|
||||
|
||||
@Post('pipelines/catalog/tables/validate')
|
||||
@ApiOperation({ summary: 'Validate tables and schemas' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||
@ApiOperation({ summary: 'Validate Table and Schema' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async validateTableAndSchema(
|
||||
@Body() payload: ValidationTableDTO,
|
||||
@User() user: RequestUser,
|
||||
@Query() query: Record<string, string>,
|
||||
@Body() validateTablesDto: ValidateTablesDTO[]
|
||||
) {
|
||||
return this.platformApiService.proxy(
|
||||
'POST',
|
||||
'/catalog/tables/validate',
|
||||
`/catalog/tables/validate`,
|
||||
user,
|
||||
validateTablesDto,
|
||||
payload,
|
||||
query,
|
||||
);
|
||||
}
|
||||
|
||||
// ==================== PIPELINE RUN ROUTES ====================
|
||||
|
||||
@Get('pipeline/:pipelineId/pipeline_run')
|
||||
@Get('pipelines/:pipelineId/pipeline_run')
|
||||
@ApiOperation({ summary: 'Get pipeline runs for a pipeline' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getPipelineRuns(
|
||||
@@ -799,7 +775,7 @@ export class PlatformApiController {
|
||||
);
|
||||
}
|
||||
|
||||
@Get('pipeline/:pipelineId/pipeline_run/:runId')
|
||||
@Get('pipelines/:pipelineId/pipeline_run/:runId')
|
||||
@ApiOperation({ summary: 'Get specific pipeline run' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getPipelineRun(
|
||||
@@ -816,7 +792,7 @@ export class PlatformApiController {
|
||||
);
|
||||
}
|
||||
|
||||
@Get('pipeline/pipeline_run/:runId/logs')
|
||||
@Get('pipelines/pipeline_run/:runId/logs')
|
||||
@ApiOperation({ summary: 'Get pipeline run logs' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getPipelineRunLogs(
|
||||
@@ -834,6 +810,68 @@ export class PlatformApiController {
|
||||
);
|
||||
}
|
||||
|
||||
@Post('pipelines/:pipelineId/pipeline_run/:runId/cancel')
|
||||
@ApiOperation({ summary: 'Cancel a running pipeline run' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||
async cancelPipelineRun(
|
||||
@Param('pipelineId') pipelineId: string,
|
||||
@Param('runId') runId: string,
|
||||
@User() user: RequestUser,
|
||||
) {
|
||||
const normalizedPipelineId = this.normalizePipelineId(pipelineId);
|
||||
const normalizedRunId = this.normalizePipelineId(runId);
|
||||
|
||||
const status = await this.platformApiService.proxy(
|
||||
'GET',
|
||||
`/pipeline/${normalizedPipelineId}/pipeline_run`,
|
||||
user,
|
||||
);
|
||||
|
||||
if (status.length === 1) {
|
||||
throw new BadRequestException('The first pipeline cannot be canceled');
|
||||
}
|
||||
|
||||
return this.platformApiService.proxy(
|
||||
'POST',
|
||||
`/pipeline/${normalizedPipelineId}/pipeline_run/${normalizedRunId}/cancel`,
|
||||
user,
|
||||
);
|
||||
}
|
||||
|
||||
@Get('pipelines/:pipelineId/pipeline_run/:runId/jobs')
|
||||
@ApiOperation({
|
||||
summary: 'Get pipeline run jobs',
|
||||
description: 'Proxies platform-api DB-backed job runs and returns `{ jobs: [...] }`.',
|
||||
})
|
||||
@ApiOkResponse({
|
||||
description: 'DB-backed job runs for the selected pipeline run.',
|
||||
schema: {
|
||||
type: 'object',
|
||||
properties: {
|
||||
jobs: {
|
||||
type: 'array',
|
||||
items: { type: 'object' },
|
||||
},
|
||||
},
|
||||
required: ['jobs'],
|
||||
},
|
||||
})
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getPipelineRunJobs(
|
||||
@Param('pipelineId') pipelineId: string,
|
||||
@Param('runId') runId: string,
|
||||
@User() user: RequestUser,
|
||||
) {
|
||||
const normalizedPipelineId = this.normalizePipelineId(pipelineId);
|
||||
const decodedRunId = this.decodePathParam(runId);
|
||||
|
||||
return this.platformApiService.proxy(
|
||||
'GET',
|
||||
`/pipeline/${normalizedPipelineId}/pipeline_run/${decodedRunId}/jobs`,
|
||||
user,
|
||||
);
|
||||
}
|
||||
|
||||
// ==================== JOBS - COLUMN EDITING ROUTES ====================
|
||||
|
||||
@Put('jobs/:jobId/input')
|
||||
@@ -933,41 +971,54 @@ export class PlatformApiController {
|
||||
);
|
||||
}
|
||||
|
||||
// ==================== JOBS - JDBC SYNC MODE ROUTES ====================
|
||||
|
||||
@Get('jobs/jdbc/:jobId')
|
||||
@ApiOperation({ summary: 'Get JDBC job details' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getJdbcJob(@Param('jobId') jobId: string, @User() user: RequestUser) {
|
||||
// Normalize job ID for Platform API (replace - with _)
|
||||
const normalizedJobId = this.normalizeJobId(jobId);
|
||||
return this.platformApiService.proxy('GET', `/jobs/jdbc/${normalizedJobId}`, user);
|
||||
}
|
||||
|
||||
@Post('jobs/jdbc/:jobId/sync-mode')
|
||||
@ApiOperation({ summary: 'Update JDBC job sync mode' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||
async updateJdbcSyncMode(
|
||||
@Param('jobId') jobId: string,
|
||||
@Body() body: any,
|
||||
@Delete('pipelines/:pipelineId/inputs/:inputId')
|
||||
@ApiOperation({ summary: 'Mark a table as deleted and delete its associated job via platform-api' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE)
|
||||
@UseGuards(PipelineExecutionGuard)
|
||||
async deleteTable(
|
||||
@Param('pipelineId') pipelineId: string,
|
||||
@Param('inputId') inputId: string,
|
||||
@Body() body: { table_name: string },
|
||||
@User() user: RequestUser,
|
||||
) {
|
||||
// Normalize job ID for Platform API (replace - with _)
|
||||
const normalizedJobId = this.normalizeJobId(jobId);
|
||||
const tableName = body.table_name;
|
||||
const info = {
|
||||
customer_id: user.customer_id,
|
||||
customer: user.customer_name,
|
||||
user_id: user.user_id,
|
||||
};
|
||||
|
||||
const result = await this.platformApiService.proxy(
|
||||
'POST',
|
||||
`/jobs/jdbc/${normalizedJobId}/sync-mode`,
|
||||
user,
|
||||
body,
|
||||
);
|
||||
this.logger.info('deleteTable: marking table as deleted', { inputId, tableName });
|
||||
const updatedInput: any = await this.inputsService.markTableDeleted({ input_id: inputId, table_name: tableName, info });
|
||||
this.logger.info('deleteTable: table marked as deleted', { inputId, tableName });
|
||||
|
||||
// Sync to DynamoDB (pass raw jobId for pipeline extraction)
|
||||
await this.syncJdbcSyncModeToDynamoDB(jobId, body, user);
|
||||
try {
|
||||
const normalizedPipelineId = this.normalizePipelineId(pipelineId);
|
||||
this.logger.info('deleteTable: fetching pipeline from platform-api', { pipelineId, normalizedPipelineId });
|
||||
const platformPipeline = await this.platformApiService.proxy('GET', `/pipeline/${normalizedPipelineId}`, user);
|
||||
this.logger.info('deleteTable: pipeline fetched', { jobCount: platformPipeline?.jobs?.length });
|
||||
|
||||
return result;
|
||||
const job = platformPipeline?.jobs?.find((j: any) => j.input?.table_name === tableName);
|
||||
if (!job) throw new NotFoundException(`Job for table '${tableName}' not found in pipeline`);
|
||||
|
||||
this.logger.info('deleteTable: deleting job from platform-api', { jobId: job.job_id });
|
||||
await this.platformApiService.proxy('DELETE', `/jobs/${job.job_id}`, user);
|
||||
this.logger.info('deleteTable: job deleted', { jobId: job.job_id });
|
||||
|
||||
return { name: tableName, is_deleted: updatedInput.is_deleted ?? true, deleted_at: updatedInput.deleted_at };
|
||||
} catch (error) {
|
||||
this.logger.error('deleteTable: platform-api delete failed, attempting rollback', { tableName, error: error.message });
|
||||
try {
|
||||
await this.inputsService.unmarkTableDeleted({ input_id: inputId, table_name: tableName, info });
|
||||
} catch (rollbackError) {
|
||||
this.logger.error('deleteTable: rollback failed', { tableName, error: rollbackError.message });
|
||||
}
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
// ==================== JOBS - JDBC SYNC MODE ROUTES ====================
|
||||
|
||||
@Get('jobs/jdbc/configs/allowed_datatypes')
|
||||
@ApiOperation({ summary: 'Get allowed datatypes for JDBC' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
@@ -979,52 +1030,6 @@ export class PlatformApiController {
|
||||
);
|
||||
}
|
||||
|
||||
// ==================== JOBS - SINGER REPLICATION ROUTES ====================
|
||||
|
||||
@Get('jobs/singer/:jobId')
|
||||
@ApiOperation({ summary: 'Get Singer job details' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getSingerJob(@Param('jobId') jobId: string, @User() user: RequestUser) {
|
||||
// Normalize job ID for Platform API (replace - with _)
|
||||
const normalizedJobId = this.normalizeJobId(jobId);
|
||||
return this.platformApiService.proxy('GET', `/jobs/singer/${normalizedJobId}`, user);
|
||||
}
|
||||
|
||||
@Post('jobs/singer/:jobId/sync-mode')
|
||||
@ApiOperation({ summary: 'Update Singer job sync mode' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||
async updateSingerSyncMode(
|
||||
@Param('jobId') jobId: string,
|
||||
@Body() body: any,
|
||||
@User() user: RequestUser,
|
||||
) {
|
||||
// Normalize job ID for Platform API (replace - with _)
|
||||
const normalizedJobId = this.normalizeJobId(jobId);
|
||||
|
||||
const result = await this.platformApiService.proxy(
|
||||
'POST',
|
||||
`/jobs/singer/${normalizedJobId}/sync-mode`,
|
||||
user,
|
||||
body,
|
||||
);
|
||||
|
||||
// Sync to DynamoDB (pass raw jobId for pipeline extraction)
|
||||
await this.syncSingerSyncModeToDynamoDB(jobId, body, user);
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
// ==================== JOBS - S3 ROUTES ====================
|
||||
|
||||
@Get('jobs/s3/:jobId')
|
||||
@ApiOperation({ summary: 'Get S3 job details' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getS3Job(@Param('jobId') jobId: string, @User() user: RequestUser) {
|
||||
// Normalize job ID for Platform API (replace - with _)
|
||||
const normalizedJobId = this.normalizeJobId(jobId);
|
||||
return this.platformApiService.proxy('GET', `/jobs/s3/${normalizedJobId}`, user);
|
||||
}
|
||||
|
||||
// ==================== HEALTH ROUTE ====================
|
||||
|
||||
@Get('health')
|
||||
|
||||
@@ -0,0 +1,9 @@
|
||||
import { ApiProperty } from "@nestjs/swagger";
|
||||
|
||||
export class ValidationTableDTO {
|
||||
@ApiProperty()
|
||||
tables: Array<{
|
||||
table_name: string;
|
||||
table_schema: string;
|
||||
}>
|
||||
}
|
||||
@@ -7,9 +7,11 @@ 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],
|
||||
imports: [ElasticsearchModule, DynamoDBModule, CustomersModule, CatalogModule, InputsModule],
|
||||
controllers: [PlatformApiController],
|
||||
providers: [PlatformApiService, DadosferaLogger],
|
||||
exports: [PlatformApiService],
|
||||
|
||||
@@ -89,6 +89,12 @@ export class PlatformApiService {
|
||||
|
||||
// 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);
|
||||
}
|
||||
|
||||
@@ -101,6 +107,8 @@ export class PlatformApiService {
|
||||
method: method.toUpperCase(),
|
||||
});
|
||||
|
||||
this.logger.error(error)
|
||||
|
||||
if (error instanceof HttpException) {
|
||||
throw error;
|
||||
}
|
||||
|
||||
@@ -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,20 @@
|
||||
import { Test, TestingModule } from '@nestjs/testing';
|
||||
import { ReleaseNoteController } from './release_note.controller';
|
||||
import { ReleaseNoteService } from './release_note.service';
|
||||
|
||||
describe('ReleaseNoteController', () => {
|
||||
let controller: ReleaseNoteController;
|
||||
|
||||
beforeEach(async () => {
|
||||
const module: TestingModule = await Test.createTestingModule({
|
||||
controllers: [ReleaseNoteController],
|
||||
providers: [ReleaseNoteService],
|
||||
}).compile();
|
||||
|
||||
controller = module.get<ReleaseNoteController>(ReleaseNoteController);
|
||||
});
|
||||
|
||||
it('should be defined', () => {
|
||||
expect(controller).toBeDefined();
|
||||
});
|
||||
});
|
||||
@@ -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,18 @@
|
||||
import { Test, TestingModule } from '@nestjs/testing';
|
||||
import { ReleaseNoteService } from './release_note.service';
|
||||
|
||||
describe('ReleaseNoteService', () => {
|
||||
let service: ReleaseNoteService;
|
||||
|
||||
beforeEach(async () => {
|
||||
const module: TestingModule = await Test.createTestingModule({
|
||||
providers: [ReleaseNoteService],
|
||||
}).compile();
|
||||
|
||||
service = module.get<ReleaseNoteService>(ReleaseNoteService);
|
||||
});
|
||||
|
||||
it('should be defined', () => {
|
||||
expect(service).toBeDefined();
|
||||
});
|
||||
});
|
||||
@@ -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()}`);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
+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;
|
||||
|
||||
@@ -125,9 +125,13 @@ export class DynamoDBService {
|
||||
},
|
||||
});
|
||||
|
||||
const { Item } = await this.documentClient.send(getCommand);
|
||||
|
||||
return Item as InputDocument | null;
|
||||
try {
|
||||
const { Item } = await this.documentClient.send(getCommand);
|
||||
return Item as InputDocument | null;
|
||||
} catch (error) {
|
||||
this.logger.error('DynamoDB: findInput failed', { inputId, clientId, error: error.message });
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
async deleteInput(clientId: string, inputId: string): Promise<void> {
|
||||
@@ -203,7 +207,6 @@ export class DynamoDBService {
|
||||
updatedTable.reference_column = changes.reference_column;
|
||||
}
|
||||
}
|
||||
|
||||
tables[tableIndex] = updatedTable;
|
||||
|
||||
// Save updated document
|
||||
@@ -231,4 +234,5 @@ export class DynamoDBService {
|
||||
throw error;
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -358,6 +358,81 @@ export class ElasticsearchService {
|
||||
}
|
||||
}
|
||||
|
||||
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,
|
||||
|
||||
@@ -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
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
}
|
||||
@@ -9,6 +9,7 @@ interface IMetadata {
|
||||
details?: string;
|
||||
sensitive?: string;
|
||||
roles?: string[];
|
||||
customer_modules?: string[];
|
||||
is_data_manager?: boolean;
|
||||
access_token?: string;
|
||||
host?: string;
|
||||
|
||||
@@ -0,0 +1,13 @@
|
||||
/**
|
||||
* Jest setupFiles hook: provide dummy values for env vars that modules read
|
||||
* at *import time* (module-load side effects), so unit specs that transitively
|
||||
* import those modules don't crash before a single test runs.
|
||||
*
|
||||
* These are never exercised by unit tests — the services that use them are
|
||||
* stubbed via DI — but the values must exist because the reads happen at
|
||||
* module evaluation, before any mock is installed.
|
||||
*
|
||||
* DUC_URL: read at load time in src/modules/duc/client.config.ts, reached via
|
||||
* the AuthClientService import chain.
|
||||
*/
|
||||
process.env.DUC_URL = process.env.DUC_URL || 'duc:localhost:0';
|
||||
Reference in New Issue
Block a user