FEAT: guard to prevent pipeline update when pipeline is running

This commit is contained in:
marcos-silva-rodrigues
2026-04-13 14:23:44 -03:00
parent ff8eb3197d
commit 0d627b0451
4 changed files with 99 additions and 295 deletions
@@ -12,6 +12,7 @@ import {
BadRequestException,
HttpException,
NotFoundException,
UseGuards,
} from '@nestjs/common';
import { ApiTags, ApiOperation } from '@nestjs/swagger';
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
@@ -31,6 +32,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 { PipelineExecutionGuard } from 'src/guards/pipeline-execution.guard';
type ValidateTablesDTO = {
@@ -132,7 +134,7 @@ export class PlatformApiController {
}
private readonly VALID_CONNECTORS = ['jdbc', 'singer', 's3'];
private readonly MAX_MEMORY_MB = 12000; // 12GB maximum memory per pipeline/job
private readonly MAX_MEMORY_MB = 12000; // 12GB maximum memory per pipelines/job
/**
* Validate that connector is provided and is a valid type.
@@ -403,55 +405,9 @@ export class PlatformApiController {
}
}
/**
* Sync sync-mode changes to DynamoDB for JDBC connectors.
* Always passes both target_load_type and incremental_column_name to ensure proper sync.
*/
private async syncJdbcSyncModeToDynamoDB(
jobId: string,
body: any,
user: RequestUser,
): Promise<void> {
// JDBC sync mode uses target_load_type field
const changes: any = {};
if ('target_load_type' in body) {
changes.target_load_type = body.target_load_type;
}
// Handle incremental_column_name:
// - If provided in body, use that value
// - If changing to full_load, explicitly clear it
if ('incremental_column_name' in body) {
changes.incremental_column_name = body.incremental_column_name;
changes.incremental_column_type = body.incremental_column_type;
} else if (body.target_load_type === 'full_load') {
// Changing to full_load without specifying incremental_column - clear it
changes.incremental_column_name = null;
}
await this.syncJobInputToDynamoDB(jobId, changes, user, 'jdbc');
}
/**
* Sync sync-mode changes to DynamoDB for Singer connectors.
*/
private async syncSingerSyncModeToDynamoDB(
jobId: string,
body: any,
user: RequestUser,
): Promise<void> {
// Singer sync mode uses replication_method field
// Map to DynamoDB type: FULL_TABLE -> full_load, INCREMENTAL -> incremental
if ('replication_method' in body) {
const type = body.replication_method === 'INCREMENTAL' ? 'incremental' : 'full_load';
await this.syncJobInputToDynamoDB(jobId, { load_type: type }, user, 'singer');
}
}
// ==================== PIPELINE ROUTES ====================
@Post('pipeline')
@Post('pipelines')
@ApiOperation({ summary: 'Create a new pipeline' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.CREATE)
async createPipeline(@Body() body: any, @User() user: RequestUser) {
@@ -572,7 +528,7 @@ export class PlatformApiController {
);
}
@Get('pipeline/:pipelineId')
@Get('pipelines/:pipelineId')
@ApiOperation({ summary: 'Get pipeline by ID' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
async getPipeline(
@@ -583,7 +539,7 @@ export class PlatformApiController {
return this.platformApiService.proxy('GET', `/pipeline/${normalizedId}`, user);
}
@Patch('pipeline/:pipelineId')
@Patch('pipelines/:pipelineId')
@ApiOperation({ summary: 'Update pipeline by ID' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async updatePipeline(
@@ -634,7 +590,7 @@ export class PlatformApiController {
return result;
}
@Delete('pipeline/:pipelineId')
@Delete('pipelines/:pipelineId')
@ApiOperation({ summary: 'Delete pipeline by ID' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE)
async deletePipeline(
@@ -664,7 +620,7 @@ export class PlatformApiController {
return result;
}
@Post('pipeline/execute')
@Post('pipelines/execute')
@ApiOperation({ summary: 'Execute a pipeline' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async executePipeline(@Body() body: any, @User() user: RequestUser) {
@@ -674,10 +630,10 @@ export class PlatformApiController {
...body,
customer_id: user.customer_name,
};
return this.platformApiService.proxy('POST', '/pipeline/execute', user, enrichedBody);
return this.platformApiService.proxy('POST', '/pipelines/execute', user, enrichedBody);
}
@Post('pipeline/pause')
@Post('pipelines/pause')
@ApiOperation({ summary: 'Pause a pipeline' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async pausePipeline(@Body() body: any, @User() user: RequestUser) {
@@ -690,7 +646,7 @@ export class PlatformApiController {
return this.platformApiService.proxy('POST', '/pipeline/pause', user, enrichedBody);
}
@Post('pipeline/unpause')
@Post('pipelines/unpause')
@ApiOperation({ summary: 'Unpause a pipeline' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async unpausePipeline(@Body() body: any, @User() user: RequestUser) {
@@ -703,7 +659,7 @@ export class PlatformApiController {
return this.platformApiService.proxy('POST', '/pipeline/unpause', user, enrichedBody);
}
@Put('pipeline/:pipelineId/memory')
@Put('pipelines/:pipelineId/memory')
@ApiOperation({ summary: 'Update pipeline memory configuration' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async updatePipelineMemory(
@@ -726,7 +682,7 @@ export class PlatformApiController {
// ==================== PIPELINE METADATA ROUTES ====================
@Put('pipeline/:pipelineId/metadata')
@Put('pipelines/:pipelineId/metadata')
@ApiOperation({ summary: 'Update pipeline metadata' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async updatePipelineMetadata(
@@ -795,7 +751,7 @@ export class PlatformApiController {
// ==================== PIPELINE RUN ROUTES ====================
@Get('pipeline/:pipelineId/pipeline_run')
@Get('pipelines/:pipelineId/pipeline_run')
@ApiOperation({ summary: 'Get pipeline runs for a pipeline' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
async getPipelineRuns(
@@ -813,7 +769,7 @@ export class PlatformApiController {
);
}
@Get('pipeline/:pipelineId/pipeline_run/:runId')
@Get('pipelines/:pipelineId/pipeline_run/:runId')
@ApiOperation({ summary: 'Get specific pipeline run' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
async getPipelineRun(
@@ -830,7 +786,7 @@ export class PlatformApiController {
);
}
@Get('pipeline/pipeline_run/:runId/logs')
@Get('pipelines/pipeline_run/:runId/logs')
@ApiOperation({ summary: 'Get pipeline run logs' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
async getPipelineRunLogs(
@@ -848,9 +804,10 @@ export class PlatformApiController {
);
}
@Post('pipeline/:pipelineId/pipeline_run/:runId/cancel')
@Post('pipelines/:pipelineId/pipeline_run/:runId/cancel')
@ApiOperation({ summary: 'Cancel a running pipeline run' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
@UseGuards(PipelineExecutionGuard)
async cancelPipelineRun(
@Param('pipelineId') pipelineId: string,
@Param('runId') runId: string,
@@ -968,6 +925,7 @@ export class PlatformApiController {
@Delete('pipelines/:pipelineId/inputs/:inputId')
@ApiOperation({ summary: 'Mark a table as deleted and delete its associated job via platform-api' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE)
@UseGuards(PipelineExecutionGuard)
async deleteTable(
@Param('pipelineId') pipelineId: string,
@Param('inputId') inputId: string,