mirror of
https://github.com/dadosfera/maestro.git
synced 2026-09-30 13:49:08 +00:00
FEAT: update job deletion and mark associated table as deleted in DynamoDB. Update protospack-v2
This commit is contained in:
@@ -11,6 +11,7 @@ import {
|
||||
Inject,
|
||||
BadRequestException,
|
||||
HttpException,
|
||||
NotFoundException,
|
||||
} from '@nestjs/common';
|
||||
import { ApiTags, ApiOperation } from '@nestjs/swagger';
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
@@ -29,6 +30,7 @@ import { validateCronAgainstScheduleLimit } from '../../utils/cron-validation';
|
||||
import { CatalogService } from '../catalog/catalog.service';
|
||||
import { PackTheMetadata } from '../../utils/PackTheMetadata';
|
||||
import { ValidationTableDTO } from './platform-api.dto';
|
||||
import { InputsService } from '../inputs/inputs.service';
|
||||
|
||||
|
||||
type ValidateTablesDTO = {
|
||||
@@ -54,6 +56,7 @@ export class PlatformApiController {
|
||||
private readonly dynamoDBService: DynamoDBService,
|
||||
private readonly customersService: CustomersService,
|
||||
private readonly catalogService: CatalogService,
|
||||
private readonly inputsService: InputsService,
|
||||
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
|
||||
) {
|
||||
this.logger = dadosferaLogger.logger;
|
||||
@@ -963,7 +966,7 @@ export class PlatformApiController {
|
||||
}
|
||||
|
||||
@Delete('jobs/:jobId')
|
||||
@ApiOperation({ summary: 'Delete a job and mark its table as deleted in DynamoDB' })
|
||||
@ApiOperation({ summary: 'Delete a job and mark its table as deleted via in-factory' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE)
|
||||
async deleteJob(
|
||||
@Param('jobId') jobId: string,
|
||||
@@ -972,47 +975,36 @@ export class PlatformApiController {
|
||||
const normalizedJobId = this.normalizeJobId(jobId);
|
||||
const esPipelineId = this.extractPipelineIdFromJobId(jobId);
|
||||
|
||||
// Fetch job details and ES pipeline in parallel — both are reads with no mutual dependency
|
||||
const [jobResult, esPipelineResult] = await Promise.allSettled([
|
||||
this.platformApiService.proxy('GET', `/jobs/${normalizedJobId}`, user),
|
||||
this.elasticsearchService.getPipeline(user.customer_name, esPipelineId),
|
||||
]);
|
||||
// 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;
|
||||
|
||||
const tableName: string | undefined = jobResult.status === 'fulfilled'
|
||||
? jobResult.value?.source_config?.table_name
|
||||
: undefined;
|
||||
|
||||
if (jobResult.status === 'rejected') {
|
||||
this.logger.warn('deleteJob: could not fetch job before deletion', { jobId, error: jobResult.reason?.message });
|
||||
if (!tableName) {
|
||||
throw new NotFoundException(`Job ${jobId} not found or has no associated table`);
|
||||
}
|
||||
|
||||
// Main operation
|
||||
// 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);
|
||||
|
||||
// Sync deleted status to DynamoDB — follows syncJobInputToDynamoDB pattern
|
||||
// config.tables in ES stores the DynamoDB input ID (legacy field naming)
|
||||
const dynamoInputId: string | undefined = esPipelineResult.status === 'fulfilled'
|
||||
? esPipelineResult.value?.config?.tables
|
||||
: undefined;
|
||||
// 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,
|
||||
};
|
||||
|
||||
if (esPipelineResult.status === 'rejected') {
|
||||
this.logger.warn('deleteJob: could not fetch ES pipeline for DynamoDB sync', { jobId, error: esPipelineResult.reason?.message });
|
||||
}
|
||||
|
||||
if (tableName && dynamoInputId) {
|
||||
await this.dynamoDBService.updateInputTable(
|
||||
user.customer_id,
|
||||
dynamoInputId,
|
||||
tableName,
|
||||
{ status: 'deleted' },
|
||||
).catch((error) => {
|
||||
this.logger.error('Failed to sync job deletion to DynamoDB', { jobId, tableName, error: error.message });
|
||||
});
|
||||
} else if (!tableName) {
|
||||
this.logger.warn('deleteJob: table name not resolved, skipping DynamoDB sync', { jobId });
|
||||
} else {
|
||||
this.logger.warn('deleteJob: no input ID in ES pipeline, skipping DynamoDB sync', { jobId });
|
||||
}
|
||||
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 });
|
||||
});
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user