diff --git a/.github/workflows/validate-k8s.yml b/.github/workflows/validate-k8s.yml index 8a67a9c..ae6d383 100644 --- a/.github/workflows/validate-k8s.yml +++ b/.github/workflows/validate-k8s.yml @@ -71,6 +71,11 @@ jobs: sudo mv helmfile /usr/local/bin/ helmfile --version + - name: Install Helm Diff plugin + run: | + helm plugin install https://github.com/databus23/helm-diff --version v3.9.3 + helm diff version + - name: Debug Helm env run: | helm env @@ -102,4 +107,5 @@ jobs: - name: Run Helmfile Diff env: ENV: ${{ needs.extract_environment.outputs.environment }} + HELM_PLUGINS: /home/runner/.local/share/helm/plugins run: helmfile -f deploy/helmfiles/${ENV}.yaml diff diff --git a/deploy/helm-chart/templates/deployment.yaml b/deploy/helm-chart/templates/deployment.yaml index 518fd7b..2cdfb2c 100644 --- a/deploy/helm-chart/templates/deployment.yaml +++ b/deploy/helm-chart/templates/deployment.yaml @@ -111,6 +111,8 @@ spec: value: "{{ .Values.maestro.redis_tls }}" - name: PLATFORM_API_URL value: {{ .Values.maestro.platform_api_url }} + - name: CONNECTIONS_API_URL + value: {{ .Values.maestro.connections_api_url | default "" | quote }} - name: STORAGE_EXPLORER_API_URL value: {{ .Values.maestro.storage_explorer_api_url | quote }} - name: FIREBASE_BASE_URL diff --git a/deploy/helm-chart/values-stg.yaml b/deploy/helm-chart/values-stg.yaml index 70d2170..9515b1c 100644 --- a/deploy/helm-chart/values-stg.yaml +++ b/deploy/helm-chart/values-stg.yaml @@ -9,6 +9,7 @@ maestro: cookie_secret: "ff7bc13823edb2ae50d248e5780bddc9d4b31c36" redis_database: "1" platform_api_url: https://xs2hkhq07k.execute-api.us-east-1.amazonaws.com + connections_api_url: https://iy40eans64.execute-api.us-east-1.amazonaws.com storage_explorer_api_url: "http://storage-explorer-{customer}.data-apps.svc.cluster.local:8000/api" firebase_base_url: https://feature-flag-25bf6-default-rtdb.firebaseio.com/stg diff --git a/src/modules/auth/auth.controller.ts b/src/modules/auth/auth.controller.ts index 40eaad9..5eb2b0f 100644 --- a/src/modules/auth/auth.controller.ts +++ b/src/modules/auth/auth.controller.ts @@ -37,6 +37,7 @@ import { RequireAllPermissions, } from 'src/decorators/authentication.decorator'; import { AuthClientService } from './auth.service'; +import { UserDTO } from './dtos/login'; import { DadosferaLogger } from '@dadosfera/dadosfera-logs'; import { GrpcToHttpExceptionFilter } from '../../error/grpc-to-http-exception.filter'; import { RequestUser, User } from 'src/decorators/user.decorator'; @@ -486,7 +487,7 @@ export class AuthController { this.logger.info('Authenticating via X-Api-key header'); const { api_key } = await this.apiKeyService.get(apiKey); - const userDto = { + const userDto: UserDTO = { id: api_key.user_id, name: api_key.username, email: api_key.username, @@ -494,7 +495,8 @@ export class AuthController { id: api_key.customer_id, name: api_key.customer_name, tier: api_key.customer_tier, - } + }, + permissions: [], }; return res.status(200).json(userDto); diff --git a/src/modules/auth/auth.service.ts b/src/modules/auth/auth.service.ts index 8539d0e..c82aeef 100644 --- a/src/modules/auth/auth.service.ts +++ b/src/modules/auth/auth.service.ts @@ -447,6 +447,9 @@ export class AuthClientService implements OnModuleInit { name: payload.customer_name, tier: payload.customer_tier, }, + // Raw permission seqids from the JWT. Consumers own the seqid->meaning + // mapping (e.g. Orchest's auth-server); Maestro reports them as-is. + permissions: payload.permissions ?? [], }; return userDto; diff --git a/src/modules/auth/dtos/login.ts b/src/modules/auth/dtos/login.ts index 2e22363..1ea22f0 100644 --- a/src/modules/auth/dtos/login.ts +++ b/src/modules/auth/dtos/login.ts @@ -152,5 +152,6 @@ export type UserDTO = { id: string, name: string, tier: string, - } + }, + permissions: number[], } diff --git a/src/modules/catalog/catalog.controller.ts b/src/modules/catalog/catalog.controller.ts index ef8ab78..2b2bc74 100644 --- a/src/modules/catalog/catalog.controller.ts +++ b/src/modules/catalog/catalog.controller.ts @@ -9,6 +9,7 @@ import { Inject, NotFoundException, Param, + Patch, Post, Put, Query, @@ -51,6 +52,7 @@ import { IUpdateDataRequest, TriggerCatalogReq, TriggerCatalogRes, + UpdateColumnsMetadataRequest, } from './dtos'; import { GrpcToHttpExceptionFilter } from 'src/error/grpc-to-http-exception.filter'; import { Language } from 'src/decorators/language.decorator'; @@ -283,6 +285,23 @@ export class CatalogController { } + @Get('custom-properties') + @RequireSomePermission( + PERMISSIONS_GROUPS.CATALOG.permissions.GET, + PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER, + ) + async getCustomPropertyDefinitions(@User() user: RequestUser) { + const { customer_id, customer_name, user_id, username } = user; + const metadata = PackTheMetadata({ + customer_id, + customer_name, + user_id, + username, + }); + + return this.catalogService.getCustomPropertyDefinitions(metadata); + } + @Get('data-asset/:id') @RequireSomePermission( PERMISSIONS_GROUPS.CATALOG.permissions.GET, @@ -431,6 +450,42 @@ export class CatalogController { return { columns_metadata }; } + @Patch('data-asset/:id/columns-metadata') + @RequireSomePermission( + PERMISSIONS_GROUPS.CATALOG.permissions.UPDATE, + PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER, + ) + async updateColumnsMetadata( + @User() user: RequestUser, + @Language() language: LanguageEnum, + @Param('id') id: string, + @Body(new ValidationPipe()) body: UpdateColumnsMetadataRequest, + ): Promise<{ success: boolean }> { + const { customer_name, customer_id, user_id, username } = user; + + this.logger.info(`/catalog - update columns metadata`, { + user_id, + customer_name, + columns_count: body.columns.length, + }); + + const metadata = PackTheMetadata({ + customer_name, + customer_id, + user_id, + username, + language, + }); + + await this.catalogService.updateColumnsDescriptions( + id, + body.columns, + metadata, + ); + + return { success: true }; + } + @Get('data-asset/:id/preview') @RequireSomePermission( PERMISSIONS_GROUPS.CATALOG.permissions.GET, @@ -1056,4 +1111,4 @@ export class CatalogController { this.logger.error(error.message); } } -} \ No newline at end of file +} diff --git a/src/modules/catalog/catalog.service.ts b/src/modules/catalog/catalog.service.ts index f0fced8..d50e819 100644 --- a/src/modules/catalog/catalog.service.ts +++ b/src/modules/catalog/catalog.service.ts @@ -120,6 +120,10 @@ class CatalogService implements OnModuleInit { } } + async getCustomPropertyDefinitions(metadata: Metadata) { + return lastValueFrom(this.catalogReadService.GetCustomPropertyDefinitions({}, metadata)); + } + async createDataAsset(data: Messages.CreateDataAssetRequest, metadata) { this.logger.info('CatalogService - Manage Data assets permissions'); if (!data.embed) data.embed = undefined; @@ -463,6 +467,16 @@ class CatalogService implements OnModuleInit { return result; } + async updateColumnsDescriptions( + id: string, + columns: { column_name: string; description: string }[], + metadata: Metadata, + ) { + await lastValueFrom( + this.catalogWriteService.UpdateColumnDescriptions({ id, columns }, metadata), + ); + } + async createDataDocs(body: CreateDataDocsDTO, metadata: Metadata) { if (body.asset_type === 'table' || body.asset_type === 'view') { return this.createDataDocsViaNimbus(body); @@ -828,6 +842,8 @@ class CatalogService implements OnModuleInit { data_asset_id: table_metadata_id.toString(), customer_name: customer_name, data_asset_type: 'dataset', + column_metadata: [], + data_preview: '', }, ], }, diff --git a/src/modules/catalog/dtos/index.ts b/src/modules/catalog/dtos/index.ts index 770b01c..c462deb 100644 --- a/src/modules/catalog/dtos/index.ts +++ b/src/modules/catalog/dtos/index.ts @@ -1,5 +1,13 @@ import { ApiProperty, ApiPropertyOptional, PickType } from '@nestjs/swagger'; -import { IsEnum } from 'class-validator'; +import { + ArrayNotEmpty, + IsArray, + IsEnum, + IsNotEmpty, + IsString, + ValidateNested, +} from 'class-validator'; +import { Type } from 'class-transformer'; import { CreateDataAssetRequest } from '@dadosfera/protospack-v2/dist/lib/Catalog/interfaces/messages'; export enum DataAssetShareType { @@ -198,6 +206,27 @@ export class IData { day_opening: number; } + +export enum CustomPropertyType { + TEXT = 'text', + NUMBER = 'number', + DATE = 'date', + BOOLEAN = 'boolean', +} + +export class CustomPropertyDto { + @ApiProperty() + key: string; + @ApiProperty() + value: string; + @ApiProperty({ enum: CustomPropertyType }) + type: CustomPropertyType; + @ApiPropertyOptional() + color?: string; + @ApiPropertyOptional() + emoji?: string; +} + export class IUpdateDataRequest { @ApiProperty() name: string; @@ -211,6 +240,8 @@ export class IUpdateDataRequest { share_type?: DataAssetShareType; @ApiPropertyOptional() docs?: string; + @ApiPropertyOptional({ type: [CustomPropertyDto] }) + custom_properties?: CustomPropertyDto[]; } export class IUpdateCertificationStatusRequest { @@ -219,6 +250,26 @@ export class IUpdateCertificationStatusRequest { certification_status: CertificationStatus; } +export class ColumnDescriptionDto { + @ApiProperty() + @IsString() + @IsNotEmpty() + column_name: string; + + @ApiProperty() + @IsString() + description: string; +} + +export class UpdateColumnsMetadataRequest { + @ApiProperty({ type: [ColumnDescriptionDto] }) + @IsArray() + @ArrayNotEmpty() + @ValidateNested({ each: true }) + @Type(() => ColumnDescriptionDto) + columns: ColumnDescriptionDto[]; +} + export class ICreateDataAsset implements CreateDataAssetRequest { @ApiProperty() display_name: string; @@ -372,4 +423,4 @@ export type CreateDataDocsDTO = { docs: string; asset_type: string; -} \ No newline at end of file +} diff --git a/src/modules/connection-test/connection-test.controller.ts b/src/modules/connection-test/connection-test.controller.ts index 939aded..bf6d383 100644 --- a/src/modules/connection-test/connection-test.controller.ts +++ b/src/modules/connection-test/connection-test.controller.ts @@ -24,6 +24,9 @@ import { GetTableMetadataReq, ValidateCdcPrerequisitesReq, ValidateCdcPrerequisitesRes, + RefreshCatalogReq, + RefreshCatalogRes, + RefreshCatalogStatusReq, } from './dto/connection-test'; import { DadosferaLogger } from '@dadosfera/dadosfera-logs'; import { Authenticated, RequireModule } from 'src/decorators/authentication.decorator'; @@ -90,7 +93,7 @@ export class ConnectionTestController { }); return this.connectionTestService.connectionTestListSchemas( body, - user.customer_name, + user, ); } @@ -107,7 +110,7 @@ export class ConnectionTestController { }); return this.connectionTestService.connectionTestListTables( body, - user.customer_name, + user, ); } @@ -124,7 +127,7 @@ export class ConnectionTestController { }); return this.connectionTestService.getTableMetadata( body, - user.customer_name, + user, ); } @@ -144,4 +147,35 @@ export class ConnectionTestController { user.customer_name, ); } + + @Post('refresh-catalog') + @ApiOkResponse({ type: RefreshCatalogRes }) + @HttpCode(HttpStatus.ACCEPTED) + async refreshCatalog( + @User() user: RequestUser, + @Body(new ValidationPipe()) body: RefreshCatalogReq, + ) { + this.logger.info('/connection-test/refresh-catalog', { + user: user.user_id, + customer: user.customer_name, + connection: body.connection_id, + }); + return this.connectionTestService.refreshCatalog(body, user); + } + + @Post('refresh-catalog/status') + @ApiOkResponse({ type: RefreshCatalogRes }) + @HttpCode(HttpStatus.OK) + async refreshCatalogStatus( + @User() user: RequestUser, + @Body(new ValidationPipe()) body: RefreshCatalogStatusReq, + ) { + this.logger.info('/connection-test/refresh-catalog/status', { + user: user.user_id, + customer: user.customer_name, + connection: body.connection_id, + session: body.session_id, + }); + return this.connectionTestService.refreshCatalogStatus(body, user); + } } diff --git a/src/modules/connection-test/connection-test.module.ts b/src/modules/connection-test/connection-test.module.ts index 1eb4099..c184391 100644 --- a/src/modules/connection-test/connection-test.module.ts +++ b/src/modules/connection-test/connection-test.module.ts @@ -5,10 +5,17 @@ import { DadosferaLogger } from '@dadosfera/dadosfera-logs'; import { ClientsModule } from '@nestjs/microservices'; import { ConnectionTestClientConfiguration } from './connection-test-client.config'; import { ConnectionModule } from '../connection/connection.module'; +import { ConnectionsApiModule } from '../connections-api/connections-api.module'; +import { PlatformApiModule } from '../platform-api/platform-api.module'; const client = new ConnectionTestClientConfiguration(); @Module({ controllers: [ConnectionTestController], providers: [ConnectionTestService, DadosferaLogger], - imports: [ClientsModule.register([client.providerOptions]), ConnectionModule], + imports: [ + ClientsModule.register([client.providerOptions]), + ConnectionModule, + ConnectionsApiModule, + PlatformApiModule, + ], }) export class ConnectionTestModule {} diff --git a/src/modules/connection-test/connection-test.service.spec.ts b/src/modules/connection-test/connection-test.service.spec.ts new file mode 100644 index 0000000..8398d48 --- /dev/null +++ b/src/modules/connection-test/connection-test.service.spec.ts @@ -0,0 +1,223 @@ +import { ConnectionTestService } from './connection-test.service'; +import { RequestUser } from 'src/decorators/user.decorator'; + +describe('ConnectionTestService catalog cache', () => { + const user: RequestUser = { + user_id: 'user-id', + username: 'user@example.com', + permissions: [], + customer_id: 'customer-id', + customer_name: 'customer-name', + customer_tier: 'standard', + access_token: 'token', + customer_modules: [], + roles: [], + }; + const grpcClient = { getService: jest.fn().mockReturnValue({}) }; + const connectionsService = {}; + const connectionsApiService = { proxy: jest.fn() }; + const platformApiService = { proxy: jest.fn() }; + let service: ConnectionTestService; + + beforeEach(() => { + jest.clearAllMocks(); + service = new ConnectionTestService( + grpcClient as any, + connectionsService as any, + connectionsApiService as any, + platformApiService as any, + ); + }); + + it('keeps the existing schemas response contract', async () => { + connectionsApiService.proxy.mockResolvedValue({ + schemas: [{ schema_name: 'analytics' }, { schema_name: 'public' }], + }); + + await expect( + service.connectionTestListSchemas( + { connection_id: 'config-id', plugin: 'postgresql' }, + user, + ), + ).resolves.toEqual({ + operation_result: true, + schema_list: ['analytics', 'public'], + }); + }); + + it('lists tables and enriches each with its cached primary keys', async () => { + connectionsApiService.proxy + // list-tables call (names only from the catalog cache) + .mockResolvedValueOnce({ + tables: [{ table_name: 'customers' }, { table_name: 'orders' }], + }) + // per-table columns calls: customers has a PK, orders has none + .mockResolvedValueOnce({ + columns: [ + { column_name: 'id', data_type: 'bigint', is_primary_key: true }, + { column_name: 'name', data_type: 'text', is_primary_key: false }, + ], + }) + .mockResolvedValueOnce({ + columns: [ + { column_name: 'total', data_type: 'numeric', is_primary_key: false }, + ], + }); + + await expect( + service.connectionTestListTables( + { + connection_id: 'config-id', + plugin: 'postgresql', + schema: 'public', + }, + user, + ), + ).resolves.toEqual({ + operation_result: true, + table_list: ['customers', 'orders'], + tables: [ + { table_name: 'customers', primary_keys: ['id'] }, + { table_name: 'orders', primary_keys: [] }, + ], + }); + }); + + it('maps cached columns to the existing table metadata contract', async () => { + connectionsApiService.proxy.mockResolvedValue({ + columns: [ + { + column_name: 'id', + data_type: 'bigint', + is_primary_key: true, + }, + ], + }); + + await expect( + service.getTableMetadata( + { + connection_id: 'config-id', + plugin: 'postgresql', + schema: 'public', + table_list: ['customers'], + }, + user, + ), + ).resolves.toEqual({ + operation_result: true, + tables_metadata: [ + { + table_name: 'customers', + columns: [ + { + name: 'id', + type: 'bigint', + is_primary_key: true, + }, + ], + references: [], + }, + ], + }); + expect(connectionsApiService.proxy).toHaveBeenCalledWith( + 'GET', + '/connection_catalog/config-id/schemas/public/tables/customers/columns', + user, + ); + }); + + it('submits a catalog refresh without holding the request open', async () => { + platformApiService.proxy.mockResolvedValue({ + session_id: 'session-id', + date: '20260731', + }); + + await expect( + service.refreshCatalog( + { connection_id: 'config-id', plugin: 'postgresql' }, + user, + ), + ).resolves.toEqual({ + operation_result: true, + status: 'PENDING', + session_id: 'session-id', + date: '20260731', + }); + + expect(platformApiService.proxy).toHaveBeenCalledWith( + 'POST', + '/connection_test', + user, + { + customer_id: user.customer_name, + plugin: 'postgresql', + task: { + task_type: 'refresh_catalog', + connection: { + provider: 'connection_manager', + config_id: 'config-id', + }, + }, + }, + ); + }); + + it('keeps polling without changing the catalog pointer while pending', async () => { + platformApiService.proxy.mockResolvedValue({ status: 'PENDING' }); + + await expect( + service.refreshCatalogStatus( + { + connection_id: 'config-id', + plugin: 'postgresql', + session_id: 'session-id', + date: '20260731', + }, + user, + ), + ).resolves.toEqual({ + operation_result: false, + status: 'PENDING', + session_id: 'session-id', + date: '20260731', + }); + + expect(connectionsApiService.proxy).not.toHaveBeenCalled(); + }); + + it('publishes the catalog pointer after the refresh finishes', async () => { + platformApiService.proxy.mockResolvedValue({ status: 'DONE' }); + connectionsApiService.proxy.mockResolvedValue({ + last_catalog_refresh_status: 'SUCCESS', + }); + + await expect( + service.refreshCatalogStatus( + { + connection_id: 'config/id', + plugin: 'postgresql', + session_id: 'session-id', + date: '20260731', + }, + user, + ), + ).resolves.toEqual({ + operation_result: true, + status: 'DONE', + session_id: 'session-id', + date: '20260731', + }); + + expect(connectionsApiService.proxy).toHaveBeenCalledWith( + 'PUT', + '/connection_config/config%2Fid/catalog_metadata', + user, + { + last_catalog_refresh_status: 'SUCCESS', + last_catalog_connection_test_date: '20260731', + last_catalog_connection_test_session_id: 'session-id', + }, + ); + }); +}); diff --git a/src/modules/connection-test/connection-test.service.ts b/src/modules/connection-test/connection-test.service.ts index 096a34e..f14efc1 100644 --- a/src/modules/connection-test/connection-test.service.ts +++ b/src/modules/connection-test/connection-test.service.ts @@ -1,4 +1,4 @@ -import { Inject, Injectable } from '@nestjs/common'; +import { HttpException, HttpStatus, Inject, Injectable } from '@nestjs/common'; import { ClientGrpc } from '@nestjs/microservices'; import { ConnectionTest } from '@dadosfera/protospack-v2'; import { lastValueFrom } from 'rxjs'; @@ -15,6 +15,9 @@ import { GetTableMetadataRes, ValidateCdcPrerequisitesReq, ValidateCdcPrerequisitesRes, + RefreshCatalogReq, + RefreshCatalogRes, + RefreshCatalogStatusReq, } from './dto/connection-test'; import { ConnectionClientService } from '../connection/client.service'; import { @@ -23,6 +26,8 @@ import { } from '../connection/dtos/connection'; import { RequestUser } from 'src/decorators/user.decorator'; import { PackTheMetadata } from 'src/utils/PackTheMetadata'; +import { ConnectionsApiService } from '../connections-api/connections-api.service'; +import { PlatformApiService } from '../platform-api/platform-api.service'; @Injectable() export class ConnectionTestService { @@ -30,6 +35,8 @@ export class ConnectionTestService { constructor( @Inject('ConnectionTestGrpcClient') private readonly grpcClient: ClientGrpc, private connectionsService: ConnectionClientService, + private connectionsApiService: ConnectionsApiService, + private platformApiService: PlatformApiService, ) { this.connectionTestReadClient = grpcClient.getService( @@ -149,46 +156,162 @@ export class ConnectionTestService { } async connectionTestListSchemas( body: ConnectionTestListSchemasReq, - customer_name: string, + user: RequestUser, ): Promise { - const { connection_id, plugin } = body; - return lastValueFrom( - this.connectionTestReadClient.ListSchemas({ - connection_id, - customer_name, - plugin, - }), + const result = await this.connectionsApiService.proxy( + 'GET', + `/connection_catalog/${encodeURIComponent(body.connection_id)}/schemas`, + user, ); + return { + operation_result: true, + schema_list: result.schemas.map((schema) => schema.schema_name), + }; } + async connectionTestListTables( body: ConnectionTestListTablesReq, - customer_name: string, + user: RequestUser, ): Promise { - const { connection_id, plugin, schema } = body; - return lastValueFrom( - this.connectionTestReadClient.ListTables({ - connection_id, - customer_name, - plugin, - schema, + const result = await this.connectionsApiService.proxy( + 'GET', + `/connection_catalog/${encodeURIComponent(body.connection_id)}` + + `/schemas/${encodeURIComponent(body.schema)}/tables`, + user, + ); + const table_names: string[] = result.tables.map((table) => table.table_name); + // CDC create needs the primary keys per table (used to build the deduped + // Iceberg identifier-fields). The catalog-cache list-tables endpoint returns + // only names, so fetch each table's columns from the cache and keep the ones + // flagged is_primary_key. Reads hit the stored catalog snapshot (populated by + // refresh-catalog), never the live connection. + const tables = await Promise.all( + table_names.map(async (table_name) => { + const columns = await this.connectionsApiService.proxy( + 'GET', + `/connection_catalog/${encodeURIComponent(body.connection_id)}` + + `/schemas/${encodeURIComponent(body.schema)}` + + `/tables/${encodeURIComponent(table_name)}/columns`, + user, + ); + return { + table_name, + primary_keys: columns.columns + .filter((column) => column.is_primary_key) + .map((column) => column.column_name), + }; }), ); + return { + operation_result: true, + table_list: table_names, + tables, + }; } async getTableMetadata( body: GetTableMetadataReq, - customer_name: string, + user: RequestUser, ): Promise { - const { schema, plugin, table_list, connection_id } = body; - return lastValueFrom( - this.connectionTestReadClient.GetTableMetadata({ - connection_id, - customer_name, - plugin, - schema, - table_list, + const tables_metadata = await Promise.all( + body.table_list.map(async (table_name) => { + const result = await this.connectionsApiService.proxy( + 'GET', + `/connection_catalog/${encodeURIComponent(body.connection_id)}` + + `/schemas/${encodeURIComponent(body.schema)}` + + `/tables/${encodeURIComponent(table_name)}/columns`, + user, + ); + return { + table_name, + columns: result.columns.map((column) => ({ + name: column.column_name, + type: column.data_type, + is_primary_key: column.is_primary_key, + })), + references: [], + }; }), ); + return { operation_result: true, tables_metadata }; + } + + async refreshCatalog( + body: RefreshCatalogReq, + user: RequestUser, + ): Promise { + const task = await this.platformApiService.proxy( + 'POST', + '/connection_test', + user, + { + customer_id: user.customer_name, + plugin: body.plugin, + task: { + task_type: 'refresh_catalog', + connection: { + provider: 'connection_manager', + config_id: body.connection_id, + }, + }, + }, + ); + + if (!task.session_id || !task.date) { + throw new HttpException( + 'Platform API did not return a catalog refresh task identifier', + HttpStatus.BAD_GATEWAY, + ); + } + + return { + operation_result: true, + status: 'PENDING', + session_id: task.session_id, + date: task.date, + }; + } + + async refreshCatalogStatus( + body: RefreshCatalogStatusReq, + user: RequestUser, + ): Promise { + const result = await this.platformApiService.proxy( + 'POST', + '/connection_test/status', + user, + { + session_id: body.session_id, + date: body.date, + }, + ); + + if (result.status === 'DONE') { + await this.connectionsApiService.proxy( + 'PUT', + `/connection_config/${encodeURIComponent( + body.connection_id, + )}/catalog_metadata`, + user, + { + last_catalog_refresh_status: 'SUCCESS', + last_catalog_connection_test_date: body.date, + last_catalog_connection_test_session_id: body.session_id, + }, + ); + } else if (result.status === 'ERROR' || result.status === 'EXPIRED') { + throw new HttpException( + `Catalog refresh finished with status ${result.status}`, + HttpStatus.BAD_GATEWAY, + ); + } + + return { + operation_result: result.status === 'DONE', + status: result.status, + session_id: body.session_id, + date: body.date, + }; } async validateCdcPrerequisites( diff --git a/src/modules/connection-test/dto/connection-test.ts b/src/modules/connection-test/dto/connection-test.ts index f7b1228..cb10c4c 100644 --- a/src/modules/connection-test/dto/connection-test.ts +++ b/src/modules/connection-test/dto/connection-test.ts @@ -1,5 +1,5 @@ import { ApiProperty, ApiPropertyOptional, OmitType } from '@nestjs/swagger'; -import { IsString, IsOptional } from 'class-validator'; +import { IsIn, IsString, IsOptional } from 'class-validator'; import { DatabaseConnectionPropertiesDto } from 'src/modules/connection/dtos/connection'; import { CreateConnectionDto } from 'src/modules/connection/dtos/connection'; export class ColumnDto { @@ -7,6 +7,8 @@ export class ColumnDto { name: string; @ApiProperty() type: string; + @ApiProperty() + is_primary_key: boolean; } export class TableMetadataDto { @ApiProperty() @@ -168,3 +170,37 @@ export class ValidateCdcPrerequisitesRes { @ApiProperty({ type: [CdcCheckDto] }) checks: CdcCheckDto[]; } + +export class RefreshCatalogReq { + @ApiProperty() + @IsString() + connection_id: string; + + @ApiProperty({ enum: ['oracle', 'mysql', 'postgresql', 'sqlserver'] }) + @IsIn(['oracle', 'mysql', 'postgresql', 'sqlserver']) + plugin: string; +} + +export class RefreshCatalogStatusReq extends RefreshCatalogReq { + @ApiProperty() + @IsString() + session_id: string; + + @ApiProperty() + @IsString() + date: string; +} + +export class RefreshCatalogRes { + @ApiProperty() + operation_result: boolean; + + @ApiProperty() + status: string; + + @ApiProperty() + session_id: string; + + @ApiProperty() + date: string; +} diff --git a/src/modules/connections-api/connections-api.config.ts b/src/modules/connections-api/connections-api.config.ts new file mode 100644 index 0000000..0e6e450 --- /dev/null +++ b/src/modules/connections-api/connections-api.config.ts @@ -0,0 +1,11 @@ +export const CONNECTIONS_API_CONFIG = { + getUrl: (): string => { + const url = process.env.CONNECTIONS_API_URL; + if (!url) { + throw new Error('CONNECTIONS_API_URL environment variable is not set'); + } + return url; + }, + region: process.env.AWS_REGION || 'us-east-1', + timeout: parseInt(process.env.CONNECTIONS_API_TIMEOUT || '30000', 10), +}; diff --git a/src/modules/connections-api/connections-api.module.ts b/src/modules/connections-api/connections-api.module.ts new file mode 100644 index 0000000..1835186 --- /dev/null +++ b/src/modules/connections-api/connections-api.module.ts @@ -0,0 +1,10 @@ +import { Module } from '@nestjs/common'; +import { DadosferaLogger } from '@dadosfera/dadosfera-logs'; + +import { ConnectionsApiService } from './connections-api.service'; + +@Module({ + providers: [ConnectionsApiService, DadosferaLogger], + exports: [ConnectionsApiService], +}) +export class ConnectionsApiModule {} diff --git a/src/modules/connections-api/connections-api.service.ts b/src/modules/connections-api/connections-api.service.ts new file mode 100644 index 0000000..ae2dfd4 --- /dev/null +++ b/src/modules/connections-api/connections-api.service.ts @@ -0,0 +1,99 @@ +import { Injectable, Inject, HttpException } from '@nestjs/common'; +import { SignatureV4 } from '@aws-sdk/signature-v4'; +import { Sha256 } from '@aws-crypto/sha256-js'; +import { defaultProvider } from '@aws-sdk/credential-provider-node'; +import axios, { AxiosResponse, Method } from 'axios'; +import { DadosferaLogger } from '@dadosfera/dadosfera-logs'; + +import { RequestUser } from '../../decorators/user.decorator'; +import { CONNECTIONS_API_CONFIG } from './connections-api.config'; + +@Injectable() +export class ConnectionsApiService { + private signer: SignatureV4; + private logger: any; + + constructor(@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger) { + this.logger = dadosferaLogger.logger; + this.signer = new SignatureV4({ + service: 'execute-api', + region: CONNECTIONS_API_CONFIG.region, + credentials: defaultProvider(), + sha256: Sha256, + }); + } + + async proxy( + method: string, + path: string, + user: RequestUser, + body?: any, + query?: Record, + ): Promise { + const baseUrl = CONNECTIONS_API_CONFIG.getUrl(); + const url = new URL(`${baseUrl}${path}`); + + if (query) { + Object.entries(query).forEach(([key, value]) => { + if (value !== undefined && value !== null) { + url.searchParams.set(key, String(value)); + } + }); + } + const headers: Record = { + host: url.hostname, + 'content-type': 'application/json', + customer_name: user.customer_name || '', + customer_id: user.customer_id || '', + 'x-user-id': user.user_id || '', + 'x-username': user.username || '', + 'x-customer-tier': user.customer_tier || '', + 'x-customer-id': user.customer_id || '', + }; + const requestToSign = { + method: method.toUpperCase(), + protocol: url.protocol, + hostname: url.hostname, + port: url.port ? parseInt(url.port, 10) : undefined, + path: url.pathname + url.search, + headers, + body: body ? JSON.stringify(body) : undefined, + }; + + try { + const signedRequest = await this.signer.sign(requestToSign); + const response: AxiosResponse = await axios({ + method: method as Method, + url: url.href, + headers: signedRequest.headers as Record, + data: body, + timeout: CONNECTIONS_API_CONFIG.timeout, + validateStatus: () => true, + }); + + if (response.status >= 400) { + throw new HttpException(response.data, response.status); + } + return response.data; + } catch (error) { + this.logger.error('Connections API proxy error', { + error: error.message, + path, + method: method.toUpperCase(), + }); + if (error instanceof HttpException) { + throw error; + } + if (error.response) { + throw new HttpException(error.response.data, error.response.status); + } + if (error.code === 'ECONNREFUSED') { + throw new HttpException('Connections API service unavailable', 503); + } + if (error.code === 'ETIMEDOUT' || error.code === 'ECONNABORTED') { + throw new HttpException('Connections API request timeout', 504); + } + throw new HttpException('Internal server error', 500); + } + } +} diff --git a/src/modules/platform-api/platform-api.controller.ts b/src/modules/platform-api/platform-api.controller.ts index 6ec8c45..44e1ac9 100644 --- a/src/modules/platform-api/platform-api.controller.ts +++ b/src/modules/platform-api/platform-api.controller.ts @@ -14,7 +14,7 @@ import { NotFoundException, UseGuards, } from '@nestjs/common'; -import { ApiTags, ApiOperation } from '@nestjs/swagger'; +import { ApiTags, ApiOperation, ApiOkResponse } from '@nestjs/swagger'; import { DadosferaLogger } from '@dadosfera/dadosfera-logs'; import { @@ -76,6 +76,10 @@ export class PlatformApiController { return id?.replace(/-/g, '_') || ''; } + private decodePathParam(value: string): string { + return value ? decodeURIComponent(value) : ''; + } + /** * Denormalize ID back to UUID format (replace _ with -). * Used when we receive a normalized ID but need the original UUID. @@ -850,6 +854,40 @@ export class PlatformApiController { ); } + @Get('pipelines/:pipelineId/pipeline_run/:runId/jobs') + @ApiOperation({ + summary: 'Get pipeline run jobs', + description: 'Proxies platform-api DB-backed job runs and returns `{ jobs: [...] }`.', + }) + @ApiOkResponse({ + description: 'DB-backed job runs for the selected pipeline run.', + schema: { + type: 'object', + properties: { + jobs: { + type: 'array', + items: { type: 'object' }, + }, + }, + required: ['jobs'], + }, + }) + @RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET) + async getPipelineRunJobs( + @Param('pipelineId') pipelineId: string, + @Param('runId') runId: string, + @User() user: RequestUser, + ) { + const normalizedPipelineId = this.normalizePipelineId(pipelineId); + const decodedRunId = this.decodePathParam(runId); + + return this.platformApiService.proxy( + 'GET', + `/pipeline/${normalizedPipelineId}/pipeline_run/${decodedRunId}/jobs`, + user, + ); + } + // ==================== JOBS - COLUMN EDITING ROUTES ==================== @Put('jobs/:jobId/input') diff --git a/src/modules/release_note/release_note.controller.spec.ts b/src/modules/release_note/release_note.controller.spec.ts index 6086196..921eca1 100644 --- a/src/modules/release_note/release_note.controller.spec.ts +++ b/src/modules/release_note/release_note.controller.spec.ts @@ -1,7 +1,6 @@ import { Test, TestingModule } from '@nestjs/testing'; import { ReleaseNoteController } from './release_note.controller'; import { ReleaseNoteService } from './release_note.service'; -import DadosferaLogger from '@dadosfera/dadosfera-logs'; describe('ReleaseNoteController', () => { let controller: ReleaseNoteController; @@ -9,7 +8,7 @@ describe('ReleaseNoteController', () => { beforeEach(async () => { const module: TestingModule = await Test.createTestingModule({ controllers: [ReleaseNoteController], - providers: [ReleaseNoteService, DadosferaLogger], + providers: [ReleaseNoteService], }).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 266146d..72add9d 100644 --- a/src/modules/release_note/release_note.service.spec.ts +++ b/src/modules/release_note/release_note.service.spec.ts @@ -1,13 +1,12 @@ import { Test, TestingModule } from '@nestjs/testing'; import { ReleaseNoteService } from './release_note.service'; -import DadosferaLogger from '@dadosfera/dadosfera-logs'; describe('ReleaseNoteService', () => { let service: ReleaseNoteService; beforeEach(async () => { const module: TestingModule = await Test.createTestingModule({ - providers: [ReleaseNoteService, DadosferaLogger], + providers: [ReleaseNoteService], }).compile(); service = module.get(ReleaseNoteService); diff --git a/src/utils/ErrorBuilder.ts b/src/utils/ErrorBuilder.ts index 49171ee..7e48cb5 100644 --- a/src/utils/ErrorBuilder.ts +++ b/src/utils/ErrorBuilder.ts @@ -315,6 +315,14 @@ export function EnrichErrorCode(code: string) { 'Tente realizar a ação novamente. Caso o erro persista, entre em contato com o suporte', code, }; + case ErrorCodes.CATALOG.COLUMN_NOT_FOUND: + return { + statusCode: HttpStatus.NOT_FOUND, + error: 'Coluna não encontrada', + message: + 'Uma ou mais colunas informadas não existem neste ativo. Verifique os nomes e tente novamente.', + code, + }; case ErrorCodes.CATALOG.PREVIEW_TOO_BIG: return { statusCode: HttpStatus.INTERNAL_SERVER_ERROR, diff --git a/src/utils/errorCodes.ts b/src/utils/errorCodes.ts index a5856bd..c6f41f9 100644 --- a/src/utils/errorCodes.ts +++ b/src/utils/errorCodes.ts @@ -74,6 +74,7 @@ const CATALOG = { DATA_ASSET_NOT_FOUND: 'CATALOG.DATA_ASSET_NOT_FOUND', PREVIEW_TOO_BIG: 'CATALOG.PREVIEW_TOO_BIG', METADATA_TOO_BIG: 'CATALOG.METADATA_TOO_BIG', + COLUMN_NOT_FOUND: 'CATALOG.COLUMN_NOT_FOUND', }; const IDENTITY_PROVIDER = { INVALID_RESPONSE: 'IDENTITY_PROVIDER.INVALID_RESPONSE',