diff --git a/docsfera.json b/docsfera.json index c78fbec..7fb3e1a 100644 --- a/docsfera.json +++ b/docsfera.json @@ -4705,11 +4705,19 @@ ] } }, - "/platform/jobs/{jobId}": { + "/platform/inputs/{inputId}/jobs/{jobId}": { "delete": { "operationId": "PlatformApiController_deleteJob", "summary": "Delete a job and mark its table as deleted via in-factory", "parameters": [ + { + "name": "inputId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + }, { "name": "jobId", "required": true, diff --git a/src/modules/inputs/inputs.service.ts b/src/modules/inputs/inputs.service.ts index 5ead744..96df891 100644 --- a/src/modules/inputs/inputs.service.ts +++ b/src/modules/inputs/inputs.service.ts @@ -292,4 +292,8 @@ export class InputsService { async markTableDeleted(data: { input_id: string; table_name: string; info: Info }) { return lastValueFrom(this.inputWriteService.MarkTableDeleted(data)); } + + async unmarkTableDeleted(data: { input_id: string; table_name: string; info: Info }) { + return lastValueFrom((this.inputWriteService as any).UnmarkTableDeleted(data)); + } } diff --git a/src/modules/pipelinesV2/pipelines.controller.ts b/src/modules/pipelinesV2/pipelines.controller.ts index 37b3a47..71091ad 100644 --- a/src/modules/pipelinesV2/pipelines.controller.ts +++ b/src/modules/pipelinesV2/pipelines.controller.ts @@ -48,7 +48,6 @@ import { import { GrpcToHttpExceptionFilter } from 'src/error/grpc-to-http-exception.filter'; import { LanguageEnum } from 'src/utils/languages.enum'; import { Language } from 'src/decorators/language.decorator'; -import { PlatformApiService } from '../platform-api/platform-api.service'; import { ApiInternalOnlyEndpoint } from 'src/decorators/swagger.decorator'; import { TableColumns } from '../inputs/dtos/input.model'; import { UpdateInputRequest } from '../inputs/dtos/old_interfaces'; @@ -69,7 +68,6 @@ export class PipelinesController { private pipelinesClientService: PipelinesService, private oldPipelinesService: OldPipelineService, - private platformApiService: PlatformApiService, ) { this.logger = dadosferaLogger.logger; } @@ -228,22 +226,14 @@ export class PipelinesController { language, }); - const [pipelineRes, platformPipeline] = await Promise.all([ - this.pipelinesClientService.findOne({ id }, metadata), - this.platformApiService.proxy('GET', `/pipeline/${id.replace(/-/g, '_')}`, user).catch(() => null), - ]); + const pipelineRes = await this.pipelinesClientService.findOne({ id }, metadata); const parsed: PipelineTablesConfig = JSON.parse(pipelineRes.pipeline.config.tables); const input_id = (parsed as { input_id?: string; tables?: PipelineTable[] })?.input_id; - let tables: PipelineTable[] = (parsed as { tables?: PipelineTable[] })?.tables ?? (parsed as PipelineTable[]); - - if (platformPipeline?.jobs?.length) { - tables = tables.map((table) => { - const job = (platformPipeline.jobs as Array<{ job_id: string; input?: { table_name: string } }>) - .find((j) => j.input?.table_name === table.name); - return job ? { ...table, job_id: job.job_id } : table; - }); - } + const normalizedId = id.replace(/-/g, '_'); + const tables: PipelineTable[] = ( + (parsed as { tables?: PipelineTable[] })?.tables ?? (parsed as PipelineTable[]) + ).map((table, index) => ({ ...table, job_id: `${normalizedId}_${index}` })); Object.assign(pipelineRes.pipeline, { transformations: pipelineRes.pipeline.transformations diff --git a/src/modules/platform-api/platform-api.controller.ts b/src/modules/platform-api/platform-api.controller.ts index efe6592..b880a9b 100644 --- a/src/modules/platform-api/platform-api.controller.ts +++ b/src/modules/platform-api/platform-api.controller.ts @@ -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 ====================