FEAT: update job deletion endpoint to include inputId in path and implement rollback for table deletion

This commit is contained in:
viniciusgadea
2026-04-01 14:26:55 -03:00
parent 877cb9d281
commit 968f75b688
4 changed files with 35 additions and 36 deletions
@@ -965,17 +965,16 @@ export class PlatformApiController {
);
}
@Delete('jobs/:jobId')
@Delete('inputs/:inputId/jobs/:jobId')
@ApiOperation({ summary: 'Delete a job and mark its table as deleted via in-factory' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE)
async deleteJob(
@Param('inputId') inputId: string,
@Param('jobId') jobId: string,
@User() user: RequestUser,
) {
const normalizedJobId = this.normalizeJobId(jobId);
const esPipelineId = this.extractPipelineIdFromJobId(jobId);
// Fetch job details first — fail fast if job doesn't exist or jobId is invalid
const jobDetails = await this.platformApiService.proxy('GET', `/jobs/${normalizedJobId}`, user);
const tableName: string | undefined = jobDetails?.source_config?.table_name;
@@ -983,30 +982,28 @@ export class PlatformApiController {
throw new NotFoundException(`Job ${jobId} not found or has no associated table`);
}
// Fetch ES pipeline to resolve the DynamoDB input ID
const esPipeline = await this.elasticsearchService.getPipeline(user.customer_name, esPipelineId);
// config.tables in ES stores the DynamoDB input ID (legacy field naming)
const inputId: string | undefined = esPipeline?.config?.tables;
if (!inputId) {
throw new NotFoundException(`Pipeline data not found for job ${jobId}`);
}
// Delete the job on platform-api
const result = await this.platformApiService.proxy('DELETE', `/jobs/${normalizedJobId}`, user);
// Mark table as deleted in DynamoDB via in-factory gRPC
const info = {
customer_id: user.customer_id,
customer: user.customer_name,
user_id: user.user_id,
};
await this.inputsService.markTableDeleted({ input_id: inputId, table_name: tableName, info }).catch((error) => {
this.logger.error('deleteJob: failed to mark table as deleted via in-factory', { jobId, tableName, error: error.message });
});
const updatedInput: any = await this.inputsService.markTableDeleted({ input_id: inputId, table_name: tableName, info });
return result;
try {
await this.platformApiService.proxy('DELETE', `/jobs/${normalizedJobId}`, user);
const updatedTable = updatedInput.tables?.find((t: any) => t.name === tableName);
return updatedTable ?? { name: tableName, is_deleted: true };
} catch (error) {
this.logger.error('deleteJob: platform-api delete failed, attempting rollback', { jobId, tableName, error: error.message });
try {
await this.inputsService.unmarkTableDeleted({ input_id: inputId, table_name: tableName, info });
} catch (rollbackError) {
this.logger.error('deleteJob: rollback failed', { jobId, tableName, error: rollbackError.message });
}
throw error;
}
}
// ==================== JOBS - JDBC SYNC MODE ROUTES ====================