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(); + }); +});