diff --git a/docsfera.json b/docsfera.json index a18dad5..3ae0f4c 100644 --- a/docsfera.json +++ b/docsfera.json @@ -42,6 +42,33 @@ ] } }, + "/auth/sign-out": { + "post": { + "operationId": "AuthController_signOut", + "parameters": [ + { + "name": "dadosfera-lang", + "in": "header", + "required": false, + "schema": { + "enum": [ + "pt-br", + "en-us" + ], + "type": "string" + } + } + ], + "responses": { + "204": { + "description": "" + } + }, + "tags": [ + "Auth" + ] + } + }, "/auth/refresh-access-token": { "post": { "operationId": "AuthController_refreshAccessToken", @@ -3415,9 +3442,6 @@ "PipelinesV2" ], "security": [ - { - "access-token": [] - }, { "access-token": [] } @@ -3467,9 +3491,6 @@ "PipelinesV2" ], "security": [ - { - "access-token": [] - }, { "access-token": [] } @@ -3519,9 +3540,6 @@ "PipelinesV2" ], "security": [ - { - "access-token": [] - }, { "access-token": [] } @@ -6760,6 +6778,863 @@ } } }, + "/platform/pipeline": { + "post": { + "operationId": "PlatformApiController_createPipeline", + "summary": "Create a new pipeline", + "parameters": [], + "responses": { + "201": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/pipelines": { + "get": { + "operationId": "PlatformApiController_getPipelines", + "summary": "List all pipelines for customer", + "parameters": [], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/pipeline/{pipelineId}": { + "get": { + "operationId": "PlatformApiController_getPipeline", + "summary": "Get pipeline by ID", + "parameters": [ + { + "name": "pipelineId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + }, + "patch": { + "operationId": "PlatformApiController_updatePipeline", + "summary": "Update pipeline by ID", + "parameters": [ + { + "name": "pipelineId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + }, + "delete": { + "operationId": "PlatformApiController_deletePipeline", + "summary": "Delete pipeline by ID", + "parameters": [ + { + "name": "pipelineId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/pipeline/execute": { + "post": { + "operationId": "PlatformApiController_executePipeline", + "summary": "Execute a pipeline", + "parameters": [], + "responses": { + "201": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/pipeline/pause": { + "post": { + "operationId": "PlatformApiController_pausePipeline", + "summary": "Pause a pipeline", + "parameters": [], + "responses": { + "201": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/pipeline/unpause": { + "post": { + "operationId": "PlatformApiController_unpausePipeline", + "summary": "Unpause a pipeline", + "parameters": [], + "responses": { + "201": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/pipeline/{pipelineId}/memory": { + "put": { + "operationId": "PlatformApiController_updatePipelineMemory", + "summary": "Update pipeline memory configuration", + "parameters": [ + { + "name": "pipelineId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/pipeline/{pipelineId}/metadata": { + "put": { + "operationId": "PlatformApiController_updatePipelineMetadata", + "summary": "Update pipeline metadata", + "parameters": [ + { + "name": "pipelineId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/pipelines/metadata": { + "get": { + "operationId": "PlatformApiController_getPipelinesMetadata", + "summary": "Get all pipelines metadata", + "parameters": [], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/pipeline/pipeline_run": { + "post": { + "operationId": "PlatformApiController_createPipelineRun", + "summary": "Create a pipeline run", + "parameters": [], + "responses": { + "201": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/pipeline/{pipelineId}/pipeline_run": { + "get": { + "operationId": "PlatformApiController_getPipelineRuns", + "summary": "Get pipeline runs for a pipeline", + "parameters": [ + { + "name": "pipelineId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/pipeline/{pipelineId}/pipeline_run/{runId}": { + "get": { + "operationId": "PlatformApiController_getPipelineRun", + "summary": "Get specific pipeline run", + "parameters": [ + { + "name": "pipelineId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + }, + { + "name": "runId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/pipeline/pipeline_run/{runId}/logs": { + "get": { + "operationId": "PlatformApiController_getPipelineRunLogs", + "summary": "Get pipeline run logs", + "parameters": [ + { + "name": "runId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/jobs/{jobId}/input": { + "put": { + "operationId": "PlatformApiController_updateJobInput", + "summary": "Update job input columns", + "parameters": [ + { + "name": "jobId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + }, + "patch": { + "operationId": "PlatformApiController_patchJobInput", + "summary": "Partial update job input columns", + "parameters": [ + { + "name": "jobId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/jobs/{jobId}/memory": { + "put": { + "operationId": "PlatformApiController_updateJobMemory", + "summary": "Update job memory configuration", + "parameters": [ + { + "name": "jobId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/jobs/{jobId}/reset-state": { + "post": { + "operationId": "PlatformApiController_resetJobState", + "summary": "Reset job state", + "parameters": [ + { + "name": "jobId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "201": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/jobs/jdbc/{jobId}": { + "get": { + "operationId": "PlatformApiController_getJdbcJob", + "summary": "Get JDBC job details", + "parameters": [ + { + "name": "jobId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/jobs/jdbc/{jobId}/sync-mode": { + "post": { + "operationId": "PlatformApiController_updateJdbcSyncMode", + "summary": "Update JDBC job sync mode", + "parameters": [ + { + "name": "jobId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "201": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/jobs/jdbc/configs/allowed_datatypes": { + "get": { + "operationId": "PlatformApiController_getJdbcAllowedDatatypes", + "summary": "Get allowed datatypes for JDBC", + "parameters": [], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/jobs/singer/{jobId}": { + "get": { + "operationId": "PlatformApiController_getSingerJob", + "summary": "Get Singer job details", + "parameters": [ + { + "name": "jobId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/jobs/singer/{jobId}/sync-mode": { + "post": { + "operationId": "PlatformApiController_updateSingerSyncMode", + "summary": "Update Singer job sync mode", + "parameters": [ + { + "name": "jobId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "201": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/jobs/s3/{jobId}": { + "get": { + "operationId": "PlatformApiController_getS3Job", + "summary": "Get S3 job details", + "parameters": [ + { + "name": "jobId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, + "/platform/health": { + "get": { + "operationId": "PlatformApiController_healthCheck", + "summary": "Platform API health check", + "parameters": [], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, "/health": { "get": { "operationId": "HealthController_check", diff --git a/src/main.ts b/src/main.ts index 70161e2..fc90db0 100644 --- a/src/main.ts +++ b/src/main.ts @@ -22,6 +22,7 @@ async function bootstrap() { logger, cors: { origin: [ + 'http://localhost:4200', 'https://app.stg.dadosfera.ai', 'https://app.dadosfera.ai', 'https://unimed.dadosfera.ai', diff --git a/src/modules/platform-api/platform-api.controller.ts b/src/modules/platform-api/platform-api.controller.ts index d9006ef..3be5530 100644 --- a/src/modules/platform-api/platform-api.controller.ts +++ b/src/modules/platform-api/platform-api.controller.ts @@ -49,7 +49,52 @@ export class PlatformApiController { return id?.replace(/-/g, '_') || ''; } + /** + * Denormalize ID back to UUID format (replace _ with -). + * Used when we receive a normalized ID but need the original UUID. + */ + private denormalizeId(id: string): string { + return id?.replace(/_/g, '-') || ''; + } + + /** + * Normalize job ID to match Platform-API format. + * Platform-API replaces '-' with '_' in job IDs. + * + * Example: "2ccf5481-59f5-4036-8a94-7d5f28f4f899-0" -> "2ccf5481_59f5_4036_8a94_7d5f28f4f899_0" + */ + private normalizeJobId(jobId: string): string { + return jobId?.replace(/-/g, '_') || ''; + } + + /** + * Extract the pipeline ID (base UUID) from a job ID. + * Job IDs have format "uuid-suffix" where suffix is the job index (e.g., "0", "1"). + * Handles both hyphenated and underscored formats, always returns hyphenated UUID for ES. + * + * Examples: + * - "2ccf5481-59f5-4036-8a94-7d5f28f4f899-0" -> "2ccf5481-59f5-4036-8a94-7d5f28f4f899" + * - "2ccf5481_59f5_4036_8a94_7d5f28f4f899_0" -> "2ccf5481-59f5-4036-8a94-7d5f28f4f899" + */ + private extractPipelineIdFromJobId(jobId: string): string { + if (!jobId) return ''; + + // Determine the separator used in the jobId + const hasUnderscores = jobId.includes('_'); + const separator = hasUnderscores ? '_' : '-'; + + const parts = jobId.split(separator); + // UUID has 5 parts (8-4-4-4-12), job suffix is the 6th part + if (parts.length >= 6) { + // Always return hyphenated format for Elasticsearch lookup + return parts.slice(0, 5).join('-'); + } + // If no suffix found, return the ID in hyphenated format + return hasUnderscores ? jobId.replace(/_/g, '-') : jobId; + } + private readonly VALID_CONNECTORS = ['jdbc', 'singer', 's3']; + private readonly MAX_MEMORY_MB = 12000; // 12GB maximum memory per pipeline/job /** * Validate that connector is provided and is a valid type. @@ -62,6 +107,17 @@ export class PlatformApiController { } } + /** + * Validate memory allocation against maximum limit. + */ + private validateMemory(memoryMb: number): void { + if (memoryMb > this.MAX_MEMORY_MB) { + throw new BadRequestException( + `Memory limit exceeded. Maximum allowed: ${this.MAX_MEMORY_MB}MB (12GB)`, + ); + } + } + /** * Validate cron expression against customer's schedule limit. * Fetches current scheduleLimit from DUC to ensure up-to-date configuration. @@ -107,7 +163,7 @@ export class PlatformApiController { name: string; type: string; columns?: string[]; - reference_column?: { name: string; type: string }; + reference_column?: string; }> { if (!jobs || jobs.length === 0) return []; @@ -115,7 +171,7 @@ export class PlatformApiController { name: string; type: string; columns?: string[]; - reference_column?: { name: string; type: string }; + reference_column?: string; }> = []; for (const job of jobs) { @@ -128,7 +184,7 @@ export class PlatformApiController { name: string; type: string; columns?: string[]; - reference_column?: { name: string; type: string }; + reference_column?: string; } = { name: input.table_name || '', type: input.load_type || 'full_load', @@ -137,10 +193,8 @@ export class PlatformApiController { table.columns = input.column_include_list; } if (input.incremental_column_name) { - table.reference_column = { - name: input.incremental_column_name, - type: input.incremental_column_type || 'timestamp', - }; + // reference_column is stored as a string (column name) + table.reference_column = input.incremental_column_name; } tables.push(table); } else if (connector === 'singer' || connector === 's3') { @@ -205,6 +259,148 @@ export class PlatformApiController { return properties; } + /** + * Sync job input changes to DynamoDB for a specific connector type. + * Extracts pipeline ID from job ID, fetches ES document to find input ID, + * then updates the table entry in DynamoDB. + * + * Job ID transformations: + * - Raw format (from endpoint): "2ccf5481-59f5-4036-8a94-7d5f28f4f899-0" + * - Platform API format: "2ccf5481_59f5_4036_8a94_7d5f28f4f899_0" (underscores) + * - Elasticsearch pipeline ID: "2ccf5481-59f5-4036-8a94-7d5f28f4f899" (UUID only, hyphens) + * + * @param connectorType - The connector type ('jdbc', 'singer', 's3') for the Platform API endpoint + */ + private async syncJobInputToDynamoDB( + jobId: string, + body: any, + user: RequestUser, + connectorType: 'jdbc' | 'singer' | 's3', + ): Promise { + try { + // Normalize job ID for Platform API GET (replace - with _) + const normalizedJobId = this.normalizeJobId(jobId); + + // Get job details using connector-specific endpoint to find table_name + const jobResult = await this.platformApiService.proxy( + 'GET', + `/jobs/${connectorType}/${normalizedJobId}`, + user, + ); + + // Extract the pipeline ID (base UUID) from the raw job ID for ES lookup + const esPipelineId = this.extractPipelineIdFromJobId(jobId); + const tableName = body.table_name || jobResult.source_config?.table_name; + + if (!esPipelineId || !tableName) { + this.logger.warn('Cannot sync job input: missing pipeline_id or table_name', { + jobId, + esPipelineId, + tableName, + }); + return; + } + + // Get pipeline from ES to find input ID (stored in config.tables) + const pipeline = await this.elasticsearchService.getPipeline( + user.customer_name, + esPipelineId, + ); + + const inputId = pipeline?.config?.tables; + if (!inputId) { + this.logger.warn('Cannot sync job input: no input ID in ES', { + jobId, + esPipelineId, + }); + return; + } + + // Build changes for DynamoDB table entry + // reference_column is stored as a string (column name), not an object + const changes: { + type?: string; + columns?: string[]; + reference_column?: string | null; + } = {}; + + + if ('target_load_type' in body) { + changes.type = body.target_load_type; + } + if ('column_include_list' in body) { + changes.columns = body.column_include_list; + } + if ('incremental_column_name' in body) { + // reference_column is just the column name as a string + changes.reference_column = body.incremental_column_name || null; + } + + // Update DynamoDB if there are changes + if (Object.keys(changes).length > 0) { + await this.dynamoDBService.updateInputTable( + user.customer_id, + inputId, + tableName, + changes, + ); + } + } catch (error) { + this.logger.error('Failed to sync job input to DynamoDB', { + jobId, + connectorType, + error: error.message, + }); + // Don't throw - Platform API update succeeded, just log the sync error + } + } + + /** + * Sync sync-mode changes to DynamoDB for JDBC connectors. + * Always passes both target_load_type and incremental_column_name to ensure proper sync. + */ + private async syncJdbcSyncModeToDynamoDB( + jobId: string, + body: any, + user: RequestUser, + ): Promise { + // JDBC sync mode uses target_load_type field + const changes: any = {}; + + if ('target_load_type' in body) { + changes.target_load_type = body.target_load_type; + } + + // Handle incremental_column_name: + // - If provided in body, use that value + // - If changing to full_load, explicitly clear it + if ('incremental_column_name' in body) { + changes.incremental_column_name = body.incremental_column_name; + changes.incremental_column_type = body.incremental_column_type; + } else if (body.target_load_type === 'full_load') { + // Changing to full_load without specifying incremental_column - clear it + changes.incremental_column_name = null; + } + + await this.syncJobInputToDynamoDB(jobId, changes, user, 'jdbc'); + } + + /** + * Sync sync-mode changes to DynamoDB for Singer connectors. + */ + private async syncSingerSyncModeToDynamoDB( + jobId: string, + body: any, + user: RequestUser, + ): Promise { + // Singer sync mode uses replication_method field + // Map to DynamoDB type: FULL_TABLE -> full_load, INCREMENTAL -> incremental + if ('replication_method' in body) { + const type = body.replication_method === 'INCREMENTAL' ? 'incremental' : 'full_load'; + await this.syncJobInputToDynamoDB(jobId, { load_type: type }, user, 'singer'); + } + } + // ==================== PIPELINE ROUTES ==================== @Post('pipeline') @@ -446,6 +642,11 @@ export class PlatformApiController { @Body() body: any, @User() user: RequestUser, ) { + // Validate memory limit + if (body.amount) { + this.validateMemory(body.amount); + } + return this.platformApiService.proxy( 'PUT', `/pipeline/${pipelineId}/memory`, @@ -561,12 +762,23 @@ export class PlatformApiController { @Body() body: any, @User() user: RequestUser, ) { - return this.platformApiService.proxy( + // Normalize job ID for Platform API (replace - with _) + const normalizedJobId = this.normalizeJobId(jobId); + + const result = await this.platformApiService.proxy( 'PUT', - `/jobs/${jobId}/input`, + `/jobs/${normalizedJobId}/input`, user, body, ); + + // Sync to DynamoDB if connector type is provided + const connectorType = body.connector as 'jdbc' | 'singer' | 's3' | undefined; + if (connectorType && this.VALID_CONNECTORS.includes(connectorType)) { + await this.syncJobInputToDynamoDB(jobId, body, user, connectorType); + } + + return result; } @Patch('jobs/:jobId/input') @@ -577,12 +789,23 @@ export class PlatformApiController { @Body() body: any, @User() user: RequestUser, ) { - return this.platformApiService.proxy( + // Normalize job ID for Platform API (replace - with _) + const normalizedJobId = this.normalizeJobId(jobId); + + const result = await this.platformApiService.proxy( 'PATCH', - `/jobs/${jobId}/input`, + `/jobs/${normalizedJobId}/input`, user, body, ); + + // Sync to DynamoDB if connector type is provided + const connectorType = body.connector as 'jdbc' | 'singer' | 's3' | undefined; + if (connectorType && this.VALID_CONNECTORS.includes(connectorType)) { + await this.syncJobInputToDynamoDB(jobId, body, user, connectorType); + } + + return result; } @Put('jobs/:jobId/memory') @@ -593,9 +816,17 @@ export class PlatformApiController { @Body() body: any, @User() user: RequestUser, ) { + // Validate memory limit + if (body.amount) { + this.validateMemory(body.amount); + } + + // Normalize job ID for Platform API (replace - with _) + const normalizedJobId = this.normalizeJobId(jobId); + return this.platformApiService.proxy( 'PUT', - `/jobs/${jobId}/memory`, + `/jobs/${normalizedJobId}/memory`, user, body, ); @@ -609,9 +840,12 @@ export class PlatformApiController { @Body() body: any, @User() user: RequestUser, ) { + // Normalize job ID for Platform API (replace - with _) + const normalizedJobId = this.normalizeJobId(jobId); + return this.platformApiService.proxy( 'POST', - `/jobs/${jobId}/reset-state`, + `/jobs/${normalizedJobId}/reset-state`, user, body, ); @@ -623,10 +857,12 @@ export class PlatformApiController { @ApiOperation({ summary: 'Get JDBC job details' }) @RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET) async getJdbcJob(@Param('jobId') jobId: string, @User() user: RequestUser) { - return this.platformApiService.proxy('GET', `/jobs/jdbc/${jobId}`, user); + // Normalize job ID for Platform API (replace - with _) + const normalizedJobId = this.normalizeJobId(jobId); + return this.platformApiService.proxy('GET', `/jobs/jdbc/${normalizedJobId}`, user); } - @Put('jobs/jdbc/:jobId/sync-mode') + @Post('jobs/jdbc/:jobId/sync-mode') @ApiOperation({ summary: 'Update JDBC job sync mode' }) @RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE) async updateJdbcSyncMode( @@ -634,12 +870,20 @@ export class PlatformApiController { @Body() body: any, @User() user: RequestUser, ) { - return this.platformApiService.proxy( - 'PUT', - `/jobs/jdbc/${jobId}/sync-mode`, + // Normalize job ID for Platform API (replace - with _) + const normalizedJobId = this.normalizeJobId(jobId); + + const result = await this.platformApiService.proxy( + 'POST', + `/jobs/jdbc/${normalizedJobId}/sync-mode`, user, body, ); + + // Sync to DynamoDB (pass raw jobId for pipeline extraction) + await this.syncJdbcSyncModeToDynamoDB(jobId, body, user); + + return result; } @Get('jobs/jdbc/configs/allowed_datatypes') @@ -659,10 +903,12 @@ export class PlatformApiController { @ApiOperation({ summary: 'Get Singer job details' }) @RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET) async getSingerJob(@Param('jobId') jobId: string, @User() user: RequestUser) { - return this.platformApiService.proxy('GET', `/jobs/singer/${jobId}`, user); + // Normalize job ID for Platform API (replace - with _) + const normalizedJobId = this.normalizeJobId(jobId); + return this.platformApiService.proxy('GET', `/jobs/singer/${normalizedJobId}`, user); } - @Put('jobs/singer/:jobId/sync-mode') + @Post('jobs/singer/:jobId/sync-mode') @ApiOperation({ summary: 'Update Singer job sync mode' }) @RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE) async updateSingerSyncMode( @@ -670,12 +916,20 @@ export class PlatformApiController { @Body() body: any, @User() user: RequestUser, ) { - return this.platformApiService.proxy( - 'PUT', - `/jobs/singer/${jobId}/sync-mode`, + // Normalize job ID for Platform API (replace - with _) + const normalizedJobId = this.normalizeJobId(jobId); + + const result = await this.platformApiService.proxy( + 'POST', + `/jobs/singer/${normalizedJobId}/sync-mode`, user, body, ); + + // Sync to DynamoDB (pass raw jobId for pipeline extraction) + await this.syncSingerSyncModeToDynamoDB(jobId, body, user); + + return result; } // ==================== JOBS - S3 ROUTES ==================== @@ -684,7 +938,9 @@ export class PlatformApiController { @ApiOperation({ summary: 'Get S3 job details' }) @RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET) async getS3Job(@Param('jobId') jobId: string, @User() user: RequestUser) { - return this.platformApiService.proxy('GET', `/jobs/s3/${jobId}`, user); + // Normalize job ID for Platform API (replace - with _) + const normalizedJobId = this.normalizeJobId(jobId); + return this.platformApiService.proxy('GET', `/jobs/s3/${normalizedJobId}`, user); } // ==================== HEALTH ROUTE ==================== diff --git a/src/services/dynamodb/dynamodb.service.ts b/src/services/dynamodb/dynamodb.service.ts index aead208..099eb99 100644 --- a/src/services/dynamodb/dynamodb.service.ts +++ b/src/services/dynamodb/dynamodb.service.ts @@ -24,7 +24,7 @@ export interface InputDocument { name: string; type: string; columns?: string[]; - reference_column?: { name: string; type: string }; + reference_column?: string; }>; credentials?: Record; } @@ -62,7 +62,7 @@ export class DynamoDBService { name: string; type: string; columns?: string[]; - reference_column?: { name: string; type: string }; + reference_column?: string; }>; }, ): Promise { @@ -145,4 +145,85 @@ export class DynamoDBService { this.logger.info('DynamoDB: Input deleted successfully', { inputId }); } + + /** + * Update a specific table entry in the input document. + * Fetches the current document, updates the matching table, and saves. + */ + async updateInputTable( + clientId: string, + inputId: string, + tableName: string, + changes: { + type?: string; + columns?: string[]; + reference_column?: string | null; + }, + ): Promise { + const dynamoTableName = DYNAMODB_CONFIG.inputsTable(); + + this.logger.info('DynamoDB: Updating input table', { + inputId, + tableName, + changes: Object.keys(changes), + }); + + // Get current document + const current = await this.findInput(clientId, inputId); + if (!current) { + this.logger.warn('DynamoDB: Input not found for update', { inputId }); + return; + } + + // Find and update the matching table + const tables = current.tables || []; + const tableIndex = tables.findIndex((t) => t.name === tableName); + + if (tableIndex === -1) { + this.logger.warn('DynamoDB: Table not found in input', { + inputId, + tableName, + }); + return; + } + + // Merge changes into the table entry + const updatedTable = { ...tables[tableIndex] }; + if ('type' in changes) updatedTable.type = changes.type; + if ('columns' in changes) updatedTable.columns = changes.columns; + if ('reference_column' in changes) { + if (changes.reference_column === null) { + delete updatedTable.reference_column; + } else { + updatedTable.reference_column = changes.reference_column; + } + } + + tables[tableIndex] = updatedTable; + + // Save updated document + const putCommand = new PutCommand({ + TableName: dynamoTableName, + Item: { + ...current, + tables, + updated_at: new Date().toISOString(), + }, + }); + + try { + await this.documentClient.send(putCommand); + this.logger.info('DynamoDB: Input table updated successfully', { + inputId, + tableName, + }); + } catch (error) { + this.logger.error('DynamoDB: Failed to update input table', { + inputId, + tableName, + error: error.message, + }); + throw error; + } + } } diff --git a/src/services/elasticsearch/elasticsearch.service.ts b/src/services/elasticsearch/elasticsearch.service.ts index 281b5d6..e2c1857 100644 --- a/src/services/elasticsearch/elasticsearch.service.ts +++ b/src/services/elasticsearch/elasticsearch.service.ts @@ -293,6 +293,33 @@ export class ElasticsearchService { } } + async getPipeline( + customerName: string, + pipelineId: string, + ): Promise { + const index = this.getIndex(customerName); + + this.logger.info('Elasticsearch: Getting pipeline', { + index, + pipelineId, + }); + + try { + const response = await this.client.get(`${index}/_doc/${pipelineId}`); + return response.data._source as PipelineDocument; + } catch (error) { + if (error instanceof AxiosError && error.response?.status === 404) { + this.logger.warn('Elasticsearch: Pipeline not found', { + pipelineId, + index, + }); + return null; + } + this.handleError('getPipeline', error, { pipelineId, index }); + throw error; + } + } + async deletePipeline( customerName: string, pipelineId: string,