From a597c410702986c880caa315f3f77737bb62a2d0 Mon Sep 17 00:00:00 2001 From: Rafael Date: Sat, 15 Aug 2026 19:45:09 -0300 Subject: [PATCH 01/22] feat(cdc): expose CDC routes over REST (on main + beta protospack) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Re-applied onto fresh origin/main. pipelinesV2: live-status/pause/unpause/ restart/reset-state routes (per-route @RequireSomePermission — GET for status, UPDATE for mutations — matching main's current auth convention). connection-test: POST /connection-test/cdc-prerequisites. inputs: POST /inputs/cdc -> InputCreateCdc. Consumes protospack 3.40.0-beta.15-cdc.0 tarball; CI guard added. Co-Authored-By: WOZCODE --- .github/workflows/test.yml | 12 + docsfera.json | 437 ++++++++++++++++++ package-lock.json | 9 +- package.json | 2 +- scripts/check-protospack-dep.js | 76 +++ .../connection-test.controller.ts | 19 + .../connection-test.service.ts | 16 + .../connection-test/dto/connection-test.ts | 28 ++ src/modules/inputs/dtos/input.model.ts | 20 + src/modules/inputs/inputs.controller.ts | 21 + src/modules/inputs/inputs.service.ts | 26 +- .../pipelinesV2/pipelines.controller.ts | 74 +++ src/modules/pipelinesV2/pipelines.service.ts | 56 +++ 13 files changed, 790 insertions(+), 6 deletions(-) create mode 100644 scripts/check-protospack-dep.js diff --git a/.github/workflows/test.yml b/.github/workflows/test.yml index 977c0fb..d293ed9 100644 --- a/.github/workflows/test.yml +++ b/.github/workflows/test.yml @@ -2,9 +2,21 @@ name: Test on: pull_request: branches: + - beta - main jobs: + # Blocks a local (file:/tarball/overlay) protospack-v2 dependency from + # reaching staging (beta) or prod (main). + protospack-dep-guard: + if: github.base_ref == 'beta' || github.base_ref == 'main' + runs-on: [self-hosted, prd] + steps: + - name: Checkout + uses: actions/checkout@v4 + - name: Check protospack-v2 is consumed from the registry + run: node scripts/check-protospack-dep.js + test: runs-on: [self-hosted, prd] env: diff --git a/docsfera.json b/docsfera.json index b3a85b5..5bc7b70 100644 --- a/docsfera.json +++ b/docsfera.json @@ -3001,6 +3001,45 @@ ] } }, + "/inputs/cdc": { + "post": { + "operationId": "InputsController_createCdc", + "parameters": [], + "requestBody": { + "required": true, + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/CreateCdcInputReq" + } + } + } + }, + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/Input" + } + } + } + }, + "201": { + "description": "" + } + }, + "tags": [ + "Inputs" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, "/inputs/{id}": { "get": { "operationId": "InputsController_findOne", @@ -3960,6 +3999,264 @@ ] } }, + "/pipelinesV2/{id}/live-status": { + "get": { + "operationId": "PipelinesController_getLiveStatus", + "parameters": [ + { + "name": "dadosfera-lang", + "in": "header", + "required": false, + "schema": { + "enum": [ + "pt-br", + "en-us" + ], + "type": "string" + } + }, + { + "name": "id", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "PipelinesV2" + ], + "security": [ + { + "access-token": [] + }, + { + "access-token": [] + } + ] + } + }, + "/pipelinesV2/{id}/pause": { + "post": { + "operationId": "PipelinesController_pause", + "parameters": [ + { + "name": "dadosfera-lang", + "in": "header", + "required": false, + "schema": { + "enum": [ + "pt-br", + "en-us" + ], + "type": "string" + } + }, + { + "name": "id", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "PipelinesV2" + ], + "security": [ + { + "access-token": [] + }, + { + "access-token": [] + } + ] + } + }, + "/pipelinesV2/{id}/unpause": { + "post": { + "operationId": "PipelinesController_unpause", + "parameters": [ + { + "name": "dadosfera-lang", + "in": "header", + "required": false, + "schema": { + "enum": [ + "pt-br", + "en-us" + ], + "type": "string" + } + }, + { + "name": "id", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "PipelinesV2" + ], + "security": [ + { + "access-token": [] + }, + { + "access-token": [] + } + ] + } + }, + "/pipelinesV2/{id}/restart": { + "post": { + "operationId": "PipelinesController_restart", + "parameters": [ + { + "name": "dadosfera-lang", + "in": "header", + "required": false, + "schema": { + "enum": [ + "pt-br", + "en-us" + ], + "type": "string" + } + }, + { + "name": "id", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "PipelinesV2" + ], + "security": [ + { + "access-token": [] + }, + { + "access-token": [] + } + ] + } + }, + "/pipelinesV2/{id}/jobs/{jobId}/reset-state": { + "post": { + "operationId": "PipelinesController_resetJobState", + "parameters": [ + { + "name": "dadosfera-lang", + "in": "header", + "required": false, + "schema": { + "enum": [ + "pt-br", + "en-us" + ], + "type": "string" + } + }, + { + "name": "id", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + }, + { + "name": "jobId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "PipelinesV2" + ], + "security": [ + { + "access-token": [] + }, + { + "access-token": [] + } + ] + } + }, "/transformations": { "post": { "operationId": "TransformationsController_create", @@ -7361,6 +7658,45 @@ ] } }, + "/connection-test/cdc-prerequisites": { + "post": { + "operationId": "ConnectionTestController_validateCdcPrerequisites", + "parameters": [], + "requestBody": { + "required": true, + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/ValidateCdcPrerequisitesReq" + } + } + } + }, + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/ValidateCdcPrerequisitesRes" + } + } + } + } + }, + "tags": [ + "Connection Test" + ], + "security": [ + { + "access-token": [] + }, + { + "access-token": [] + } + ] + } + }, "/oauth/hubspot": { "get": { "operationId": "OauthController_oauthHubspot", @@ -10504,6 +10840,51 @@ "updated_at" ] }, + "CdcTableReq": { + "type": "object", + "properties": { + "name": { + "type": "string" + }, + "table_schema": { + "type": "string" + }, + "primary_keys": { + "type": "array", + "items": { + "type": "string" + } + } + }, + "required": [ + "name" + ] + }, + "CreateCdcInputReq": { + "type": "object", + "properties": { + "name": { + "type": "string" + }, + "plugin": { + "type": "string" + }, + "tables": { + "type": "array", + "items": { + "$ref": "#/components/schemas/CdcTableReq" + } + }, + "read_only": { + "type": "boolean" + } + }, + "required": [ + "name", + "plugin", + "tables" + ] + }, "ICreatePipelineV2Req": { "type": "object", "properties": { @@ -11730,6 +12111,62 @@ "tables_metadata" ] }, + "ValidateCdcPrerequisitesReq": { + "type": "object", + "properties": { + "plugin": { + "type": "string" + }, + "connection_id": { + "type": "string" + } + }, + "required": [ + "plugin", + "connection_id" + ] + }, + "CdcCheckDto": { + "type": "object", + "properties": { + "name": { + "type": "string" + }, + "expected": { + "type": "string" + }, + "actual": { + "type": "string" + }, + "passed": { + "type": "boolean" + } + }, + "required": [ + "name", + "expected", + "actual", + "passed" + ] + }, + "ValidateCdcPrerequisitesRes": { + "type": "object", + "properties": { + "operation_result": { + "type": "boolean" + }, + "checks": { + "type": "array", + "items": { + "$ref": "#/components/schemas/CdcCheckDto" + } + } + }, + "required": [ + "operation_result", + "checks" + ] + }, "INote": { "type": "object", "properties": { diff --git a/package-lock.json b/package-lock.json index dcb4ea9..6007a70 100644 --- a/package-lock.json +++ b/package-lock.json @@ -16,7 +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-v2": "^3.40.0-beta.10", + "@dadosfera/protospack-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.0.tgz", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", @@ -1735,9 +1735,10 @@ } }, "node_modules/@dadosfera/protospack-v2": { - "version": "3.40.0-beta.10", - "resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.40.0-beta.10.tgz", - "integrity": "sha512-F45dSEIKG+gwwDMYHayA242bFwhFTJbZm26KesQbGhf4I45ur6hYZD8dK0voNuVQInLMQhtm6+7th6/jJ8xpTQ==", + "version": "3.40.0-beta.15-cdc.0", + "resolved": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.0.tgz", + "integrity": "sha512-ySucTiWmZ9eiw/uujSY8MFaj2+V/J+42y54P737PlrIXJgTSlLZalhwyKNkHuwo1Az8Uj5gC1J9AN4/sn9mIyg==", + "license": "ISC", "dependencies": { "@grpc/grpc-js": "^1.9.3", "rxjs": "^7.5.5" diff --git a/package.json b/package.json index fdbc61e..6af93ac 100644 --- a/package.json +++ b/package.json @@ -34,7 +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-v2": "^3.40.0-beta.10", + "@dadosfera/protospack-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.0.tgz", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", diff --git a/scripts/check-protospack-dep.js b/scripts/check-protospack-dep.js new file mode 100644 index 0000000..1d44c0c --- /dev/null +++ b/scripts/check-protospack-dep.js @@ -0,0 +1,76 @@ +#!/usr/bin/env node +/* + * CI guard: fail if @dadosfera/protospack-v2 is consumed from a LOCAL ref + * (file:/link:/git/relative path/bare tarball) instead of the CodeArtifact + * registry. + * + * Only local consumption is blocked. Versions published to CodeArtifact — + * including alpha/beta/rc prereleases produced by the alpha/beta branches — + * are fine; those resolve to a registry URL in the lockfile. The thing that + * must NOT reach beta (staging) or main (prod) is a dependency wired to a + * local `npm pack` tarball / overlay. Runs in the PR test workflow for PRs + * targeting beta/main and exits non-zero on any local ref. + */ +const fs = require('fs'); +const path = require('path'); + +const PKG = '@dadosfera/protospack-v2'; +const root = path.resolve(__dirname, '..'); +const pkg = JSON.parse(fs.readFileSync(path.join(root, 'package.json'), 'utf8')); + +const problems = []; + +// A dependency SPEC is local if it's a filesystem path, symlink, git ref, or a +// bare tarball path. A plain semver (incl. prereleases like 3.35.0-beta.1) +// resolves from the registry and is allowed. +function isLocalSpec(spec) { + return /^(file:|link:|git[:+]|\.\.?\/|\/|~\/)/.test(spec) || spec.endsWith('.tgz'); +} + +const spec = + (pkg.dependencies && pkg.dependencies[PKG]) || + (pkg.devDependencies && pkg.devDependencies[PKG]); + +if (!spec) { + problems.push(`${PKG} is not listed as a dependency at all.`); +} else if (isLocalSpec(spec)) { + problems.push(`${PKG} points at a local path/tarball/git ref: "${spec}".`); +} + +// Also catch a lockfile resolved to a LOCAL ref even if package.json looks +// clean. A registry URL (https://.../-/*.tgz) is the normal published +// resolution and is fine — only file: refs and bare local tarball paths +// (no http host) are blocked. Prerelease VERSIONS are not flagged: an +// alpha/beta/rc published to CodeArtifact resolves to a registry URL. +const lockPath = path.join(root, 'package-lock.json'); +if (fs.existsSync(lockPath)) { + const lock = JSON.parse(fs.readFileSync(lockPath, 'utf8')); + const nodes = { ...(lock.packages || {}), ...(lock.dependencies || {}) }; + for (const [name, node] of Object.entries(nodes)) { + if (!name.includes('protospack-v2') || !node) continue; + const resolved = node.resolved || ''; + const isLocal = + resolved.startsWith('file:') || + (resolved.endsWith('.tgz') && !/^https?:\/\//.test(resolved)); + if (isLocal) { + problems.push( + `package-lock.json resolves ${PKG} to a local ref: "${resolved}".`, + ); + } + } +} + +if (problems.length) { + console.error('✗ protospack-v2 dependency guard FAILED:'); + for (const p of problems) console.error(' - ' + p); + console.error( + '\nMerging to beta/main requires ' + + PKG + + ' to come from CodeArtifact, not a local tarball/overlay. Publish ' + + 'protospack-v2 (a beta prerelease is fine for the beta branch) and ' + + 'repoint this dependency before merging.', + ); + process.exit(1); +} + +console.log(`✓ ${PKG} is consumed from the registry: "${spec}"`); diff --git a/src/modules/connection-test/connection-test.controller.ts b/src/modules/connection-test/connection-test.controller.ts index 029da68..939aded 100644 --- a/src/modules/connection-test/connection-test.controller.ts +++ b/src/modules/connection-test/connection-test.controller.ts @@ -22,6 +22,8 @@ import { ConnectionTestListTablesRes, GetTableMetadataRes, GetTableMetadataReq, + ValidateCdcPrerequisitesReq, + ValidateCdcPrerequisitesRes, } from './dto/connection-test'; import { DadosferaLogger } from '@dadosfera/dadosfera-logs'; import { Authenticated, RequireModule } from 'src/decorators/authentication.decorator'; @@ -125,4 +127,21 @@ export class ConnectionTestController { user.customer_name, ); } + + @Post('cdc-prerequisites') + @ApiOkResponse({ type: ValidateCdcPrerequisitesRes }) + @HttpCode(HttpStatus.OK) + async validateCdcPrerequisites( + @User() user: RequestUser, + @Body(new ValidationPipe()) body: ValidateCdcPrerequisitesReq, + ) { + this.logger.info('/connection-test/cdc-prerequisites', { + user: user.user_id, + customer: user.customer_name, + }); + return this.connectionTestService.validateCdcPrerequisites( + body, + user.customer_name, + ); + } } diff --git a/src/modules/connection-test/connection-test.service.ts b/src/modules/connection-test/connection-test.service.ts index 3b92e12..096a34e 100644 --- a/src/modules/connection-test/connection-test.service.ts +++ b/src/modules/connection-test/connection-test.service.ts @@ -13,6 +13,8 @@ import { ConnectionTestPingRes, GetTableMetadataReq, GetTableMetadataRes, + ValidateCdcPrerequisitesReq, + ValidateCdcPrerequisitesRes, } from './dto/connection-test'; import { ConnectionClientService } from '../connection/client.service'; import { @@ -188,4 +190,18 @@ export class ConnectionTestService { }), ); } + + async validateCdcPrerequisites( + body: ValidateCdcPrerequisitesReq, + customer_name: string, + ): Promise { + const { plugin, connection_id } = body; + return lastValueFrom( + this.connectionTestReadClient.ValidateCdcPrerequisites({ + connection_id, + customer_name, + plugin, + }), + ); + } } diff --git a/src/modules/connection-test/dto/connection-test.ts b/src/modules/connection-test/dto/connection-test.ts index b6bb4c4..fd9ae91 100644 --- a/src/modules/connection-test/dto/connection-test.ts +++ b/src/modules/connection-test/dto/connection-test.ts @@ -131,3 +131,31 @@ export class GetTableMetadataRes { @ApiProperty({ type: [TableMetadataDto] }) tables_metadata: TableMetadataDto[]; } + +export class CdcCheckDto { + @ApiProperty() + name: string; + @ApiProperty() + expected: string; + @ApiProperty() + actual: string; + @ApiProperty() + passed: boolean; +} + +export class ValidateCdcPrerequisitesReq { + @ApiProperty() + @IsString() + plugin: string; + + @ApiProperty() + @IsString() + connection_id: string; +} + +export class ValidateCdcPrerequisitesRes { + @ApiProperty() + operation_result: boolean; + @ApiProperty({ type: [CdcCheckDto] }) + checks: CdcCheckDto[]; +} diff --git a/src/modules/inputs/dtos/input.model.ts b/src/modules/inputs/dtos/input.model.ts index c026f8d..fa73fd1 100644 --- a/src/modules/inputs/dtos/input.model.ts +++ b/src/modules/inputs/dtos/input.model.ts @@ -65,3 +65,23 @@ export class CreateInputReq extends OmitType(Input, [ 'created_at', 'updated_at', ]) {} + +export class CdcTableReq { + @ApiProperty() + name: string; + @ApiPropertyOptional() + table_schema?: string; + @ApiPropertyOptional({ type: [String] }) + primary_keys?: string[]; +} + +export class CreateCdcInputReq { + @ApiProperty() + name: string; + @ApiProperty() + plugin: string; // mysql_cdc (v1) + @ApiProperty({ type: [CdcTableReq] }) + tables: CdcTableReq[]; + @ApiPropertyOptional() + read_only?: boolean; +} diff --git a/src/modules/inputs/inputs.controller.ts b/src/modules/inputs/inputs.controller.ts index a8cc143..3d753fd 100644 --- a/src/modules/inputs/inputs.controller.ts +++ b/src/modules/inputs/inputs.controller.ts @@ -15,6 +15,7 @@ import { PERMISSIONS_GROUPS } from '../../authentication/permissions.enum'; import { AuthenticateCondition } from 'src/decorators/authentication.decorator'; import { ApiOkResponse, ApiTags } from '@nestjs/swagger'; import { + CreateCdcInputReq, CreateInputReq, GetAvailableEntitiesReq, GetAvailableEntitiesRes, @@ -105,6 +106,26 @@ export class InputsController { return response; } + @Post('cdc') + @ApiInternalOnlyEndpoint() + @ApiOkResponse({ type: Input }) + async createCdc( + @Body() body: CreateCdcInputReq, + @User() user: RequestUser, + ) { + const info: Info = { + user_id: user.user_id, + customer: user.customer_name, + customer_id: user.customer_id, + }; + this.logger.info(`/inputs/cdc - ON CREATE CDC INPUT ROUTE`, { + user: info.user_id, + customer: info.customer, + }); + + return this.inputService.createCdc({ body, info }); + } + @ApiInternalOnlyEndpoint() @Get() async findAll(@User() user: RequestUser) { diff --git a/src/modules/inputs/inputs.service.ts b/src/modules/inputs/inputs.service.ts index 99a6753..a70b608 100644 --- a/src/modules/inputs/inputs.service.ts +++ b/src/modules/inputs/inputs.service.ts @@ -15,6 +15,7 @@ import { Input } from '@dadosfera/protospack-v2'; import { GetAvailableEntitiesRequest, InputCreateGenericRequest, + InputCreateCdcRequest, InputCreateS3Request, InputNewCreateRequest, InputUpdateResponse, @@ -22,7 +23,7 @@ import { 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 { CreateCdcInputReq, CreateInputReq } from './dtos/input.model'; import { Metadata } from '@grpc/grpc-js'; @Injectable() @@ -183,6 +184,29 @@ export class InputsService { return { input: adjustedInput }; } + async createCdc(data: { body: CreateCdcInputReq; info: Info }) { + const { body, info } = data; + + const inputCreateCdcRequest: InputCreateCdcRequest = { + input: { + name: body.name, + plugin: body.plugin, + read_only: body.read_only ?? true, + tables: body.tables.map((t) => ({ + table_schema: t.table_schema, + table_name: t.name, + primary_keys: t.primary_keys ?? [], + })), + }, + info, + }; + + const { input } = await lastValueFrom( + this.inputWriteService.InputCreateCdc(inputCreateCdcRequest), + ); + return { input }; + } + async getAvailableEntities(data: GetAvailableEntitiesRequest) { return lastValueFrom(this.inputReadService.GetAvailableEntities(data)); } diff --git a/src/modules/pipelinesV2/pipelines.controller.ts b/src/modules/pipelinesV2/pipelines.controller.ts index c24c090..6c529c5 100644 --- a/src/modules/pipelinesV2/pipelines.controller.ts +++ b/src/modules/pipelinesV2/pipelines.controller.ts @@ -510,4 +510,78 @@ export class PipelinesController { return response; } + + // ---- CDC pipeline operations (Kafka Connect backed) ---- + // These operate on an existing pipeline, so they require UPDATE (not CREATE). + + @Get(':id/live-status') + @ApiInternalOnlyEndpoint() + @RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.GET) + async getLiveStatus( + @Param('id') id: string, + @User() user: RequestUser, + ): Promise { + this.logger.info('PipelinesController - getLiveStatus', { id }); + const metadata = PackTheMetadata(user); + return this.pipelinesClientService.getLiveStatus(id, metadata); + } + + @Post(':id/pause') + @HttpCode(HttpStatus.OK) + @ApiInternalOnlyEndpoint() + @RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE) + async pause( + @Param('id') id: string, + @User() user: RequestUser, + ): Promise { + this.logger.info('PipelinesController - pause', { id }); + const metadata = PackTheMetadata(user); + return this.pipelinesClientService.pause(id, metadata); + } + + @Post(':id/unpause') + @HttpCode(HttpStatus.OK) + @ApiInternalOnlyEndpoint() + @RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE) + async unpause( + @Param('id') id: string, + @User() user: RequestUser, + ): Promise { + this.logger.info('PipelinesController - unpause', { id }); + const metadata = PackTheMetadata(user); + return this.pipelinesClientService.unpause(id, metadata); + } + + @Post(':id/restart') + @HttpCode(HttpStatus.OK) + @ApiInternalOnlyEndpoint() + @RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE) + async restart( + @Param('id') id: string, + @User() user: RequestUser, + ): Promise { + this.logger.info('PipelinesController - restart', { id }); + const metadata = PackTheMetadata(user); + return this.pipelinesClientService.restart(id, metadata); + } + + @Post(':id/jobs/:jobId/reset-state') + @HttpCode(HttpStatus.OK) + @ApiInternalOnlyEndpoint() + @RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE) + async resetJobState( + @Param('id') id: string, + @Param('jobId') jobId: string, + @Body() body: { schedule_minutes?: number }, + @User() user: RequestUser, + ): Promise { + this.logger.info('PipelinesController - resetJobState', { id, jobId }); + const metadata = PackTheMetadata(user); + return this.pipelinesClientService.resetJobState( + id, + jobId, + body?.schedule_minutes, + metadata, + ); + } } diff --git a/src/modules/pipelinesV2/pipelines.service.ts b/src/modules/pipelinesV2/pipelines.service.ts index 41d9354..b875cb7 100644 --- a/src/modules/pipelinesV2/pipelines.service.ts +++ b/src/modules/pipelinesV2/pipelines.service.ts @@ -160,6 +160,62 @@ export class PipelinesService implements OnModuleInit { return findOnePipelineResponse; } + // CDC lifecycle operations (Kafka Connect backed). + async pause( + id: string, + metadata, + ): Promise { + this.logger.info('PipelinesClientService - Pause'); + return lastValueFrom( + this.pipelineWriteService.PipelineV2Pause({ id }, metadata), + ); + } + + async unpause( + id: string, + metadata, + ): Promise { + this.logger.info('PipelinesClientService - Unpause'); + return lastValueFrom( + this.pipelineWriteService.PipelineV2Unpause({ id }, metadata), + ); + } + + async restart( + id: string, + metadata, + ): Promise { + this.logger.info('PipelinesClientService - Restart'); + return lastValueFrom( + this.pipelineWriteService.PipelineV2Restart({ id }, metadata), + ); + } + + async resetJobState( + id: string, + job_id: string, + schedule_minutes: number | undefined, + metadata, + ): Promise { + this.logger.info('PipelinesClientService - ResetJobState'); + return lastValueFrom( + this.pipelineWriteService.PipelineV2ResetJobState( + { id, job_id, schedule_minutes }, + metadata, + ), + ); + } + + async getLiveStatus( + id: string, + metadata, + ): Promise { + this.logger.info('PipelinesClientService - GetLiveStatus'); + return lastValueFrom( + this.pipelineReadService.PipelineV2GetLiveStatus({ id }, metadata), + ); + } + async update( UpdatePipelineRequest: Messages.PipelineV2UpdateRequest, metadata, From 27e0dedea4ed4c7ff45489fddc10fd30aafe76fd Mon Sep 17 00:00:00 2001 From: Rafael Date: Sun, 16 Aug 2026 16:57:29 -0300 Subject: [PATCH 02/22] feat(maestro): expose tables[].primary_keys on /connection-test/tables Co-Authored-By: WOZCODE --- package-lock.json | 8 ++++---- package.json | 2 +- src/modules/connection-test/dto/connection-test.ts | 9 +++++++++ 3 files changed, 14 insertions(+), 5 deletions(-) diff --git a/package-lock.json b/package-lock.json index 6007a70..c46cadb 100644 --- a/package-lock.json +++ b/package-lock.json @@ -16,7 +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-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.0.tgz", + "@dadosfera/protospack-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.1.tgz", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", @@ -1735,9 +1735,9 @@ } }, "node_modules/@dadosfera/protospack-v2": { - "version": "3.40.0-beta.15-cdc.0", - "resolved": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.0.tgz", - "integrity": "sha512-ySucTiWmZ9eiw/uujSY8MFaj2+V/J+42y54P737PlrIXJgTSlLZalhwyKNkHuwo1Az8Uj5gC1J9AN4/sn9mIyg==", + "version": "3.40.0-beta.15-cdc.1", + "resolved": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.1.tgz", + "integrity": "sha512-ANHb7b7HcjXl3FMI5HRJUVGzJWTOyiepxHGPs9z1NDiXtP3KxFPZuSQBNXuJ78R/0MEGm9ALeEi4/8wyQ+fOFw==", "license": "ISC", "dependencies": { "@grpc/grpc-js": "^1.9.3", diff --git a/package.json b/package.json index 6af93ac..027c754 100644 --- a/package.json +++ b/package.json @@ -34,7 +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-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.0.tgz", + "@dadosfera/protospack-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.1.tgz", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", diff --git a/src/modules/connection-test/dto/connection-test.ts b/src/modules/connection-test/dto/connection-test.ts index fd9ae91..f7b1228 100644 --- a/src/modules/connection-test/dto/connection-test.ts +++ b/src/modules/connection-test/dto/connection-test.ts @@ -100,11 +100,20 @@ export class ConnectionTestListTablesReq { schema: string; } +export class ConnectionTestListTablesEntry { + @ApiProperty() + table_name: string; + @ApiProperty({ type: [String] }) + primary_keys: string[]; +} + export class ConnectionTestListTablesRes { @ApiProperty() operation_result: boolean; @ApiProperty() table_list: string[]; + @ApiProperty({ type: [ConnectionTestListTablesEntry] }) + tables: ConnectionTestListTablesEntry[]; } export class GetTableMetadataReq { From fc4e1c27e415af6f51c0de71af6b44fe10d137fb Mon Sep 17 00:00:00 2001 From: Rafael Date: Sun, 16 Aug 2026 16:58:00 -0300 Subject: [PATCH 03/22] chore(maestro): regenerate swagger spec for tables[].primary_keys Co-Authored-By: WOZCODE --- docsfera.json | 27 ++++++++++++++++++++++++++- 1 file changed, 26 insertions(+), 1 deletion(-) diff --git a/docsfera.json b/docsfera.json index 5bc7b70..ff2d395 100644 --- a/docsfera.json +++ b/docsfera.json @@ -12009,6 +12009,24 @@ "schema" ] }, + "ConnectionTestListTablesEntry": { + "type": "object", + "properties": { + "table_name": { + "type": "string" + }, + "primary_keys": { + "type": "array", + "items": { + "type": "string" + } + } + }, + "required": [ + "table_name", + "primary_keys" + ] + }, "ConnectionTestListTablesRes": { "type": "object", "properties": { @@ -12020,11 +12038,18 @@ "items": { "type": "string" } + }, + "tables": { + "type": "array", + "items": { + "$ref": "#/components/schemas/ConnectionTestListTablesEntry" + } } }, "required": [ "operation_result", - "table_list" + "table_list", + "tables" ] }, "GetTableMetadataReq": { From 017fd145a912ed51c169b706672ae06bc0546cd2 Mon Sep 17 00:00:00 2001 From: Rafael Date: Sun, 16 Aug 2026 21:56:34 -0300 Subject: [PATCH 04/22] CHORE: bump protospack-v2 to cdc.2 + pass name through createCdc CdcTable.name is now required; map it (== table_name) in maestro createCdc. Co-Authored-By: WOZCODE --- package-lock.json | 8 ++++---- package.json | 2 +- src/modules/inputs/inputs.service.ts | 1 + 3 files changed, 6 insertions(+), 5 deletions(-) diff --git a/package-lock.json b/package-lock.json index c46cadb..2cc5927 100644 --- a/package-lock.json +++ b/package-lock.json @@ -16,7 +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-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.1.tgz", + "@dadosfera/protospack-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.2.tgz", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", @@ -1735,9 +1735,9 @@ } }, "node_modules/@dadosfera/protospack-v2": { - "version": "3.40.0-beta.15-cdc.1", - "resolved": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.1.tgz", - "integrity": "sha512-ANHb7b7HcjXl3FMI5HRJUVGzJWTOyiepxHGPs9z1NDiXtP3KxFPZuSQBNXuJ78R/0MEGm9ALeEi4/8wyQ+fOFw==", + "version": "3.40.0-beta.15-cdc.2", + "resolved": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.2.tgz", + "integrity": "sha512-FiZXH+bSv2Skxy9YJkczQOPnZxYbca5JlJSNbSzNxuk5e044i6Mp79LAA7SQDGt5yyAEfNHq2bAorHUxf0Q0TA==", "license": "ISC", "dependencies": { "@grpc/grpc-js": "^1.9.3", diff --git a/package.json b/package.json index 027c754..882ad39 100644 --- a/package.json +++ b/package.json @@ -34,7 +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-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.1.tgz", + "@dadosfera/protospack-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.2.tgz", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", diff --git a/src/modules/inputs/inputs.service.ts b/src/modules/inputs/inputs.service.ts index a70b608..29048ed 100644 --- a/src/modules/inputs/inputs.service.ts +++ b/src/modules/inputs/inputs.service.ts @@ -195,6 +195,7 @@ export class InputsService { tables: body.tables.map((t) => ({ table_schema: t.table_schema, table_name: t.name, + name: t.name, // canonical identity == table_name (in-factory also backfills) primary_keys: t.primary_keys ?? [], })), }, From 7bd269950f329b6b81e3b7ab3300079ea733d073 Mon Sep 17 00:00:00 2001 From: Rafael Date: Sun, 16 Aug 2026 22:00:39 -0300 Subject: [PATCH 05/22] FIX: CDC table removal reconfigures the Debezium connector deleteTable called the DB-row-only DELETE /jobs/{id}, leaving the CDC connector still replicating a removed table. For CDC jobs, call DELETE /pipeline/{id}/jobs (RemoveJobsUsecase) with delete_snowflake_tables false so replication stops but landed data is kept. Batch path unchanged. Co-Authored-By: WOZCODE --- .../platform-api.controller.spec.ts | 107 ++++++++++++++++++ .../platform-api/platform-api.controller.ts | 13 ++- 2 files changed, 119 insertions(+), 1 deletion(-) create mode 100644 src/modules/platform-api/platform-api.controller.spec.ts diff --git a/src/modules/platform-api/platform-api.controller.spec.ts b/src/modules/platform-api/platform-api.controller.spec.ts new file mode 100644 index 0000000..c235dd3 --- /dev/null +++ b/src/modules/platform-api/platform-api.controller.spec.ts @@ -0,0 +1,107 @@ +// These service modules pull in gRPC client-config modules that read +// process.env at load time; mock them (hoisted before imports) so the spec +// needs no runtime env. Each mock severs an entire import subtree and still +// provides a class usable as a DI token. +jest.mock('../customers/customers.service', () => ({ CustomersService: class {} })); +jest.mock('../catalog/catalog.service', () => ({ CatalogService: class {} })); +jest.mock('../inputs/inputs.service', () => ({ InputsService: class {} })); + +import { Test, TestingModule } from '@nestjs/testing'; +import { DadosferaLogger } from '@dadosfera/dadosfera-logs'; + +import { PlatformApiController } from './platform-api.controller'; +import { PlatformApiService } from './platform-api.service'; +import { ElasticsearchService } from '../../services/elasticsearch'; +import { DynamoDBService } from '../../services/dynamodb'; +import { CustomersService } from '../customers/customers.service'; +import { CatalogService } from '../catalog/catalog.service'; +import { InputsService } from '../inputs/inputs.service'; + +const logger = { + info: (...args) => args, + error: (...args) => args, +}; + +const mockUser: any = { + customer_id: 'c1', + customer_name: 'cust', + user_id: 'u1', +}; + +describe('PlatformApiController - deleteTable', () => { + let controller: PlatformApiController; + let platformApiService: { proxy: jest.Mock }; + let inputsService: { markTableDeleted: jest.Mock; unmarkTableDeleted: jest.Mock }; + + beforeEach(async () => { + platformApiService = { proxy: jest.fn() }; + inputsService = { + markTableDeleted: jest.fn(), + unmarkTableDeleted: jest.fn(), + }; + + const module: TestingModule = await Test.createTestingModule({ + controllers: [PlatformApiController], + providers: [ + { provide: PlatformApiService, useValue: platformApiService }, + { provide: ElasticsearchService, useValue: {} }, + { provide: DynamoDBService, useValue: {} }, + { provide: CustomersService, useValue: {} }, + { provide: CatalogService, useValue: {} }, + { provide: InputsService, useValue: inputsService }, + { provide: DadosferaLogger, useValue: { logger } }, + ], + }).compile(); + + controller = module.get(PlatformApiController); + }); + + it('CDC removal reconfigures the connector', async () => { + inputsService.markTableDeleted.mockResolvedValue({ is_deleted: true, deleted_at: 't' }); + platformApiService.proxy.mockImplementation((method: string, path: string) => { + if (method === 'GET') { + return Promise.resolve({ + jobs: [{ job_id: 'p_0', input: { connector: 'cdc', table_name: 'pedidos' } }], + }); + } + return Promise.resolve({}); + }); + + await controller.deleteTable('pid', 'iid', { table_name: 'pedidos' }, mockUser); + + expect(platformApiService.proxy).toHaveBeenCalledWith( + 'DELETE', + '/pipeline/pid/jobs', + mockUser, + { job_ids: ['p_0'], delete_snowflake_tables: false }, + ); + const deletedViaJobsRoute = platformApiService.proxy.mock.calls.some( + ([method, path]: any[]) => method === 'DELETE' && path === '/jobs/p_0', + ); + expect(deletedViaJobsRoute).toBe(false); + }); + + it('batch removal unchanged', async () => { + inputsService.markTableDeleted.mockResolvedValue({ is_deleted: true, deleted_at: 't' }); + platformApiService.proxy.mockImplementation((method: string, path: string) => { + if (method === 'GET') { + return Promise.resolve({ + jobs: [{ job_id: 'p_0', input: { connector: 'jdbc', table_name: 'pedidos' } }], + }); + } + return Promise.resolve({}); + }); + + await controller.deleteTable('pid', 'iid', { table_name: 'pedidos' }, mockUser); + + expect(platformApiService.proxy).toHaveBeenCalledWith( + 'DELETE', + '/jobs/p_0', + mockUser, + ); + const reconfiguredConnector = platformApiService.proxy.mock.calls.some( + ([method, path]: any[]) => method === 'DELETE' && path === '/pipeline/pid/jobs', + ); + expect(reconfiguredConnector).toBe(false); + }); +}); diff --git a/src/modules/platform-api/platform-api.controller.ts b/src/modules/platform-api/platform-api.controller.ts index 503d117..b10b09e 100644 --- a/src/modules/platform-api/platform-api.controller.ts +++ b/src/modules/platform-api/platform-api.controller.ts @@ -964,7 +964,18 @@ export class PlatformApiController { 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); + if (job.input?.connector === 'cdc') { + // CDC: reconfigure the Debezium/Kafka-Connect connector (stop + // replicating this table); keep the landed Snowflake data. + await this.platformApiService.proxy( + 'DELETE', + `/pipeline/${normalizedPipelineId}/jobs`, + user, + { job_ids: [job.job_id], delete_snowflake_tables: false }, + ); + } else { + 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 }; From 0d771d1e4ac76ee6238cf9c6789799c17c22ce83 Mon Sep 17 00:00:00 2001 From: Rafael Date: Sun, 16 Aug 2026 22:04:00 -0300 Subject: [PATCH 06/22] FIX: skip batch job-updates for CDC in updatePipelineInput MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit updatePlatformJobs pushes batch sync_mode/memory to jobs by positional index — meaningless for CDC and corrupting. CDC add/remove use dedicated endpoints, so skip updatePlatformJobs for CDC inputs. The input record update still runs. Co-Authored-By: WOZCODE --- .../pipelinesV2/pipelines.service.spec.ts | 95 +++++++++++++++++++ src/modules/pipelinesV2/pipelines.service.ts | 27 +++--- 2 files changed, 111 insertions(+), 11 deletions(-) create mode 100644 src/modules/pipelinesV2/pipelines.service.spec.ts diff --git a/src/modules/pipelinesV2/pipelines.service.spec.ts b/src/modules/pipelinesV2/pipelines.service.spec.ts new file mode 100644 index 0000000..564639e --- /dev/null +++ b/src/modules/pipelinesV2/pipelines.service.spec.ts @@ -0,0 +1,95 @@ +// These imported modules pull in gRPC client-config / service modules that read +// process.env at load time; mock them (hoisted before imports) so the spec needs +// no runtime env. Each mock severs an entire import subtree while still providing +// a class usable as a value/DI token. Mirrors platform-api.controller.spec.ts. +jest.mock('./pipelines-client', () => ({ PipelinesClientConfiguration: class {} })); +jest.mock('../connector/client.service', () => ({ ConnectorClientService: class {} })); +jest.mock('../inputs/inputs.service', () => ({ InputsService: class {} })); +jest.mock('../transformations/transformations.service', () => ({ TransformationsService: class {} })); +jest.mock('../platform-api/platform-api.service', () => ({ PlatformApiService: class {} })); +jest.mock('src/services/nimbus/nimbus.service', () => ({ NimbusService: class {} })); +jest.mock('../catalog/catalog.service', () => ({ CatalogService: class {} })); + +import { PipelinesService } from './pipelines.service'; + +const logger = { + info: jest.fn(), + error: jest.fn(), +}; + +const cdcOldInput = { input: { plugin: 'mysql_cdc', tables: [] } }; +const batchOldInput = { input: { plugin: 'mysql', type: 'database', tables: [] } }; + +const updateResponse = { + input: { type: 'database' }, + tablesUpdate: [], + dataAssetUpdate: [], +}; + +const user: any = { customer_modules: [] }; +const updateInputDTO: any = { tables: [] }; +const info: any = { customer: 'cust' }; +const metadata: any = {}; + +function buildService(oldInput: any) { + const inputsService: any = { + findOne: jest.fn().mockResolvedValue(oldInput), + update: jest.fn().mockResolvedValue(updateResponse), + rollbackUpdate: jest.fn().mockResolvedValue({}), + }; + const nimbusService: any = { renameTable: jest.fn().mockResolvedValue({}) }; + + const service = new PipelinesService( + { logger } as any, // dadosferaLogger + {} as any, // grpcClient + {} as any, // connectorService + inputsService, // inputsService + {} as any, // transformationsService + {} as any, // platformAPI + nimbusService, // nimbusService + {} as any, // catalogService + ); + + const updatePlatformJobsSpy = jest + .spyOn(service, 'updatePlatformJobs') + .mockResolvedValue(undefined as any); + + return { service, inputsService, updatePlatformJobsSpy }; +} + +describe('PipelinesService - updatePipelineInput', () => { + afterEach(() => jest.clearAllMocks()); + + it('CDC input skips updatePlatformJobs', async () => { + const { service, inputsService, updatePlatformJobsSpy } = buildService(cdcOldInput); + + const result = await service.updatePipelineInput( + 'pipeline-id', + 'input-id', + updateInputDTO, + info, + user, + metadata, + ); + + expect(updatePlatformJobsSpy).not.toHaveBeenCalled(); + expect(inputsService.update).toHaveBeenCalled(); + expect(result).toBe(updateResponse); + }); + + it('batch input calls updatePlatformJobs', async () => { + const { service, inputsService, updatePlatformJobsSpy } = buildService(batchOldInput); + + await service.updatePipelineInput( + 'pipeline-id', + 'input-id', + updateInputDTO, + info, + user, + metadata, + ); + + expect(updatePlatformJobsSpy).toHaveBeenCalled(); + expect(inputsService.update).toHaveBeenCalled(); + }); +}); diff --git a/src/modules/pipelinesV2/pipelines.service.ts b/src/modules/pipelinesV2/pipelines.service.ts index b875cb7..992077f 100644 --- a/src/modules/pipelinesV2/pipelines.service.ts +++ b/src/modules/pipelinesV2/pipelines.service.ts @@ -433,6 +433,7 @@ export class PipelinesService implements OnModuleInit { }); this.logger.info('Update Dynamo Reference :' + JSON.stringify(oldInput)); + const isCdc = !!oldInput.plugin?.endsWith('_cdc'); const pipelineIdFormat = pipelineId.split('-').join('_'); const rollback: RollbackPromise[] = []; @@ -493,17 +494,21 @@ export class PipelinesService implements OnModuleInit { } } - 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"); + if (!isCdc) { + 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"); + } + } else { + this.logger.info('CDC input: skipping updatePlatformJobs (batch sync_mode/memory do not apply to CDC jobs)'); } return updateInputResponse; From c85e2ac5da52757afc29f063e73920bca5c575f5 Mon Sep 17 00:00:00 2001 From: Rafael Date: Mon, 17 Aug 2026 11:33:29 -0300 Subject: [PATCH 07/22] feat: forward create body config over gRPC (cdc.3) for CDC destinations The pipeline create body carries `config.tables[].destinations`, which the CDC path in pi-factory needs to honor a user-supplied raw Snowflake table name. The gRPC PipelineV2CreateRequest previously had no `config` field, so `...body` dropped it on the wire. Bump protospack to cdc.3 (adds optional `config` string). Serialize `body.config` into the create request the same way `properties` is handled, and add `config?` to the ICreatePipelineV2Req DTO so it's typed. Add a spec asserting the gRPC request carries a stringified config with the destination intact. Co-Authored-By: WOZCODE Claude-Session: https://claude.ai/code/session_01145m1zZMfx8RSJBxAhySdg --- package-lock.json | 8 ++-- package.json | 2 +- src/modules/pipelinesV2/interfaces.ts | 2 + .../pipelinesV2/pipelines.service.spec.ts | 47 +++++++++++++++++++ src/modules/pipelinesV2/pipelines.service.ts | 4 ++ 5 files changed, 58 insertions(+), 5 deletions(-) diff --git a/package-lock.json b/package-lock.json index 2cc5927..bef9f40 100644 --- a/package-lock.json +++ b/package-lock.json @@ -16,7 +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-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.2.tgz", + "@dadosfera/protospack-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.3.tgz", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", @@ -1735,9 +1735,9 @@ } }, "node_modules/@dadosfera/protospack-v2": { - "version": "3.40.0-beta.15-cdc.2", - "resolved": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.2.tgz", - "integrity": "sha512-FiZXH+bSv2Skxy9YJkczQOPnZxYbca5JlJSNbSzNxuk5e044i6Mp79LAA7SQDGt5yyAEfNHq2bAorHUxf0Q0TA==", + "version": "3.40.0-beta.15-cdc.3", + "resolved": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.3.tgz", + "integrity": "sha512-/JG35L/4iStITRhtk4GZnzKiF0oa6Hfc5i66rC7yDr2wG/MSEogQjtiU9G4LLJK84iAMpSstm6fdX7nWsjr5hQ==", "license": "ISC", "dependencies": { "@grpc/grpc-js": "^1.9.3", diff --git a/package.json b/package.json index 882ad39..53ec490 100644 --- a/package.json +++ b/package.json @@ -34,7 +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-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.2.tgz", + "@dadosfera/protospack-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.3.tgz", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", diff --git a/src/modules/pipelinesV2/interfaces.ts b/src/modules/pipelinesV2/interfaces.ts index 635ae80..c54e74a 100644 --- a/src/modules/pipelinesV2/interfaces.ts +++ b/src/modules/pipelinesV2/interfaces.ts @@ -31,6 +31,8 @@ export class IPipelineV2 { tags?: string[]; @ApiPropertyOptional() properties?: any; + @ApiPropertyOptional() + config?: any; @ApiProperty() connector_name: string; diff --git a/src/modules/pipelinesV2/pipelines.service.spec.ts b/src/modules/pipelinesV2/pipelines.service.spec.ts index 564639e..279f567 100644 --- a/src/modules/pipelinesV2/pipelines.service.spec.ts +++ b/src/modules/pipelinesV2/pipelines.service.spec.ts @@ -10,6 +10,7 @@ jest.mock('../platform-api/platform-api.service', () => ({ PlatformApiService: c jest.mock('src/services/nimbus/nimbus.service', () => ({ NimbusService: class {} })); jest.mock('../catalog/catalog.service', () => ({ CatalogService: class {} })); +import { of } from 'rxjs'; import { PipelinesService } from './pipelines.service'; const logger = { @@ -93,3 +94,49 @@ describe('PipelinesService - updatePipelineInput', () => { expect(inputsService.update).toHaveBeenCalled(); }); }); + +describe('PipelinesService - create', () => { + afterEach(() => jest.clearAllMocks()); + + // The body's `config` (carrying CDC destinations) must reach pi-factory as a + // JSON string — the gRPC proto field is a string, so an object would be + // stripped on the wire. Mirrors how `properties` is serialized. + it('serializes the body config into the gRPC create request', async () => { + const { service } = buildService(cdcOldInput); + + let captured: any; + (service as any).pipelineWriteService = { + PipelineV2Create: (req: any) => { + captured = req; + // The service does `lastValueFrom(...)`; return a real Observable. + return of({ pipeline: {} }); + }, + }; + + const body: any = { + name: 'p', + input_id: 'i', + transformations_ids: [], + tags: [], + properties: { schema: 'cadastros' }, + config: { + cron: '@once', + tables: [ + { + name: 'pedidos', + destinations: { + raw: { table_schema: 'PUBLIC', table_name: 'pedidos_001' }, + }, + }, + ], + }, + }; + + await service.create(body, metadata); + + expect(typeof captured.config).toBe('string'); + expect(JSON.parse(captured.config).tables[0].destinations.raw.table_name).toBe( + 'pedidos_001', + ); + }); +}); diff --git a/src/modules/pipelinesV2/pipelines.service.ts b/src/modules/pipelinesV2/pipelines.service.ts index 992077f..b3d641b 100644 --- a/src/modules/pipelinesV2/pipelines.service.ts +++ b/src/modules/pipelinesV2/pipelines.service.ts @@ -105,6 +105,10 @@ export class PipelinesService implements OnModuleInit { transformations_ids: body.transformations_ids, tags: body.tags, properties: body.properties && JSON.stringify(body.properties), + // JSON-serialize the create body config so it survives the gRPC wire + // (the proto field is a string). pi-factory's CDC path reads + // config.tables[].destinations to honor a user-supplied raw table name. + config: body.config && JSON.stringify(body.config), }; const createPipelineResponse = await lastValueFrom( From f1d539a912702d8ffc364046d450d31d7768e08c Mon Sep 17 00:00:00 2001 From: Rafael Date: Mon, 17 Aug 2026 13:35:19 -0300 Subject: [PATCH 08/22] fix: PipelineExecutionGuard allows edits when pipeline has no run history MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The guard read `status[status.length-1].last_status` on the pipeline_run history. For a pipeline that has never run — which is the permanent state of CDC pipelines, since they replicate continuously and never record batch runs — pipeline_run is empty, so the last element is undefined and `.last_status` threw, surfacing as "Error checking pipeline status: Cannot read properties of undefined (reading 'last_status')". This blocked Edit Objects (and delete-table) for every CDC pipeline. Treat an empty/statusless run history as "not running" and allow the edit (undefined never satisfied 'running' anyway). Also stop the catch from double-wrapping the deliberate is-running BadRequestException, so that rejection keeps its clear message; genuine status-check failures still fail closed with the wrapped message (safe default for a destructive-op gate). Adds a guard spec: empty history -> allow, statusless -> allow, not running -> allow, running -> is-running message (not wrapped), platform error -> wrapped message. Co-Authored-By: WOZCODE Claude-Session: https://claude.ai/code/session_01145m1zZMfx8RSJBxAhySdg --- src/guards/pipeline-execution.guard.spec.ts | 71 +++++++++++++++++++++ src/guards/pipeline-execution.guard.ts | 15 ++++- 2 files changed, 84 insertions(+), 2 deletions(-) create mode 100644 src/guards/pipeline-execution.guard.spec.ts diff --git a/src/guards/pipeline-execution.guard.spec.ts b/src/guards/pipeline-execution.guard.spec.ts new file mode 100644 index 0000000..0e2dbf9 --- /dev/null +++ b/src/guards/pipeline-execution.guard.spec.ts @@ -0,0 +1,71 @@ +import { BadRequestException } from '@nestjs/common'; +import { PipelineExecutionGuard } from './pipeline-execution.guard'; + +const logger = { info: jest.fn(), error: jest.fn() }; + +function buildGuard(proxyImpl: jest.Mock) { + const platformApiService: any = { proxy: proxyImpl }; + return new PipelineExecutionGuard( + { logger } as any, + platformApiService, + ); +} + +function contextWith(pipelineId = 'abc-123') { + return { + switchToHttp: () => ({ + getRequest: () => ({ params: { pipelineId }, user: {} }), + }), + } as any; +} + +describe('PipelineExecutionGuard', () => { + afterEach(() => jest.clearAllMocks()); + + it('allows the edit when the pipeline has no run history (empty array)', async () => { + const guard = buildGuard(jest.fn().mockResolvedValue([])); + await expect(guard.canActivate(contextWith())).resolves.toBe(true); + }); + + it('allows the edit when the last run has no last_status', async () => { + const guard = buildGuard(jest.fn().mockResolvedValue([{}])); + await expect(guard.canActivate(contextWith())).resolves.toBe(true); + }); + + it('allows the edit when the pipeline is not running', async () => { + const guard = buildGuard( + jest.fn().mockResolvedValue([{ last_status: 'SUCCEEDED' }]), + ); + await expect(guard.canActivate(contextWith())).resolves.toBe(true); + }); + + it('blocks with the is-running message when the pipeline is running', async () => { + const guard = buildGuard( + jest.fn().mockResolvedValue([{ last_status: 'RUNNING' }]), + ); + await expect(guard.canActivate(contextWith())).rejects.toThrow( + 'Pipeline is running, cannot update input now', + ); + }); + + it('wraps a genuine status-check failure (fail closed)', async () => { + const guard = buildGuard( + jest.fn().mockRejectedValue(new Error('platform down')), + ); + await expect(guard.canActivate(contextWith())).rejects.toThrow( + 'Error checking pipeline status: platform down', + ); + }); + + it('does not double-wrap the is-running BadRequestException', async () => { + const guard = buildGuard( + jest.fn().mockResolvedValue([{ last_status: 'running' }]), + ); + await expect(guard.canActivate(contextWith())).rejects.toBeInstanceOf( + BadRequestException, + ); + await expect(guard.canActivate(contextWith())).rejects.not.toThrow( + /Error checking pipeline status/, + ); + }); +}); diff --git a/src/guards/pipeline-execution.guard.ts b/src/guards/pipeline-execution.guard.ts index ffd651b..1b6ebcd 100644 --- a/src/guards/pipeline-execution.guard.ts +++ b/src/guards/pipeline-execution.guard.ts @@ -46,10 +46,16 @@ export class PipelineExecutionGuard implements CanActivate { user, ); - const currentStatus = status[status.length - 1] + const currentStatus = status?.[status.length - 1]; this.logger.info('Pipeline current status response:' + JSON.stringify(currentStatus)); - + + // No run history (e.g. CDC pipelines never record batch runs) means + // nothing is executing — allow the edit rather than crash on .last_status. + if (!currentStatus?.last_status) { + return true; + } + 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'); @@ -57,6 +63,11 @@ export class PipelineExecutionGuard implements CanActivate { return true; } } catch (error) { + // Preserve the deliberate is-running rejection; only wrap genuine + // status-check failures (fail closed on those for a destructive gate). + if (error instanceof BadRequestException) { + throw error; + } this.logger.error('Error in PipelineExecutionGuard: ' + error.message); throw new BadRequestException('Error checking pipeline status: ' + error.message); } From fab2061efc9072bf94af0f9d4f8ac30c73c85d98 Mon Sep 17 00:00:00 2001 From: Rafael Date: Mon, 17 Aug 2026 16:16:32 -0300 Subject: [PATCH 09/22] chore: bump protospack-v2 to cdc.4 Co-Authored-By: WOZCODE Claude-Session: https://claude.ai/code/session_01145m1zZMfx8RSJBxAhySdg --- package-lock.json | 8 ++++---- package.json | 2 +- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/package-lock.json b/package-lock.json index bef9f40..8b9b72b 100644 --- a/package-lock.json +++ b/package-lock.json @@ -16,7 +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-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.3.tgz", + "@dadosfera/protospack-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.4.tgz", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", @@ -1735,9 +1735,9 @@ } }, "node_modules/@dadosfera/protospack-v2": { - "version": "3.40.0-beta.15-cdc.3", - "resolved": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.3.tgz", - "integrity": "sha512-/JG35L/4iStITRhtk4GZnzKiF0oa6Hfc5i66rC7yDr2wG/MSEogQjtiU9G4LLJK84iAMpSstm6fdX7nWsjr5hQ==", + "version": "3.40.0-beta.15-cdc.4", + "resolved": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.4.tgz", + "integrity": "sha512-jREB4KF76Z7holGByXyys/RC0ryu3doWikYNgzO7mnzrz5kWYi7U2S/9FF3iDZVqPmK5flYg4tlC2WR79c2ztA==", "license": "ISC", "dependencies": { "@grpc/grpc-js": "^1.9.3", diff --git a/package.json b/package.json index 53ec490..2926ef6 100644 --- a/package.json +++ b/package.json @@ -34,7 +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-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.3.tgz", + "@dadosfera/protospack-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.4.tgz", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", From 55308420ff8750e63271468dd574d1b65e9356c9 Mon Sep 17 00:00:00 2001 From: Rafael Date: Mon, 17 Aug 2026 17:01:05 -0300 Subject: [PATCH 10/22] =?UTF-8?q?feat(cdc):=20addTable=20route=20=E2=80=94?= =?UTF-8?q?=20DynamoDB-first=20+=20rollback?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docsfera.json | 53 +++++++++++ src/modules/inputs/inputs.service.ts | 8 ++ src/modules/pipelinesV2/pipelines.module.ts | 4 +- src/modules/pipelinesV2/pipelines.service.ts | 6 +- .../platform-api.controller.spec.ts | 94 +++++++++++++++++++ .../platform-api/platform-api.controller.ts | 68 ++++++++++++++ .../platform-api/platform-api.module.ts | 12 ++- 7 files changed, 240 insertions(+), 5 deletions(-) diff --git a/docsfera.json b/docsfera.json index ff2d395..da23ed4 100644 --- a/docsfera.json +++ b/docsfera.json @@ -5205,6 +5205,53 @@ ] } }, + "/platform/pipelines/{pipelineId}/inputs/{inputId}/tables": { + "post": { + "operationId": "PlatformApiController_addTable", + "summary": "Add a CDC table to an input and dispatch its platform jobs", + "parameters": [ + { + "name": "pipelineId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + }, + { + "name": "inputId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "201": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + }, + { + "access-token": [] + } + ] + } + }, "/platform/jobs/jdbc/configs/allowed_datatypes": { "get": { "operationId": "PlatformApiController_getJdbcAllowedDatatypes", @@ -10921,6 +10968,9 @@ "properties": { "type": "object" }, + "config": { + "type": "object" + }, "connector_name": { "type": "string" }, @@ -10985,6 +11035,9 @@ "properties": { "type": "object" }, + "config": { + "type": "object" + }, "connector_name": { "type": "string" }, diff --git a/src/modules/inputs/inputs.service.ts b/src/modules/inputs/inputs.service.ts index 29048ed..ad31e1c 100644 --- a/src/modules/inputs/inputs.service.ts +++ b/src/modules/inputs/inputs.service.ts @@ -322,4 +322,12 @@ export class InputsService { async unmarkTableDeleted(data: { input_id: string; table_name: string; info: Info }) { return lastValueFrom((this.inputWriteService as any).UnmarkTableDeleted(data)); } + + async addCdcTable(data: { client_id?: string; id: string; table: any; info: Info }) { + return lastValueFrom(this.inputWriteService.AddCdcTable(data as any)); + } + + async removeCdcTable(data: { client_id?: string; id: string; table_name: string; info: Info }) { + return lastValueFrom((this.inputWriteService as any).RemoveCdcTable(data)); + } } diff --git a/src/modules/pipelinesV2/pipelines.module.ts b/src/modules/pipelinesV2/pipelines.module.ts index b4f24af..efd3c74 100644 --- a/src/modules/pipelinesV2/pipelines.module.ts +++ b/src/modules/pipelinesV2/pipelines.module.ts @@ -1,4 +1,4 @@ -import { Module } from '@nestjs/common'; +import { Module, forwardRef } from '@nestjs/common'; import { ClientsModule } from '@nestjs/microservices'; import { DadosferaLogger } from '@dadosfera/dadosfera-logs'; @@ -23,7 +23,7 @@ const client = new PipelinesClientConfiguration(); ConnectorModule, InputsModule, TransformationsModule, - PlatformApiModule, + forwardRef(() => PlatformApiModule), NimbusServicesModule, CatalogModule ], diff --git a/src/modules/pipelinesV2/pipelines.service.ts b/src/modules/pipelinesV2/pipelines.service.ts index b3d641b..77c8059 100644 --- a/src/modules/pipelinesV2/pipelines.service.ts +++ b/src/modules/pipelinesV2/pipelines.service.ts @@ -19,7 +19,7 @@ import { lastValueFrom } from 'rxjs'; import { DadosferaLogger } from '@dadosfera/dadosfera-logs'; import { PipelinesClientConfiguration } from './pipelines-client'; import { ICreatePipelineV2Req, IIdRequest, UpdatePlatformInputRequest, UpdateTableDTO } from './interfaces'; -import { PipelineV2CreateRequest } from '@dadosfera/protospack-v2/dist/lib/PipelineV2/interfaces/messages'; +import { PipelineV2CreateRequest, AddCdcJobsRequest, AddCdcJobsResponse } from '@dadosfera/protospack-v2/dist/lib/PipelineV2/interfaces/messages'; import { Metadata } from '@grpc/grpc-js'; import { ConnectorClientService } from '../connector/client.service'; import { InputsService } from '../inputs/inputs.service'; @@ -727,6 +727,10 @@ export class PipelinesService implements OnModuleInit { return statusPipelineResponse; } + async addCdcJobs(data: AddCdcJobsRequest): Promise { + return lastValueFrom(this.pipelineWriteService.AddCdcJobs(data)); + } + async runPipeline({ id, info }: IIdRequest) { this.logger.info('PipelinesClientService - RunPipeline'); const statusPipelineResponse = await lastValueFrom( diff --git a/src/modules/platform-api/platform-api.controller.spec.ts b/src/modules/platform-api/platform-api.controller.spec.ts index c235dd3..25077a1 100644 --- a/src/modules/platform-api/platform-api.controller.spec.ts +++ b/src/modules/platform-api/platform-api.controller.spec.ts @@ -5,6 +5,7 @@ jest.mock('../customers/customers.service', () => ({ CustomersService: class {} })); jest.mock('../catalog/catalog.service', () => ({ CatalogService: class {} })); jest.mock('../inputs/inputs.service', () => ({ InputsService: class {} })); +jest.mock('../pipelinesV2/pipelines.service', () => ({ PipelinesService: class {} })); import { Test, TestingModule } from '@nestjs/testing'; import { DadosferaLogger } from '@dadosfera/dadosfera-logs'; @@ -16,6 +17,7 @@ import { DynamoDBService } from '../../services/dynamodb'; import { CustomersService } from '../customers/customers.service'; import { CatalogService } from '../catalog/catalog.service'; import { InputsService } from '../inputs/inputs.service'; +import { PipelinesService } from '../pipelinesV2/pipelines.service'; const logger = { info: (...args) => args, @@ -49,6 +51,7 @@ describe('PlatformApiController - deleteTable', () => { { provide: CustomersService, useValue: {} }, { provide: CatalogService, useValue: {} }, { provide: InputsService, useValue: inputsService }, + { provide: PipelinesService, useValue: {} }, { provide: DadosferaLogger, useValue: { logger } }, ], }).compile(); @@ -105,3 +108,94 @@ describe('PlatformApiController - deleteTable', () => { expect(reconfiguredConnector).toBe(false); }); }); + +describe('PlatformApiController - addTable', () => { + let controller: PlatformApiController; + let inputsService: { addCdcTable: jest.Mock; removeCdcTable: jest.Mock }; + let pipelinesClientService: { addCdcJobs: jest.Mock }; + + const body = { + table_name: 'orders', + table_schema: 'public', + primary_keys: ['id'], + destinations: { + raw: { table_schema: 'raw', table_name: 'orders' }, + qualify: { table_schema: 'qualify', table_name: 'orders' }, + }, + }; + + beforeEach(async () => { + inputsService = { + addCdcTable: jest.fn(), + removeCdcTable: jest.fn(), + }; + pipelinesClientService = { + addCdcJobs: jest.fn(), + }; + + const module: TestingModule = await Test.createTestingModule({ + controllers: [PlatformApiController], + providers: [ + { provide: PlatformApiService, useValue: { proxy: jest.fn() } }, + { provide: ElasticsearchService, useValue: {} }, + { provide: DynamoDBService, useValue: {} }, + { provide: CustomersService, useValue: {} }, + { provide: CatalogService, useValue: {} }, + { provide: InputsService, useValue: inputsService }, + { provide: PipelinesService, useValue: pipelinesClientService }, + { provide: DadosferaLogger, useValue: { logger } }, + ], + }).compile(); + + controller = module.get(PlatformApiController); + }); + + it('addTable: DynamoDB append then platform AddJobs, returns job_ids', async () => { + inputsService.addCdcTable.mockResolvedValue({ input: {} }); + pipelinesClientService.addCdcJobs.mockResolvedValue({ job_ids: ['p_2'], skipped: [] }); + + const result = await controller.addTable('pid', 'iid', body, mockUser); + + expect(inputsService.addCdcTable).toHaveBeenCalledWith({ + id: 'iid', + table: { + table_schema: 'public', + table_name: 'orders', + primary_keys: ['id'], + name: 'orders', + }, + info: { customer_id: 'c1', customer: 'cust', user_id: 'u1' }, + }); + expect(pipelinesClientService.addCdcJobs).toHaveBeenCalledWith({ + pipeline_id: 'pid', + input_id: 'iid', + tables: [{ + table_schema: 'public', + table_name: 'orders', + primary_keys: ['id'], + destinations: body.destinations, + }], + info: { customer_id: 'c1', user_id: 'u1', customer: 'cust' }, + }); + expect(result).toEqual({ job_ids: ['p_2'], skipped: [] }); + expect(inputsService.removeCdcTable).not.toHaveBeenCalled(); + + const addCdcTableOrder = inputsService.addCdcTable.mock.invocationCallOrder[0]; + const addCdcJobsOrder = pipelinesClientService.addCdcJobs.mock.invocationCallOrder[0]; + expect(addCdcTableOrder).toBeLessThan(addCdcJobsOrder); + }); + + it('addTable: rolls back the DynamoDB row when AddJobs fails', async () => { + inputsService.addCdcTable.mockResolvedValue({ input: {} }); + pipelinesClientService.addCdcJobs.mockRejectedValue(new Error('platform down')); + inputsService.removeCdcTable.mockResolvedValue({}); + + await expect(controller.addTable('pid', 'iid', body, mockUser)).rejects.toThrow('platform down'); + + expect(inputsService.removeCdcTable).toHaveBeenCalledWith({ + id: 'iid', + table_name: 'orders', + info: { customer_id: 'c1', customer: 'cust', user_id: 'u1' }, + }); + }); +}); diff --git a/src/modules/platform-api/platform-api.controller.ts b/src/modules/platform-api/platform-api.controller.ts index b10b09e..0a6c6cc 100644 --- a/src/modules/platform-api/platform-api.controller.ts +++ b/src/modules/platform-api/platform-api.controller.ts @@ -33,6 +33,7 @@ import { CatalogService } from '../catalog/catalog.service'; import { PackTheMetadata } from '../../utils/PackTheMetadata'; import { ValidationTableDTO } from './platform-api.dto'; import { InputsService } from '../inputs/inputs.service'; +import { PipelinesService } from '../pipelinesV2/pipelines.service'; import { PipelineExecutionGuard } from 'src/guards/pipeline-execution.guard'; @@ -61,6 +62,7 @@ export class PlatformApiController { private readonly customersService: CustomersService, private readonly catalogService: CatalogService, private readonly inputsService: InputsService, + private readonly pipelinesClientService: PipelinesService, @Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger, ) { this.logger = dadosferaLogger.logger; @@ -990,6 +992,72 @@ export class PlatformApiController { } } + @Post('pipelines/:pipelineId/inputs/:inputId/tables') + @ApiOperation({ summary: 'Add a CDC table to an input and dispatch its platform jobs' }) + @RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE) + @UseGuards(PipelineExecutionGuard) + async addTable( + @Param('pipelineId') pipelineId: string, + @Param('inputId') inputId: string, + @Body() body: { + table_name: string; + table_schema: string; + primary_keys: string[]; + destinations: { + raw: { table_schema: string; table_name: string }; + qualify: { table_schema: string; table_name: string }; + }; + }, + @User() user: RequestUser, + ) { + const info = { + customer_id: user.customer_id, + customer: user.customer_name, + user_id: user.user_id, + }; + + const cdcTable = { + table_schema: body.table_schema, + table_name: body.table_name, + primary_keys: body.primary_keys, + name: body.table_name, + }; + + this.logger.info('addTable: appending CDC table to DynamoDB', { inputId, tableName: body.table_name }); + await this.inputsService.addCdcTable({ id: inputId, table: cdcTable, info }); + this.logger.info('addTable: DynamoDB row appended', { inputId, tableName: body.table_name }); + + try { + const customInfo = { + customer_id: user.customer_id, + user_id: user.user_id, + customer: user.customer_name, + }; + + const res = await this.pipelinesClientService.addCdcJobs({ + pipeline_id: pipelineId, + input_id: inputId, + tables: [{ + table_schema: body.table_schema, + table_name: body.table_name, + primary_keys: body.primary_keys, + destinations: body.destinations, + }], + info: customInfo, + }); + + return res; + } catch (error) { + this.logger.error('addTable: platform AddJobs failed, rolling back the DynamoDB row', { tableName: body.table_name, error: error.message }); + try { + await this.inputsService.removeCdcTable({ id: inputId, table_name: body.table_name, info }); + } catch (rbErr) { + this.logger.error('addTable: rollback failed', { error: rbErr.message }); + } + throw error; + } + } + // ==================== JOBS - JDBC SYNC MODE ROUTES ==================== @Get('jobs/jdbc/configs/allowed_datatypes') diff --git a/src/modules/platform-api/platform-api.module.ts b/src/modules/platform-api/platform-api.module.ts index b2b0726..726e743 100644 --- a/src/modules/platform-api/platform-api.module.ts +++ b/src/modules/platform-api/platform-api.module.ts @@ -1,4 +1,4 @@ -import { Module } from '@nestjs/common'; +import { Module, forwardRef } from '@nestjs/common'; import { DadosferaLogger } from '@dadosfera/dadosfera-logs'; @@ -9,9 +9,17 @@ import { DynamoDBModule } from '../../services/dynamodb'; import { CustomersModule } from '../customers/customers.module'; import { CatalogModule } from '../catalog/catalog.module'; import { InputsModule } from '../inputs/inputs.module'; +import { PipelinesV2Module } from '../pipelinesV2/pipelines.module'; @Module({ - imports: [ElasticsearchModule, DynamoDBModule, CustomersModule, CatalogModule, InputsModule], + imports: [ + ElasticsearchModule, + DynamoDBModule, + CustomersModule, + CatalogModule, + InputsModule, + forwardRef(() => PipelinesV2Module), + ], controllers: [PlatformApiController], providers: [PlatformApiService, DadosferaLogger], exports: [PlatformApiService], From 714c1334b9d4ff52b2ad339a78111a902798d10f Mon Sep 17 00:00:00 2001 From: Rafael Date: Wed, 19 Aug 2026 09:49:43 -0300 Subject: [PATCH 11/22] chore(deps): use published protospack-v2 3.40.0-beta.16 Switch @dadosfera/protospack-v2 from the local file: tarball (removed from the protospack-v2 repo) to the CI-published CodeArtifact version. Installs and builds clean; test suites pass. Co-Authored-By: WOZCODE --- package-lock.json | 8 ++++---- package.json | 2 +- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/package-lock.json b/package-lock.json index 8b9b72b..c95fa11 100644 --- a/package-lock.json +++ b/package-lock.json @@ -16,7 +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-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.4.tgz", + "@dadosfera/protospack-v2": "3.40.0-beta.16", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", @@ -1735,9 +1735,9 @@ } }, "node_modules/@dadosfera/protospack-v2": { - "version": "3.40.0-beta.15-cdc.4", - "resolved": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.4.tgz", - "integrity": "sha512-jREB4KF76Z7holGByXyys/RC0ryu3doWikYNgzO7mnzrz5kWYi7U2S/9FF3iDZVqPmK5flYg4tlC2WR79c2ztA==", + "version": "3.40.0-beta.16", + "resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.40.0-beta.16.tgz", + "integrity": "sha512-XWQ4zV4Gx2s1nSGVqOe4CkY6j4oRHdajO0ztZkMkaC0l48+qxpBlFfpcf6rvTd1d339T3MaL0rzZOI0YKSHkFA==", "license": "ISC", "dependencies": { "@grpc/grpc-js": "^1.9.3", diff --git a/package.json b/package.json index 2926ef6..d97ac62 100644 --- a/package.json +++ b/package.json @@ -34,7 +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-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.40.0-beta.15-cdc.4.tgz", + "@dadosfera/protospack-v2": "3.40.0-beta.16", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", From ead3fa12dc080924be1c5f316c22b657fb504d8b Mon Sep 17 00:00:00 2001 From: Rafael Date: Wed, 19 Aug 2026 10:47:06 -0300 Subject: [PATCH 12/22] =?UTF-8?q?feat(cdc):=20batch=20table=20removal=20?= =?UTF-8?q?=E2=80=94=20reconfigure=20connectors=20once?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Removing N tables via Edit Objects previously looped a single-table DELETE per table (maestro deleteTable hardcodes job_ids:[one]), so the Debezium source + Snowflake sink connectors were rewritten/restarted once per table. platform-api's DELETE /pipeline/:id/jobs already batches (2 connector writes total for any N), but nothing above it used the array. New maestro DELETE /pipelines/:pipelineId/inputs/:inputId/tables takes { table_names: [] }: soft-deletes each in DynamoDB (tracking successes), resolves all table_names -> job_ids from the platform pipeline in one GET, then makes ONE DELETE /pipeline/:id/jobs with all job_ids. All-or-nothing: any failure (a later mark, an unmatched table, or the platform delete) rolls back only the marks made in this call. The single-table deleteTable route is kept (unchanged) — nothing else depends on removing it, and that's a separable cleanup. Tests: N tables -> one platform DELETE with all job_ids and no per-job call; rollback on platform failure; rollback + no delete when a later mark fails; 404 for an unmatched table. 8 controller specs pass; maestro builds. Co-Authored-By: WOZCODE --- .../platform-api.controller.spec.ts | 110 ++++++++++++++++++ .../platform-api/platform-api.controller.ts | 69 +++++++++++ 2 files changed, 179 insertions(+) diff --git a/src/modules/platform-api/platform-api.controller.spec.ts b/src/modules/platform-api/platform-api.controller.spec.ts index 25077a1..ae4c352 100644 --- a/src/modules/platform-api/platform-api.controller.spec.ts +++ b/src/modules/platform-api/platform-api.controller.spec.ts @@ -199,3 +199,113 @@ describe('PlatformApiController - addTable', () => { }); }); }); + + +describe('PlatformApiController - deleteTables (batch)', () => { + let controller: PlatformApiController; + let platformApiService: { proxy: jest.Mock }; + let inputsService: { markTableDeleted: jest.Mock; unmarkTableDeleted: jest.Mock }; + + beforeEach(async () => { + platformApiService = { proxy: jest.fn() }; + inputsService = { + markTableDeleted: jest.fn().mockResolvedValue({ is_deleted: true, deleted_at: 't' }), + unmarkTableDeleted: jest.fn().mockResolvedValue({}), + }; + + const module: TestingModule = await Test.createTestingModule({ + controllers: [PlatformApiController], + providers: [ + { provide: PlatformApiService, useValue: platformApiService }, + { provide: ElasticsearchService, useValue: {} }, + { provide: DynamoDBService, useValue: {} }, + { provide: CustomersService, useValue: {} }, + { provide: CatalogService, useValue: {} }, + { provide: InputsService, useValue: inputsService }, + { provide: PipelinesService, useValue: {} }, + { provide: DadosferaLogger, useValue: { logger } }, + ], + }).compile(); + + controller = module.get(PlatformApiController); + }); + + const pipelineWithJobs = () => ({ + jobs: [ + { job_id: 'p_0', input: { connector: 'cdc', table_name: 'pedidos' } }, + { job_id: 'p_1', input: { connector: 'cdc', table_name: 'clientes' } }, + { job_id: 'p_2', input: { connector: 'cdc', table_name: 'produtos' } }, + ], + }); + + it('removes N tables in ONE platform call (connectors reconfigured once)', async () => { + platformApiService.proxy.mockImplementation((method: string) => + Promise.resolve(method === 'GET' ? pipelineWithJobs() : {}), + ); + + await controller.deleteTables('pid', 'iid', { table_names: ['pedidos', 'produtos'] }, mockUser); + + // one mark per table + expect(inputsService.markTableDeleted).toHaveBeenCalledTimes(2); + + // exactly one DELETE to the batch endpoint, with BOTH job_ids + const deleteCalls = platformApiService.proxy.mock.calls.filter( + ([m, p]: any[]) => m === 'DELETE' && p === '/pipeline/pid/jobs', + ); + expect(deleteCalls).toHaveLength(1); + expect(deleteCalls[0][3]).toEqual({ + job_ids: ['p_0', 'p_2'], + delete_snowflake_tables: false, + }); + // never the per-job route + const perJob = platformApiService.proxy.mock.calls.some( + ([m, p]: any[]) => m === 'DELETE' && String(p).startsWith('/jobs/'), + ); + expect(perJob).toBe(false); + }); + + it('rolls back only this call\'s marks when the platform delete fails', async () => { + platformApiService.proxy.mockImplementation((method: string) => { + if (method === 'GET') return Promise.resolve(pipelineWithJobs()); + return Promise.reject(new Error('platform boom')); + }); + + await expect( + controller.deleteTables('pid', 'iid', { table_names: ['pedidos', 'clientes'] }, mockUser), + ).rejects.toThrow(); + + // both marks rolled back, nothing else + expect(inputsService.unmarkTableDeleted).toHaveBeenCalledTimes(2); + const unmarked = inputsService.unmarkTableDeleted.mock.calls.map((c: any[]) => c[0].table_name).sort(); + expect(unmarked).toEqual(['clientes', 'pedidos']); + }); + + it('rolls back the marks made so far if a later mark fails (atomic)', async () => { + // second mark fails → first must be rolled back, no platform delete attempted + inputsService.markTableDeleted + .mockResolvedValueOnce({ is_deleted: true, deleted_at: 't' }) + .mockRejectedValueOnce(new Error('dynamo boom')); + + await expect( + controller.deleteTables('pid', 'iid', { table_names: ['pedidos', 'clientes'] }, mockUser), + ).rejects.toThrow(); + + expect(inputsService.unmarkTableDeleted).toHaveBeenCalledTimes(1); + expect(inputsService.unmarkTableDeleted.mock.calls[0][0].table_name).toBe('pedidos'); + // never reached the platform delete + const attemptedDelete = platformApiService.proxy.mock.calls.some(([m]: any[]) => m === 'DELETE'); + expect(attemptedDelete).toBe(false); + }); + + it('404s when a requested table has no matching job', async () => { + platformApiService.proxy.mockImplementation((method: string) => + Promise.resolve(method === 'GET' ? pipelineWithJobs() : {}), + ); + + await expect( + controller.deleteTables('pid', 'iid', { table_names: ['pedidos', 'ghost'] }, mockUser), + ).rejects.toThrow(); + // the successful mark (pedidos) must be rolled back + expect(inputsService.unmarkTableDeleted).toHaveBeenCalled(); + }); +}); diff --git a/src/modules/platform-api/platform-api.controller.ts b/src/modules/platform-api/platform-api.controller.ts index 0a6c6cc..e547b0d 100644 --- a/src/modules/platform-api/platform-api.controller.ts +++ b/src/modules/platform-api/platform-api.controller.ts @@ -992,6 +992,75 @@ export class PlatformApiController { } } + @Delete('pipelines/:pipelineId/inputs/:inputId/tables') + @ApiOperation({ summary: 'Batch-remove tables from an input; reconfigures the CDC connectors once' }) + @RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE) + @UseGuards(PipelineExecutionGuard) + async deleteTables( + @Param('pipelineId') pipelineId: string, + @Param('inputId') inputId: string, + @Body() body: { table_names: string[] }, + @User() user: RequestUser, + ) { + const tableNames = body.table_names ?? []; + if (tableNames.length === 0) { + throw new BadRequestException('table_names must be a non-empty array'); + } + const info = { + customer_id: user.customer_id, + customer: user.customer_name, + user_id: user.user_id, + }; + + // 1) Soft-delete each table in DynamoDB, tracking which succeeded so a later + // failure only rolls back the marks made in THIS call. + const marked: string[] = []; + const rollback = async () => { + for (const name of marked) { + try { + await this.inputsService.unmarkTableDeleted({ input_id: inputId, table_name: name, info }); + } catch (rollbackError) { + this.logger.error('deleteTables: rollback failed', { tableName: name, error: rollbackError.message }); + } + } + }; + + try { + for (const name of tableNames) { + await this.inputsService.markTableDeleted({ input_id: inputId, table_name: name, info }); + marked.push(name); + } + + // 2) Resolve all table_names -> job_ids from the platform pipeline (one GET). + const normalizedPipelineId = this.normalizePipelineId(pipelineId); + const platformPipeline = await this.platformApiService.proxy('GET', `/pipeline/${normalizedPipelineId}`, user); + const jobs = platformPipeline?.jobs ?? []; + + const jobIds: string[] = []; + for (const name of tableNames) { + const job = jobs.find((j: any) => j.input?.table_name === name); + if (!job) throw new NotFoundException(`Job for table '${name}' not found in pipeline`); + jobIds.push(job.job_id); + } + + // 3) Remove them all in ONE platform call so the Debezium source + Snowflake + // sink connectors are reconfigured a single time, not once per table. + await this.platformApiService.proxy( + 'DELETE', + `/pipeline/${normalizedPipelineId}/jobs`, + user, + { job_ids: jobIds, delete_snowflake_tables: false }, + ); + this.logger.info('deleteTables: jobs deleted', { jobIds }); + + return { table_names: tableNames, deleted: true }; + } catch (error) { + this.logger.error('deleteTables: failed, rolling back this call\'s marks', { tableNames, error: error.message }); + await rollback(); + throw error; + } + } + @Post('pipelines/:pipelineId/inputs/:inputId/tables') @ApiOperation({ summary: 'Add a CDC table to an input and dispatch its platform jobs' }) @RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE) From 997eef876dcd647288f42e85c0725dd6dbe95bb3 Mon Sep 17 00:00:00 2001 From: Rafael Date: Fri, 21 Aug 2026 11:10:24 -0300 Subject: [PATCH 13/22] CHORE(cdc): point protospack at local tarball 3.41.0-cdc-iceberg.0 (dev-only) Unblocks the CdcDestination contract locally on the CDC branch. Swap to a published registry version before merge (CI guard blocks file: deps). Co-Authored-By: WOZCODE --- package-lock.json | 8 ++++---- package.json | 2 +- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/package-lock.json b/package-lock.json index c95fa11..1ae57c4 100644 --- a/package-lock.json +++ b/package-lock.json @@ -16,7 +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-v2": "3.40.0-beta.16", + "@dadosfera/protospack-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.0.tgz", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", @@ -1735,9 +1735,9 @@ } }, "node_modules/@dadosfera/protospack-v2": { - "version": "3.40.0-beta.16", - "resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.40.0-beta.16.tgz", - "integrity": "sha512-XWQ4zV4Gx2s1nSGVqOe4CkY6j4oRHdajO0ztZkMkaC0l48+qxpBlFfpcf6rvTd1d339T3MaL0rzZOI0YKSHkFA==", + "version": "3.41.0-cdc-iceberg.0", + "resolved": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.0.tgz", + "integrity": "sha512-qPs0LukP5BJnybl3JQqqbYQ9oi5SJ86+4A3zd9zisrqqzFHv3o9c0CSwcimwyT3CWm/147vq5LMQy0ulal0Ung==", "license": "ISC", "dependencies": { "@grpc/grpc-js": "^1.9.3", diff --git a/package.json b/package.json index d97ac62..17c03bf 100644 --- a/package.json +++ b/package.json @@ -34,7 +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-v2": "3.40.0-beta.16", + "@dadosfera/protospack-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.0.tgz", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", From 7de207c676b4467306ff0625b04eff8e5de1d488 Mon Sep 17 00:00:00 2001 From: Rafael Date: Fri, 21 Aug 2026 11:13:27 -0300 Subject: [PATCH 14/22] FEAT(cdc): carry iceberg destination through /inputs/cdc DTOs Co-Authored-By: WOZCODE --- src/modules/inputs/dtos/input.model.ts | 14 ++++ src/modules/inputs/inputs.controller.spec.ts | 75 ++++++++++++++++++++ src/modules/inputs/inputs.service.ts | 2 + 3 files changed, 91 insertions(+) create mode 100644 src/modules/inputs/inputs.controller.spec.ts diff --git a/src/modules/inputs/dtos/input.model.ts b/src/modules/inputs/dtos/input.model.ts index fa73fd1..259d756 100644 --- a/src/modules/inputs/dtos/input.model.ts +++ b/src/modules/inputs/dtos/input.model.ts @@ -73,6 +73,18 @@ export class CdcTableReq { table_schema?: string; @ApiPropertyOptional({ type: [String] }) primary_keys?: string[]; + @ApiPropertyOptional() + iceberg_table_name?: string; +} + +export class IcebergDestinationReq { + @ApiProperty() + namespace: string; +} + +export class CdcDestinationReq { + @ApiPropertyOptional({ type: IcebergDestinationReq }) + iceberg?: IcebergDestinationReq; } export class CreateCdcInputReq { @@ -84,4 +96,6 @@ export class CreateCdcInputReq { tables: CdcTableReq[]; @ApiPropertyOptional() read_only?: boolean; + @ApiPropertyOptional({ type: CdcDestinationReq }) + destination?: CdcDestinationReq; } diff --git a/src/modules/inputs/inputs.controller.spec.ts b/src/modules/inputs/inputs.controller.spec.ts new file mode 100644 index 0000000..241918f --- /dev/null +++ b/src/modules/inputs/inputs.controller.spec.ts @@ -0,0 +1,75 @@ +import { Test, TestingModule } from '@nestjs/testing'; +import { InputsController } from './inputs.controller'; +import { InputsService } from './inputs.service'; +import DadosferaLogger from '@dadosfera/dadosfera-logs'; +import { CreateCdcInputReq } from './dtos/input.model'; +import { RequestUser } from 'src/decorators/user.decorator'; + +describe('InputsController', () => { + let controller: InputsController; + let inputsService: { createCdc: jest.Mock }; + + beforeEach(async () => { + inputsService = { + createCdc: jest.fn().mockResolvedValue({ input: {} }), + }; + + const module: TestingModule = await Test.createTestingModule({ + controllers: [InputsController], + providers: [ + { + provide: DadosferaLogger, + useValue: { logger: { info: jest.fn() } }, + }, + { + provide: InputsService, + useValue: inputsService, + }, + ], + }).compile(); + + controller = module.get(InputsController); + }); + + it('should be defined', () => { + expect(controller).toBeDefined(); + }); + + it('forwards destination.iceberg.namespace to InputsService.createCdc', async () => { + const body: CreateCdcInputReq = { + name: 'my-cdc-input', + plugin: 'mysql_cdc', + tables: [ + { + name: 'orders', + table_schema: 'public', + iceberg_table_name: 'orders_iceberg', + }, + ], + destination: { + iceberg: { + namespace: 'my_namespace', + }, + }, + }; + const user: RequestUser = { + user_id: 'user-1', + customer_id: 'customer-1', + customer_name: 'customer', + } as RequestUser; + + await controller.createCdc(body, user); + + expect(inputsService.createCdc).toHaveBeenCalledWith( + expect.objectContaining({ + body: expect.objectContaining({ + destination: { + iceberg: { + namespace: 'my_namespace', + }, + }, + }), + }), + ); + }); +}); diff --git a/src/modules/inputs/inputs.service.ts b/src/modules/inputs/inputs.service.ts index ad31e1c..c221520 100644 --- a/src/modules/inputs/inputs.service.ts +++ b/src/modules/inputs/inputs.service.ts @@ -197,7 +197,9 @@ export class InputsService { table_name: t.name, name: t.name, // canonical identity == table_name (in-factory also backfills) primary_keys: t.primary_keys ?? [], + iceberg_table_name: t.iceberg_table_name, })), + destination: body.destination, }, info, }; From b3bbf2473aeeac7b2507923b90b4bdc2af24fb27 Mon Sep 17 00:00:00 2001 From: Rafael Date: Fri, 21 Aug 2026 11:29:54 -0300 Subject: [PATCH 15/22] FEAT(cdc): carry iceberg_table_name through the add-tables endpoint The POST /pipelines/:pipelineId/inputs/:inputId/tables route built its CdcTable payload field-by-field and silently dropped iceberg_table_name even though inputsService.addCdcTable/the gRPC AddCdcTable call (and the protospack CdcTable message) already support it. Widen the inline request body type and thread the field into the addCdcTable payload; absent for snowflake, unchanged back-compat. Co-Authored-By: WOZCODE --- .../platform-api.controller.spec.ts | 21 +++++++++++++++++++ .../platform-api/platform-api.controller.ts | 4 ++++ 2 files changed, 25 insertions(+) diff --git a/src/modules/platform-api/platform-api.controller.spec.ts b/src/modules/platform-api/platform-api.controller.spec.ts index ae4c352..57127ca 100644 --- a/src/modules/platform-api/platform-api.controller.spec.ts +++ b/src/modules/platform-api/platform-api.controller.spec.ts @@ -185,6 +185,27 @@ describe('PlatformApiController - addTable', () => { expect(addCdcTableOrder).toBeLessThan(addCdcJobsOrder); }); + it('addTable: carries iceberg_table_name on the added table through to AddCdcTable', async () => { + inputsService.addCdcTable.mockResolvedValue({ input: {} }); + pipelinesClientService.addCdcJobs.mockResolvedValue({ job_ids: ['p_3'], skipped: [] }); + + const icebergBody = { ...body, iceberg_table_name: 'cdc_raw.public__orders' }; + + await controller.addTable('pid', 'iid', icebergBody, mockUser); + + expect(inputsService.addCdcTable).toHaveBeenCalledWith({ + id: 'iid', + table: { + table_schema: 'public', + table_name: 'orders', + primary_keys: ['id'], + name: 'orders', + iceberg_table_name: 'cdc_raw.public__orders', + }, + info: { customer_id: 'c1', customer: 'cust', user_id: 'u1' }, + }); + }); + it('addTable: rolls back the DynamoDB row when AddJobs fails', async () => { inputsService.addCdcTable.mockResolvedValue({ input: {} }); pipelinesClientService.addCdcJobs.mockRejectedValue(new Error('platform down')); diff --git a/src/modules/platform-api/platform-api.controller.ts b/src/modules/platform-api/platform-api.controller.ts index e547b0d..8355ffa 100644 --- a/src/modules/platform-api/platform-api.controller.ts +++ b/src/modules/platform-api/platform-api.controller.ts @@ -1076,6 +1076,9 @@ export class PlatformApiController { raw: { table_schema: string; table_name: string }; qualify: { table_schema: string; table_name: string }; }; + // Iceberg destination only (protospack CdcTable.iceberg_table_name); + // absent for snowflake, back-compat. + iceberg_table_name?: string; }, @User() user: RequestUser, ) { @@ -1090,6 +1093,7 @@ export class PlatformApiController { table_name: body.table_name, primary_keys: body.primary_keys, name: body.table_name, + iceberg_table_name: body.iceberg_table_name, }; this.logger.info('addTable: appending CDC table to DynamoDB', { inputId, tableName: body.table_name }); From d119f0d95562de107c9764a01b94b0ffc9e574cb Mon Sep 17 00:00:00 2001 From: Rafael Date: Fri, 21 Aug 2026 11:46:45 -0300 Subject: [PATCH 16/22] test(inputs): add service-level test for CDC create field-mapping InputsService.createCdc builds InputCreateCdcRequest by enumerating fields (not spreading), so a revert of the destination/iceberg_table_name mapping lines would not be caught by the existing controller spec, which only mocks InputsService. Add a unit test at the service boundary that asserts destination and per-table iceberg_table_name reach the gRPC request, plus a back-compat case with no destination. Also mark CdcTableReq.iceberg_table_name as advisory/reserved: the platform derives the Iceberg table name itself today and does not yet consume this field. Co-Authored-By: WOZCODE --- src/modules/inputs/dtos/input.model.ts | 1 + src/modules/inputs/inputs.service.spec.ts | 84 +++++++++++++++++++++++ 2 files changed, 85 insertions(+) create mode 100644 src/modules/inputs/inputs.service.spec.ts diff --git a/src/modules/inputs/dtos/input.model.ts b/src/modules/inputs/dtos/input.model.ts index 259d756..2643521 100644 --- a/src/modules/inputs/dtos/input.model.ts +++ b/src/modules/inputs/dtos/input.model.ts @@ -73,6 +73,7 @@ export class CdcTableReq { table_schema?: string; @ApiPropertyOptional({ type: [String] }) primary_keys?: string[]; + // Advisory/reserved: the platform currently derives the Iceberg table name itself (create_iceberg_table_name); this value is persisted but not yet consumed on the create/add path. Do not treat as the authoritative table name. @ApiPropertyOptional() iceberg_table_name?: string; } diff --git a/src/modules/inputs/inputs.service.spec.ts b/src/modules/inputs/inputs.service.spec.ts new file mode 100644 index 0000000..fed634d --- /dev/null +++ b/src/modules/inputs/inputs.service.spec.ts @@ -0,0 +1,84 @@ +import { of } from 'rxjs'; +import { InputsService } from './inputs.service'; +import DadosferaLogger from '@dadosfera/dadosfera-logs/dist'; +import { CreateCdcInputReq } from './dtos/input.model'; +import { Info } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/entities'; + +const info = { customer_id: 'cid', user_id: 'u' } as unknown as Info; + +describe('InputsService.createCdc', () => { + let service: InputsService; + let inputCreateCdcMock: jest.Mock; + + beforeEach(async () => { + inputCreateCdcMock = jest + .fn() + .mockImplementation((req) => of({ input: req.input })); + + const grpcClient: any = { + getService: jest.fn().mockReturnValue({ + InputCreateCdc: inputCreateCdcMock, + }), + }; + + service = new InputsService(new DadosferaLogger(), grpcClient); + await service.onModuleInit(); + }); + + it('forwards destination and per-table iceberg_table_name to the gRPC request', async () => { + const body: CreateCdcInputReq = { + name: 'CDC Iceberg Test', + plugin: 'mysql_cdc', + read_only: true, + destination: { iceberg: { namespace: 'cdc_raw' } }, + tables: [ + { + name: 'orders', + table_schema: 'mydb', + primary_keys: ['id'], + iceberg_table_name: 'cdc_raw.mydb__orders', + }, + ], + }; + + await service.createCdc({ body, info }); + + expect(inputCreateCdcMock).toHaveBeenCalledTimes(1); + const sentRequest = inputCreateCdcMock.mock.calls[0][0]; + + expect(sentRequest.input).toEqual( + expect.objectContaining({ + destination: { iceberg: { namespace: 'cdc_raw' } }, + }), + ); + expect(sentRequest.input.tables[0]).toEqual( + expect.objectContaining({ + iceberg_table_name: 'cdc_raw.mydb__orders', + }), + ); + }); + + it('back-compat: a body with no destination sends destination undefined, not an error', async () => { + const body: CreateCdcInputReq = { + name: 'CDC Legacy Test', + plugin: 'mysql_cdc', + read_only: true, + tables: [ + { + name: 'pedidos', + table_schema: 'cadastros', + primary_keys: ['id'], + }, + ], + }; + + const result = await service.createCdc({ body, info }); + + expect(inputCreateCdcMock).toHaveBeenCalledTimes(1); + const sentRequest = inputCreateCdcMock.mock.calls[0][0]; + + expect(sentRequest.input.destination).toBeUndefined(); + expect(sentRequest.input.tables[0].iceberg_table_name).toBeUndefined(); + expect(result.input).toBeDefined(); + }); +}); From 4457c0ae628172e098f9613b348c874dedd1eed2 Mon Sep 17 00:00:00 2001 From: Rafael Date: Fri, 21 Aug 2026 18:48:11 -0300 Subject: [PATCH 17/22] feat(cdc): accept + forward source columns on /inputs/cdc Adds CdcColumnReq {name, type, is_primary_key} and columns? on CdcTableReq, forwarded through createCdc and the addTable (Edit Objects add-table) path so the column schema reaches in-factory for Iceberg deduped-table pre-create. Bumps protospack-v2 to 3.41.0-cdc-iceberg.1, which adds the matching CdcColumn field (now required on CdcTable) and regenerates docsfera.json. Co-Authored-By: WOZCODE --- docsfera.json | 88 +++++++++++++++++++ package-lock.json | 8 +- package.json | 2 +- src/modules/inputs/dtos/input.model.ts | 11 +++ src/modules/inputs/inputs.service.spec.ts | 29 ++++++ src/modules/inputs/inputs.service.ts | 1 + .../platform-api.controller.spec.ts | 32 +++++++ .../platform-api/platform-api.controller.ts | 3 + 8 files changed, 169 insertions(+), 5 deletions(-) diff --git a/docsfera.json b/docsfera.json index da23ed4..2653fe1 100644 --- a/docsfera.json +++ b/docsfera.json @@ -5206,6 +5206,44 @@ } }, "/platform/pipelines/{pipelineId}/inputs/{inputId}/tables": { + "delete": { + "operationId": "PlatformApiController_deleteTables", + "summary": "Batch-remove tables from an input; reconfigures the CDC connectors once", + "parameters": [ + { + "name": "pipelineId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + }, + { + "name": "inputId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "" + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + }, + { + "access-token": [] + } + ] + }, "post": { "operationId": "PlatformApiController_addTable", "summary": "Add a CDC table to an input and dispatch its platform jobs", @@ -10887,6 +10925,25 @@ "updated_at" ] }, + "CdcColumnReq": { + "type": "object", + "properties": { + "name": { + "type": "string" + }, + "type": { + "type": "string" + }, + "is_primary_key": { + "type": "boolean" + } + }, + "required": [ + "name", + "type", + "is_primary_key" + ] + }, "CdcTableReq": { "type": "object", "properties": { @@ -10901,12 +10958,40 @@ "items": { "type": "string" } + }, + "iceberg_table_name": { + "type": "string" + }, + "columns": { + "type": "array", + "items": { + "$ref": "#/components/schemas/CdcColumnReq" + } } }, "required": [ "name" ] }, + "IcebergDestinationReq": { + "type": "object", + "properties": { + "namespace": { + "type": "string" + } + }, + "required": [ + "namespace" + ] + }, + "CdcDestinationReq": { + "type": "object", + "properties": { + "iceberg": { + "$ref": "#/components/schemas/IcebergDestinationReq" + } + } + }, "CreateCdcInputReq": { "type": "object", "properties": { @@ -10924,6 +11009,9 @@ }, "read_only": { "type": "boolean" + }, + "destination": { + "$ref": "#/components/schemas/CdcDestinationReq" } }, "required": [ diff --git a/package-lock.json b/package-lock.json index 1ae57c4..3b8fcd3 100644 --- a/package-lock.json +++ b/package-lock.json @@ -16,7 +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-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.0.tgz", + "@dadosfera/protospack-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.1.tgz", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", @@ -1735,9 +1735,9 @@ } }, "node_modules/@dadosfera/protospack-v2": { - "version": "3.41.0-cdc-iceberg.0", - "resolved": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.0.tgz", - "integrity": "sha512-qPs0LukP5BJnybl3JQqqbYQ9oi5SJ86+4A3zd9zisrqqzFHv3o9c0CSwcimwyT3CWm/147vq5LMQy0ulal0Ung==", + "version": "3.41.0-cdc-iceberg.1", + "resolved": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.1.tgz", + "integrity": "sha512-9skEKTRbsgcnACfk6YDg2i9IePdS2FTgH5ngacaxljVv+4+oHsxcWRkIsTbvf/dC6RDAdjd38cU+oa9VGOTCEg==", "license": "ISC", "dependencies": { "@grpc/grpc-js": "^1.9.3", diff --git a/package.json b/package.json index 17c03bf..faf536f 100644 --- a/package.json +++ b/package.json @@ -34,7 +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-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.0.tgz", + "@dadosfera/protospack-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.1.tgz", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", diff --git a/src/modules/inputs/dtos/input.model.ts b/src/modules/inputs/dtos/input.model.ts index 2643521..5f7e330 100644 --- a/src/modules/inputs/dtos/input.model.ts +++ b/src/modules/inputs/dtos/input.model.ts @@ -66,6 +66,15 @@ export class CreateInputReq extends OmitType(Input, [ 'updated_at', ]) {} +export class CdcColumnReq { + @ApiProperty() + name: string; + @ApiProperty() + type: string; + @ApiProperty() + is_primary_key: boolean; +} + export class CdcTableReq { @ApiProperty() name: string; @@ -76,6 +85,8 @@ export class CdcTableReq { // Advisory/reserved: the platform currently derives the Iceberg table name itself (create_iceberg_table_name); this value is persisted but not yet consumed on the create/add path. Do not treat as the authoritative table name. @ApiPropertyOptional() iceberg_table_name?: string; + @ApiPropertyOptional({ type: [CdcColumnReq] }) + columns?: CdcColumnReq[]; } export class IcebergDestinationReq { diff --git a/src/modules/inputs/inputs.service.spec.ts b/src/modules/inputs/inputs.service.spec.ts index fed634d..46d9e38 100644 --- a/src/modules/inputs/inputs.service.spec.ts +++ b/src/modules/inputs/inputs.service.spec.ts @@ -58,6 +58,35 @@ describe('InputsService.createCdc', () => { ); }); + it('forwards per-table columns to the gRPC request', async () => { + const body: CreateCdcInputReq = { + name: 'CDC Columns Test', + plugin: 'mysql_cdc', + read_only: true, + tables: [ + { + name: 'orders', + table_schema: 'mydb', + primary_keys: ['id'], + columns: [ + { name: 'id', type: 'int', is_primary_key: true }, + { name: 'descr', type: 'varchar(255)', is_primary_key: false }, + ], + }, + ], + }; + + await service.createCdc({ body, info }); + + expect(inputCreateCdcMock).toHaveBeenCalledTimes(1); + const sentRequest = inputCreateCdcMock.mock.calls[0][0]; + + expect(sentRequest.input.tables[0].columns).toEqual([ + { name: 'id', type: 'int', is_primary_key: true }, + { name: 'descr', type: 'varchar(255)', is_primary_key: false }, + ]); + }); + it('back-compat: a body with no destination sends destination undefined, not an error', async () => { const body: CreateCdcInputReq = { name: 'CDC Legacy Test', diff --git a/src/modules/inputs/inputs.service.ts b/src/modules/inputs/inputs.service.ts index c221520..38b3a4b 100644 --- a/src/modules/inputs/inputs.service.ts +++ b/src/modules/inputs/inputs.service.ts @@ -198,6 +198,7 @@ export class InputsService { name: t.name, // canonical identity == table_name (in-factory also backfills) primary_keys: t.primary_keys ?? [], iceberg_table_name: t.iceberg_table_name, + columns: t.columns ?? [], })), destination: body.destination, }, diff --git a/src/modules/platform-api/platform-api.controller.spec.ts b/src/modules/platform-api/platform-api.controller.spec.ts index 57127ca..fd6cf78 100644 --- a/src/modules/platform-api/platform-api.controller.spec.ts +++ b/src/modules/platform-api/platform-api.controller.spec.ts @@ -163,6 +163,7 @@ describe('PlatformApiController - addTable', () => { table_name: 'orders', primary_keys: ['id'], name: 'orders', + columns: [], }, info: { customer_id: 'c1', customer: 'cust', user_id: 'u1' }, }); @@ -201,6 +202,37 @@ describe('PlatformApiController - addTable', () => { primary_keys: ['id'], name: 'orders', iceberg_table_name: 'cdc_raw.public__orders', + columns: [], + }, + info: { customer_id: 'c1', customer: 'cust', user_id: 'u1' }, + }); + }); + + it('addTable: carries columns on the added table through to AddCdcTable', async () => { + inputsService.addCdcTable.mockResolvedValue({ input: {} }); + pipelinesClientService.addCdcJobs.mockResolvedValue({ job_ids: ['p_4'], skipped: [] }); + + const columnsBody = { + ...body, + columns: [ + { name: 'id', type: 'int', is_primary_key: true }, + { name: 'descr', type: 'varchar(255)', is_primary_key: false }, + ], + }; + + await controller.addTable('pid', 'iid', columnsBody, mockUser); + + expect(inputsService.addCdcTable).toHaveBeenCalledWith({ + id: 'iid', + table: { + table_schema: 'public', + table_name: 'orders', + primary_keys: ['id'], + name: 'orders', + columns: [ + { name: 'id', type: 'int', is_primary_key: true }, + { name: 'descr', type: 'varchar(255)', is_primary_key: false }, + ], }, info: { customer_id: 'c1', customer: 'cust', user_id: 'u1' }, }); diff --git a/src/modules/platform-api/platform-api.controller.ts b/src/modules/platform-api/platform-api.controller.ts index 8355ffa..cbb9aef 100644 --- a/src/modules/platform-api/platform-api.controller.ts +++ b/src/modules/platform-api/platform-api.controller.ts @@ -1079,6 +1079,8 @@ export class PlatformApiController { // Iceberg destination only (protospack CdcTable.iceberg_table_name); // absent for snowflake, back-compat. iceberg_table_name?: string; + // Source column schema for iceberg deduped table pre-create (protospack CdcTable.columns). + columns?: { name: string; type: string; is_primary_key: boolean }[]; }, @User() user: RequestUser, ) { @@ -1094,6 +1096,7 @@ export class PlatformApiController { primary_keys: body.primary_keys, name: body.table_name, iceberg_table_name: body.iceberg_table_name, + columns: body.columns ?? [], }; this.logger.info('addTable: appending CDC table to DynamoDB', { inputId, tableName: body.table_name }); From 5de033ec193439fa15600957a8e844c0f4016b39 Mon Sep 17 00:00:00 2001 From: Rafael Date: Mon, 24 Aug 2026 09:08:36 -0300 Subject: [PATCH 18/22] feat(cdc-iceberg): thread qualify namespace/table through maestro MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Threads the new CDC→Iceberg fields from the REST DTOs to the gRPC calls: - CdcTableReq.iceberg_qualify_table_name + IcebergDestinationReq.qualify_namespace in input.model.ts - inputs.service.ts create map forwards iceberg_qualify_table_name - platform-api.controller.ts addTable body + cdcTable thread iceberg_qualify_table_name Bumps protospack to v3.41.0-cdc-iceberg.4; docsfera.json regenerated with the new /platform/iceberg/{namespaces,tables/validate} routes. Co-Authored-By: WOZCODE --- docsfera.json | 72 +++++++++++++++++++ package-lock.json | 8 +-- package.json | 2 +- src/modules/inputs/dtos/input.model.ts | 15 +++- src/modules/inputs/inputs.service.ts | 2 + .../platform-api/platform-api.controller.ts | 21 ++++++ 6 files changed, 114 insertions(+), 6 deletions(-) diff --git a/docsfera.json b/docsfera.json index 2653fe1..350903a 100644 --- a/docsfera.json +++ b/docsfera.json @@ -4571,6 +4571,66 @@ ] } }, + "/platform/iceberg/namespaces": { + "get": { + "operationId": "PlatformApiController_getIcebergNamespaces", + "summary": "List existing Polaris Iceberg namespaces (CDC destination dropdown)", + "parameters": [], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + }, + { + "access-token": [] + } + ] + } + }, + "/platform/iceberg/tables/validate": { + "post": { + "operationId": "PlatformApiController_validateIcebergTables", + "summary": "Validate CDC Iceberg raw table names against Polaris", + "parameters": [], + "responses": { + "201": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + }, + { + "access-token": [] + } + ] + } + }, "/platform/pipelines/execute": { "post": { "operationId": "PlatformApiController_executePipeline", @@ -10962,11 +11022,20 @@ "iceberg_table_name": { "type": "string" }, + "iceberg_qualify_table_name": { + "type": "string" + }, "columns": { "type": "array", "items": { "$ref": "#/components/schemas/CdcColumnReq" } + }, + "column_exclude_list": { + "type": "array", + "items": { + "type": "string" + } } }, "required": [ @@ -10978,6 +11047,9 @@ "properties": { "namespace": { "type": "string" + }, + "qualify_namespace": { + "type": "string" } }, "required": [ diff --git a/package-lock.json b/package-lock.json index 3b8fcd3..c39332e 100644 --- a/package-lock.json +++ b/package-lock.json @@ -16,7 +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-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.1.tgz", + "@dadosfera/protospack-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.4.tgz", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", @@ -1735,9 +1735,9 @@ } }, "node_modules/@dadosfera/protospack-v2": { - "version": "3.41.0-cdc-iceberg.1", - "resolved": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.1.tgz", - "integrity": "sha512-9skEKTRbsgcnACfk6YDg2i9IePdS2FTgH5ngacaxljVv+4+oHsxcWRkIsTbvf/dC6RDAdjd38cU+oa9VGOTCEg==", + "version": "3.41.0-cdc-iceberg.4", + "resolved": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.4.tgz", + "integrity": "sha512-QcRhM2wHgoRgXbHfpu+ki+Iozbz7qunKD544v2DCdD1fbeuI6kEit/S4napquSzmzP2qcMjANobyAsWabC/f2w==", "license": "ISC", "dependencies": { "@grpc/grpc-js": "^1.9.3", diff --git a/package.json b/package.json index faf536f..279b538 100644 --- a/package.json +++ b/package.json @@ -34,7 +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-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.1.tgz", + "@dadosfera/protospack-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.4.tgz", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", diff --git a/src/modules/inputs/dtos/input.model.ts b/src/modules/inputs/dtos/input.model.ts index 5f7e330..f7d3aa4 100644 --- a/src/modules/inputs/dtos/input.model.ts +++ b/src/modules/inputs/dtos/input.model.ts @@ -82,16 +82,29 @@ export class CdcTableReq { table_schema?: string; @ApiPropertyOptional({ type: [String] }) primary_keys?: string[]; - // Advisory/reserved: the platform currently derives the Iceberg table name itself (create_iceberg_table_name); this value is persisted but not yet consumed on the create/add path. Do not treat as the authoritative table name. + // Per-table raw Iceberg table name override (iceberg destination only). + // Honored on the create/add path: the platform lowercases + sanitizes it + // authoritatively; empty/absent => the platform derives tb____. @ApiPropertyOptional() iceberg_table_name?: string; + // Per-table deduped (qualify) Iceberg table name override (iceberg dest only). + // Empty/absent => the deduped table takes the same name as the raw table. + @ApiPropertyOptional() + iceberg_qualify_table_name?: string; @ApiPropertyOptional({ type: [CdcColumnReq] }) columns?: CdcColumnReq[]; + // Columns the user chose to ignore -> Debezium column.exclude.list. + @ApiPropertyOptional({ type: [String] }) + column_exclude_list?: string[]; } export class IcebergDestinationReq { @ApiProperty() namespace: string; + // Pipeline-wide deduped (qualify) namespace. Absent => the platform derives + // the sibling of `namespace` (cdc_raw -> cdc_dedup). + @ApiPropertyOptional() + qualify_namespace?: string; } export class CdcDestinationReq { diff --git a/src/modules/inputs/inputs.service.ts b/src/modules/inputs/inputs.service.ts index 38b3a4b..80eef5b 100644 --- a/src/modules/inputs/inputs.service.ts +++ b/src/modules/inputs/inputs.service.ts @@ -198,7 +198,9 @@ export class InputsService { name: t.name, // canonical identity == table_name (in-factory also backfills) primary_keys: t.primary_keys ?? [], iceberg_table_name: t.iceberg_table_name, + iceberg_qualify_table_name: t.iceberg_qualify_table_name, columns: t.columns ?? [], + column_exclude_list: t.column_exclude_list ?? [], })), destination: body.destination, }, diff --git a/src/modules/platform-api/platform-api.controller.ts b/src/modules/platform-api/platform-api.controller.ts index cbb9aef..6ec8c45 100644 --- a/src/modules/platform-api/platform-api.controller.ts +++ b/src/modules/platform-api/platform-api.controller.ts @@ -543,6 +543,20 @@ export class PlatformApiController { return this.platformApiService.proxy('GET', `/pipeline/${normalizedId}`, user); } + @Get('iceberg/namespaces') + @ApiOperation({ summary: 'List existing Polaris Iceberg namespaces (CDC destination dropdown)' }) + @RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET) + async getIcebergNamespaces(@User() user: RequestUser) { + return this.platformApiService.proxy('GET', '/iceberg/namespaces', user); + } + + @Post('iceberg/tables/validate') + @ApiOperation({ summary: 'Validate CDC Iceberg raw table names against Polaris' }) + @RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET) + async validateIcebergTables(@Body() body: any, @User() user: RequestUser) { + return this.platformApiService.proxy('POST', '/iceberg/tables/validate', user, body); + } + @Patch('pipelines/:pipelineId') @ApiOperation({ summary: 'Update pipeline by ID' }) @RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE) @@ -1079,8 +1093,13 @@ export class PlatformApiController { // Iceberg destination only (protospack CdcTable.iceberg_table_name); // absent for snowflake, back-compat. iceberg_table_name?: string; + // Per-table deduped (qualify) Iceberg table name (protospack + // CdcTable.iceberg_qualify_table_name); absent => same as the raw name. + iceberg_qualify_table_name?: string; // Source column schema for iceberg deduped table pre-create (protospack CdcTable.columns). columns?: { name: string; type: string; is_primary_key: boolean }[]; + // Columns the user chose to ignore -> Debezium column.exclude.list. + column_exclude_list?: string[]; }, @User() user: RequestUser, ) { @@ -1096,7 +1115,9 @@ export class PlatformApiController { primary_keys: body.primary_keys, name: body.table_name, iceberg_table_name: body.iceberg_table_name, + iceberg_qualify_table_name: body.iceberg_qualify_table_name, columns: body.columns ?? [], + column_exclude_list: body.column_exclude_list ?? [], }; this.logger.info('addTable: appending CDC table to DynamoDB', { inputId, tableName: body.table_name }); From 23a633badf56ec5cfbae1e855f62d7c21b0c3d4f Mon Sep 17 00:00:00 2001 From: Rafael Date: Thu, 27 Aug 2026 19:14:26 -0300 Subject: [PATCH 19/22] UPDATE: switch protospack-v2 to published @3.40.0-beta.20 MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Replaces the local file:../protospack-v2/...cdc-iceberg.4.tgz tarball reference with the published CodeArtifact version ^3.40.0-beta.20 (carries the CDC→Iceberg qualify_namespace / iceberg_qualify_table_name / column_exclude_list proto fields). Co-Authored-By: WOZCODE --- package-lock.json | 8 ++++---- package.json | 2 +- 2 files changed, 5 insertions(+), 5 deletions(-) diff --git a/package-lock.json b/package-lock.json index c39332e..30d0276 100644 --- a/package-lock.json +++ b/package-lock.json @@ -16,7 +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-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.4.tgz", + "@dadosfera/protospack-v2": "^3.40.0-beta.20", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", @@ -1735,9 +1735,9 @@ } }, "node_modules/@dadosfera/protospack-v2": { - "version": "3.41.0-cdc-iceberg.4", - "resolved": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.4.tgz", - "integrity": "sha512-QcRhM2wHgoRgXbHfpu+ki+Iozbz7qunKD544v2DCdD1fbeuI6kEit/S4napquSzmzP2qcMjANobyAsWabC/f2w==", + "version": "3.40.0-beta.20", + "resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.40.0-beta.20.tgz", + "integrity": "sha512-A12jgcVMCylfXZyXZYLuZNFJBuEBV1ZYmo3w01qhemKrFAzXRjOnPTBGHngdWNHahIaEYdocieBoIzl0n6OvFw==", "license": "ISC", "dependencies": { "@grpc/grpc-js": "^1.9.3", diff --git a/package.json b/package.json index 279b538..02633dc 100644 --- a/package.json +++ b/package.json @@ -34,7 +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-v2": "file:../protospack-v2/dadosfera-protospack-v2-3.41.0-cdc-iceberg.4.tgz", + "@dadosfera/protospack-v2": "^3.40.0-beta.20", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", From 28a11c961a2dffbf9613183757f38ad4d3512e40 Mon Sep 17 00:00:00 2001 From: Rafael Date: Thu, 27 Aug 2026 20:13:00 -0300 Subject: [PATCH 20/22] fix(cdc): allow CDC plugins in RefreshCatalogReq validator Beta's refresh-catalog RefreshCatalogReq DTO restricted plugin to oracle/mysql/postgresql/sqlserver. Under the cache-first catalog model (adopted for CDC in the beta merge), the CDC schema-fetch flow posts /connection-test/refresh-catalog with plugin=mysql_cdc, which the @IsIn rejected ("plugin must be one of: oracle, mysql, postgresql, sqlserver"). Add mysql_cdc/postgresql_cdc/oracle_cdc to the @IsIn and @ApiProperty enum, matching the platform connection-test SQS plugin set. docsfera.json regenerated. Co-Authored-By: WOZCODE --- docsfera.json | 391 +++++++++++++++++- .../connection-test/dto/connection-test.ts | 22 +- 2 files changed, 410 insertions(+), 3 deletions(-) diff --git a/docsfera.json b/docsfera.json index 350903a..eb4af4e 100644 --- a/docsfera.json +++ b/docsfera.json @@ -5071,6 +5071,65 @@ ] } }, + "/platform/pipelines/{pipelineId}/pipeline_run/{runId}/jobs": { + "get": { + "operationId": "PlatformApiController_getPipelineRunJobs", + "summary": "Get pipeline run jobs", + "description": "Proxies platform-api DB-backed job runs and returns `{ jobs: [...] }`.", + "parameters": [ + { + "name": "pipelineId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + }, + { + "name": "runId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "DB-backed job runs for the selected pipeline run.", + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "jobs": { + "type": "array", + "items": { + "type": "object" + } + } + }, + "required": [ + "jobs" + ] + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + }, + { + "access-token": [] + } + ] + } + }, "/platform/jobs/{jobId}/input": { "put": { "operationId": "PlatformApiController_updateJobInput", @@ -6104,6 +6163,48 @@ ] } }, + "/catalog/custom-properties": { + "get": { + "operationId": "CatalogController_getCustomPropertyDefinitions", + "parameters": [ + { + "name": "dadosfera-lang", + "in": "header", + "required": false, + "schema": { + "enum": [ + "pt-br", + "en-us" + ], + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Catalog" + ], + "security": [ + { + "access-token": [] + }, + { + "access-token": [] + } + ] + } + }, "/catalog/data-asset/{id}": { "get": { "operationId": "CatalogController_getDataAsset", @@ -6364,6 +6465,57 @@ "access-token": [] } ] + }, + "patch": { + "operationId": "CatalogController_updateColumnsMetadata", + "parameters": [ + { + "name": "dadosfera-lang", + "in": "header", + "required": false, + "schema": { + "enum": [ + "pt-br", + "en-us" + ], + "type": "string" + } + }, + { + "name": "id", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "requestBody": { + "required": true, + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/UpdateColumnsMetadataRequest" + } + } + } + }, + "responses": { + "200": { + "description": "" + } + }, + "tags": [ + "Catalog" + ], + "security": [ + { + "access-token": [] + }, + { + "access-token": [] + } + ] } }, "/catalog/data-asset/{id}/preview": { @@ -7842,6 +7994,94 @@ ] } }, + "/connection-test/refresh-catalog": { + "post": { + "operationId": "ConnectionTestController_refreshCatalog", + "parameters": [], + "requestBody": { + "required": true, + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/RefreshCatalogReq" + } + } + } + }, + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/RefreshCatalogRes" + } + } + } + }, + "202": { + "description": "", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/RefreshCatalogRes" + } + } + } + } + }, + "tags": [ + "Connection Test" + ], + "security": [ + { + "access-token": [] + }, + { + "access-token": [] + } + ] + } + }, + "/connection-test/refresh-catalog/status": { + "post": { + "operationId": "ConnectionTestController_refreshCatalogStatus", + "parameters": [], + "requestBody": { + "required": true, + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/RefreshCatalogStatusReq" + } + } + } + }, + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/RefreshCatalogRes" + } + } + } + } + }, + "tags": [ + "Connection Test" + ], + "security": [ + { + "access-token": [] + }, + { + "access-token": [] + } + ] + } + }, "/oauth/hubspot": { "get": { "operationId": "OauthController_oauthHubspot", @@ -11715,6 +11955,35 @@ "columns_metadata" ] }, + "ColumnDescriptionDto": { + "type": "object", + "properties": { + "column_name": { + "type": "string" + }, + "description": { + "type": "string" + } + }, + "required": [ + "column_name", + "description" + ] + }, + "UpdateColumnsMetadataRequest": { + "type": "object", + "properties": { + "columns": { + "type": "array", + "items": { + "$ref": "#/components/schemas/ColumnDescriptionDto" + } + } + }, + "required": [ + "columns" + ] + }, "IData": { "type": "object", "properties": { @@ -11817,6 +12086,37 @@ "docs" ] }, + "CustomPropertyDto": { + "type": "object", + "properties": { + "key": { + "type": "string" + }, + "value": { + "type": "string" + }, + "type": { + "type": "string", + "enum": [ + "text", + "number", + "date", + "boolean" + ] + }, + "color": { + "type": "string" + }, + "emoji": { + "type": "string" + } + }, + "required": [ + "key", + "value", + "type" + ] + }, "IUpdateDataRequest": { "type": "object", "properties": { @@ -11845,6 +12145,12 @@ }, "docs": { "type": "string" + }, + "custom_properties": { + "type": "array", + "items": { + "$ref": "#/components/schemas/CustomPropertyDto" + } } }, "required": [ @@ -12299,11 +12605,15 @@ }, "type": { "type": "string" + }, + "is_primary_key": { + "type": "boolean" } }, "required": [ "name", - "type" + "type", + "is_primary_key" ] }, "TableMetadataDto": { @@ -12405,6 +12715,85 @@ "checks" ] }, + "RefreshCatalogReq": { + "type": "object", + "properties": { + "connection_id": { + "type": "string" + }, + "plugin": { + "type": "string", + "enum": [ + "oracle", + "mysql", + "postgresql", + "sqlserver", + "mysql_cdc", + "postgresql_cdc", + "oracle_cdc" + ] + } + }, + "required": [ + "connection_id", + "plugin" + ] + }, + "RefreshCatalogRes": { + "type": "object", + "properties": { + "operation_result": { + "type": "boolean" + }, + "status": { + "type": "string" + }, + "session_id": { + "type": "string" + }, + "date": { + "type": "string" + } + }, + "required": [ + "operation_result", + "status", + "session_id", + "date" + ] + }, + "RefreshCatalogStatusReq": { + "type": "object", + "properties": { + "connection_id": { + "type": "string" + }, + "plugin": { + "type": "string", + "enum": [ + "oracle", + "mysql", + "postgresql", + "sqlserver", + "mysql_cdc", + "postgresql_cdc", + "oracle_cdc" + ] + }, + "session_id": { + "type": "string" + }, + "date": { + "type": "string" + } + }, + "required": [ + "connection_id", + "plugin", + "session_id", + "date" + ] + }, "INote": { "type": "object", "properties": { diff --git a/src/modules/connection-test/dto/connection-test.ts b/src/modules/connection-test/dto/connection-test.ts index cb10c4c..20447b9 100644 --- a/src/modules/connection-test/dto/connection-test.ts +++ b/src/modules/connection-test/dto/connection-test.ts @@ -176,8 +176,26 @@ export class RefreshCatalogReq { @IsString() connection_id: string; - @ApiProperty({ enum: ['oracle', 'mysql', 'postgresql', 'sqlserver'] }) - @IsIn(['oracle', 'mysql', 'postgresql', 'sqlserver']) + @ApiProperty({ + enum: [ + 'oracle', + 'mysql', + 'postgresql', + 'sqlserver', + 'mysql_cdc', + 'postgresql_cdc', + 'oracle_cdc', + ], + }) + @IsIn([ + 'oracle', + 'mysql', + 'postgresql', + 'sqlserver', + 'mysql_cdc', + 'postgresql_cdc', + 'oracle_cdc', + ]) plugin: string; } From 5c609a8df8fffeb975b31be0593b3f20d4856284 Mon Sep 17 00:00:00 2001 From: Rafael Date: Fri, 28 Aug 2026 09:38:20 -0300 Subject: [PATCH 21/22] test(maestro): fix connection-test merge test failures MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - platform-api.controller.spec: addCdcTable expectations now include the CDC fields the controller threads (iceberg_table_name, iceberg_qualify_table_name, column_exclude_list) which were added by the CDC-iceberg work. - release_note specs: provide DadosferaLogger mock — ReleaseNoteService gained an @Inject(DadosferaLogger) dependency (from beta) without its specs being updated, so they failed DI resolution on merge. Full suite: 52 passed, 5 skipped, 0 failed. Co-Authored-By: WOZCODE --- .../platform-api/platform-api.controller.spec.ts | 8 ++++++++ .../release_note/release_note.controller.spec.ts | 14 +++++++++++++- .../release_note/release_note.service.spec.ts | 14 +++++++++++++- 3 files changed, 34 insertions(+), 2 deletions(-) diff --git a/src/modules/platform-api/platform-api.controller.spec.ts b/src/modules/platform-api/platform-api.controller.spec.ts index fd6cf78..f3d58a3 100644 --- a/src/modules/platform-api/platform-api.controller.spec.ts +++ b/src/modules/platform-api/platform-api.controller.spec.ts @@ -163,7 +163,10 @@ describe('PlatformApiController - addTable', () => { table_name: 'orders', primary_keys: ['id'], name: 'orders', + iceberg_table_name: undefined, + iceberg_qualify_table_name: undefined, columns: [], + column_exclude_list: [], }, info: { customer_id: 'c1', customer: 'cust', user_id: 'u1' }, }); @@ -202,7 +205,9 @@ describe('PlatformApiController - addTable', () => { primary_keys: ['id'], name: 'orders', iceberg_table_name: 'cdc_raw.public__orders', + iceberg_qualify_table_name: undefined, columns: [], + column_exclude_list: [], }, info: { customer_id: 'c1', customer: 'cust', user_id: 'u1' }, }); @@ -229,10 +234,13 @@ describe('PlatformApiController - addTable', () => { table_name: 'orders', primary_keys: ['id'], name: 'orders', + iceberg_table_name: undefined, + iceberg_qualify_table_name: undefined, columns: [ { name: 'id', type: 'int', is_primary_key: true }, { name: 'descr', type: 'varchar(255)', is_primary_key: false }, ], + column_exclude_list: [], }, info: { customer_id: 'c1', customer: 'cust', user_id: 'u1' }, }); diff --git a/src/modules/release_note/release_note.controller.spec.ts b/src/modules/release_note/release_note.controller.spec.ts index 921eca1..1fd7eba 100644 --- a/src/modules/release_note/release_note.controller.spec.ts +++ b/src/modules/release_note/release_note.controller.spec.ts @@ -1,14 +1,26 @@ import { Test, TestingModule } from '@nestjs/testing'; +import { DadosferaLogger } from '@dadosfera/dadosfera-logs'; import { ReleaseNoteController } from './release_note.controller'; import { ReleaseNoteService } from './release_note.service'; +const logger = { + info: (...args) => args, + error: (...args) => args, +}; + describe('ReleaseNoteController', () => { let controller: ReleaseNoteController; beforeEach(async () => { const module: TestingModule = await Test.createTestingModule({ controllers: [ReleaseNoteController], - providers: [ReleaseNoteService], + providers: [ + ReleaseNoteService, + { + provide: DadosferaLogger, + useValue: { logger }, + }, + ], }).compile(); controller = module.get(ReleaseNoteController); diff --git a/src/modules/release_note/release_note.service.spec.ts b/src/modules/release_note/release_note.service.spec.ts index 72add9d..70d2024 100644 --- a/src/modules/release_note/release_note.service.spec.ts +++ b/src/modules/release_note/release_note.service.spec.ts @@ -1,12 +1,24 @@ import { Test, TestingModule } from '@nestjs/testing'; +import { DadosferaLogger } from '@dadosfera/dadosfera-logs'; import { ReleaseNoteService } from './release_note.service'; +const logger = { + info: (...args) => args, + error: (...args) => args, +}; + describe('ReleaseNoteService', () => { let service: ReleaseNoteService; beforeEach(async () => { const module: TestingModule = await Test.createTestingModule({ - providers: [ReleaseNoteService], + providers: [ + ReleaseNoteService, + { + provide: DadosferaLogger, + useValue: { logger }, + }, + ], }).compile(); service = module.get(ReleaseNoteService); From b0fa8d29fc7153da13908bf6db0560d5d5dae92c Mon Sep 17 00:00:00 2001 From: Rafael Date: Fri, 28 Aug 2026 10:32:56 -0300 Subject: [PATCH 22/22] test(maestro): set INFACTORY_URL in the test Docker target Beta's cache-first work added connection-test.service.spec.ts, whose import graph (connection/client.config.ts) reads process.env.INFACTORY_URL at load time. The Dockerfile test target only set DUC_URL, so that suite crashed at import ("Cannot read properties of undefined (reading 'startsWith')") in CI and in any bare `npm test` run. Add ENV INFACTORY_URL=0.0.0.0:50052 alongside the existing DUC_URL, matching the local-connection convention (0.0.0.0 => no SSL). Full suite: 52 passed, 5 skipped. Co-Authored-By: WOZCODE --- Dockerfile | 1 + 1 file changed, 1 insertion(+) diff --git a/Dockerfile b/Dockerfile index 97b68de..e6d1888 100644 --- a/Dockerfile +++ b/Dockerfile @@ -29,6 +29,7 @@ COPY . . # unit test specific build FROM ci_image AS test ENV DUC_URL=0.0.0.0:50051 +ENV INFACTORY_URL=0.0.0.0:50052 ENTRYPOINT ["npm", "run", "test"]