FEAT: add delete job endpoint and update related services; update protospack-v2 version to 3.40.0-beta.3

This commit is contained in:
viniciusgadea
2026-03-31 10:54:50 -03:00
parent 271174176b
commit 2d46ac3213
7 changed files with 1199 additions and 122 deletions
@@ -962,6 +962,80 @@ export class PlatformApiController {
);
}
@Delete('jobs/:jobId')
@ApiOperation({ summary: 'Delete a job and sync pipeline data' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE)
async deleteJob(
@Param('jobId') jobId: string,
@User() user: RequestUser,
) {
// Get job details before deletion to extract table_name and connector type
const normalizedJobId = this.normalizeJobId(jobId);
let tableName: string | undefined;
this.logger.info('[deleteJob] Start', { jobId });
try {
const job = await this.platformApiService.proxy(
'GET',
`/jobs/${normalizedJobId}`,
user,
);
tableName = job?.source_config?.table_name;
this.logger.info('[deleteJob] Job fetched', { jobId, tableName });
} catch (err) {
this.logger.warn('[deleteJob] Could not fetch job before deletion', { jobId, error: err.message });
}
const result = await this.platformApiService.proxy(
'DELETE',
`/jobs/${normalizedJobId}`,
user,
);
this.logger.info('[deleteJob] Job deleted from platform-api', { jobId });
// Mark the table as deleted via in-factory (which has the correct DynamoDB credentials).
// config.tables in ES stores the DynamoDB input ID (legacy naming).
if (tableName) {
try {
const esPipelineId = this.extractPipelineIdFromJobId(jobId);
const esPipeline = await this.elasticsearchService.getPipeline(
user.customer_name,
esPipelineId,
);
// config.tables holds the DynamoDB input ID (see syncJobInputToDynamoDB for context)
const dynamoInputId = esPipeline?.config?.tables;
if (dynamoInputId) {
const info = { customer_id: user.customer_id, user_id: user.user_id, customer: user.customer_name };
// Read current input to get full tables array (adjustInputPayload flattens the response)
const { input } = await this.inputsService.findOne({ id: dynamoInputId, info });
const currentTables: any[] = input?.tables ?? [];
const updatedTables = currentTables.map((t: any) =>
t.name === tableName ? { ...t, status: 'deleted' } : t,
);
await this.inputsService.update(dynamoInputId, { tables: updatedTables }, info);
this.logger.info('[deleteJob] Marked table as deleted via in-factory', { jobId, tableName });
} else {
this.logger.warn('[deleteJob] No input ID in ES pipeline, skipping DynamoDB update', { jobId, tableName });
}
} catch (error) {
this.logger.error('[deleteJob] Failed to mark table as deleted via in-factory', {
jobId,
tableName,
error: error.message,
});
}
} else {
this.logger.warn('[deleteJob] tableName not resolved, skipping DynamoDB update', { jobId });
}
return result;
}
// ==================== JOBS - JDBC SYNC MODE ROUTES ====================
@Get('jobs/jdbc/:jobId')