diff --git a/docsfera.json b/docsfera.json index ec823b4..c78fbec 100644 --- a/docsfera.json +++ b/docsfera.json @@ -4705,6 +4705,42 @@ ] } }, + "/platform/jobs/{jobId}": { + "delete": { + "operationId": "PlatformApiController_deleteJob", + "summary": "Delete a job and mark its table as deleted via in-factory", + "parameters": [ + { + "name": "jobId", + "required": true, + "in": "path", + "schema": { + "type": "string" + } + } + ], + "responses": { + "200": { + "description": "", + "content": { + "application/json": { + "schema": { + "type": "object" + } + } + } + } + }, + "tags": [ + "Platform API" + ], + "security": [ + { + "access-token": [] + } + ] + } + }, "/platform/jobs/jdbc/{jobId}": { "get": { "operationId": "PlatformApiController_getJdbcJob", @@ -7952,966 +7988,6 @@ } } }, - "/platform/pipeline": { - "post": { - "operationId": "PlatformApiController_createPipeline", - "summary": "Create a new pipeline", - "parameters": [], - "responses": { - "201": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/pipelines": { - "get": { - "operationId": "PlatformApiController_getPipelines", - "summary": "List all pipelines for customer", - "parameters": [], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/pipeline/{pipelineId}": { - "get": { - "operationId": "PlatformApiController_getPipeline", - "summary": "Get pipeline by ID", - "parameters": [ - { - "name": "pipelineId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - }, - "patch": { - "operationId": "PlatformApiController_updatePipeline", - "summary": "Update pipeline by ID", - "parameters": [ - { - "name": "pipelineId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - }, - "delete": { - "operationId": "PlatformApiController_deletePipeline", - "summary": "Delete pipeline by ID", - "parameters": [ - { - "name": "pipelineId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/pipeline/execute": { - "post": { - "operationId": "PlatformApiController_executePipeline", - "summary": "Execute a pipeline", - "parameters": [], - "responses": { - "201": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/pipeline/pause": { - "post": { - "operationId": "PlatformApiController_pausePipeline", - "summary": "Pause a pipeline", - "parameters": [], - "responses": { - "201": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/pipeline/unpause": { - "post": { - "operationId": "PlatformApiController_unpausePipeline", - "summary": "Unpause a pipeline", - "parameters": [], - "responses": { - "201": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/pipeline/{pipelineId}/memory": { - "put": { - "operationId": "PlatformApiController_updatePipelineMemory", - "summary": "Update pipeline memory configuration", - "parameters": [ - { - "name": "pipelineId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/pipeline/{pipelineId}/metadata": { - "put": { - "operationId": "PlatformApiController_updatePipelineMetadata", - "summary": "Update pipeline metadata", - "parameters": [ - { - "name": "pipelineId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/pipelines/metadata": { - "get": { - "operationId": "PlatformApiController_getPipelinesMetadata", - "summary": "Get all pipelines metadata", - "parameters": [], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/pipelines/catalog/tables": { - "get": { - "operationId": "PlatformApiController_getAvailableTables", - "summary": "Get all tables available", - "parameters": [], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/pipelines/catalog/schemas": { - "get": { - "operationId": "PlatformApiController_getAvailableSchemas", - "summary": "Get all schemas available", - "parameters": [], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/pipelines/catalog/tables/validate": { - "post": { - "operationId": "PlatformApiController_validateTableAndSchema", - "summary": "Validate tables and schemas", - "parameters": [], - "requestBody": { - "required": true, - "content": { - "application/json": { - "schema": { - "type": "array", - "items": { - "type": "string" - } - } - } - } - }, - "responses": { - "201": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/pipeline/{pipelineId}/pipeline_run": { - "get": { - "operationId": "PlatformApiController_getPipelineRuns", - "summary": "Get pipeline runs for a pipeline", - "parameters": [ - { - "name": "pipelineId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/pipeline/{pipelineId}/pipeline_run/{runId}": { - "get": { - "operationId": "PlatformApiController_getPipelineRun", - "summary": "Get specific pipeline run", - "parameters": [ - { - "name": "pipelineId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - }, - { - "name": "runId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/pipeline/pipeline_run/{runId}/logs": { - "get": { - "operationId": "PlatformApiController_getPipelineRunLogs", - "summary": "Get pipeline run logs", - "parameters": [ - { - "name": "runId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/jobs/{jobId}/input": { - "put": { - "operationId": "PlatformApiController_updateJobInput", - "summary": "Update job input columns", - "parameters": [ - { - "name": "jobId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - }, - "patch": { - "operationId": "PlatformApiController_patchJobInput", - "summary": "Partial update job input columns", - "parameters": [ - { - "name": "jobId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/jobs/{jobId}/memory": { - "put": { - "operationId": "PlatformApiController_updateJobMemory", - "summary": "Update job memory configuration", - "parameters": [ - { - "name": "jobId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/jobs/{jobId}/reset-state": { - "post": { - "operationId": "PlatformApiController_resetJobState", - "summary": "Reset job state", - "parameters": [ - { - "name": "jobId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "201": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/jobs/{jobId}": { - "delete": { - "operationId": "PlatformApiController_deleteJob", - "summary": "Delete a job and mark its table as deleted in DynamoDB", - "parameters": [ - { - "name": "jobId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/jobs/jdbc/{jobId}": { - "get": { - "operationId": "PlatformApiController_getJdbcJob", - "summary": "Get JDBC job details", - "parameters": [ - { - "name": "jobId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/jobs/jdbc/{jobId}/sync-mode": { - "post": { - "operationId": "PlatformApiController_updateJdbcSyncMode", - "summary": "Update JDBC job sync mode", - "parameters": [ - { - "name": "jobId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "201": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/jobs/jdbc/configs/allowed_datatypes": { - "get": { - "operationId": "PlatformApiController_getJdbcAllowedDatatypes", - "summary": "Get allowed datatypes for JDBC", - "parameters": [], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/jobs/singer/{jobId}": { - "get": { - "operationId": "PlatformApiController_getSingerJob", - "summary": "Get Singer job details", - "parameters": [ - { - "name": "jobId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/jobs/singer/{jobId}/sync-mode": { - "post": { - "operationId": "PlatformApiController_updateSingerSyncMode", - "summary": "Update Singer job sync mode", - "parameters": [ - { - "name": "jobId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "201": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/jobs/s3/{jobId}": { - "get": { - "operationId": "PlatformApiController_getS3Job", - "summary": "Get S3 job details", - "parameters": [ - { - "name": "jobId", - "required": true, - "in": "path", - "schema": { - "type": "string" - } - } - ], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, - "/platform/health": { - "get": { - "operationId": "PlatformApiController_healthCheck", - "summary": "Platform API health check", - "parameters": [], - "responses": { - "200": { - "description": "", - "content": { - "application/json": { - "schema": { - "type": "object" - } - } - } - } - }, - "tags": [ - "Platform API" - ], - "security": [ - { - "access-token": [] - } - ] - } - }, "/storage-explorer/tables/validate-name": { "post": { "operationId": "StorageExplorerController_validateTableName", @@ -11483,6 +10559,114 @@ "engine" ] }, + "ValidationTableDTO": { + "type": "object", + "properties": { + "tables": { + "type": "array", + "items": { + "type": "string" + } + } + }, + "required": [ + "tables" + ] + }, + "EnforceMfa": { + "type": "object", + "properties": { + "enabled": { + "type": "boolean" + } + }, + "required": [ + "enabled" + ] + }, + "CustomerLinkItem": { + "type": "object", + "properties": { + "href": { + "type": "string" + }, + "name": { + "type": "string" + }, + "description": { + "type": "string" + }, + "iconSrc": { + "type": "string" + } + }, + "required": [ + "href", + "name", + "description" + ] + }, + "CustomerSidebarSection": { + "type": "object", + "properties": { + "title": { + "type": "object" + }, + "items": { + "type": "array", + "items": { + "oneOf": [ + { + "$ref": "#/components/schemas/CustomerSidebarMenuItem" + }, + { + "$ref": "#/components/schemas/CustomerSidebarLinkItem" + } + ] + } + } + }, + "required": [ + "title", + "items" + ] + }, + "CustomerLinksConfig": { + "type": "object", + "properties": { + "home": { + "type": "array", + "items": { + "$ref": "#/components/schemas/CustomerLinkItem" + } + }, + "sidebar": { + "type": "array", + "items": { + "$ref": "#/components/schemas/CustomerSidebarSection" + } + } + } + }, + "CustomerLinksResponse": { + "type": "object", + "properties": { + "links": { + "$ref": "#/components/schemas/CustomerLinksConfig" + } + } + }, + "CustomerLinkRequest": { + "type": "object", + "properties": { + "links": { + "$ref": "#/components/schemas/CustomerLinksConfig" + } + }, + "required": [ + "links" + ] + }, "IRole": { "type": "object", "properties": { @@ -12072,100 +11256,6 @@ "updated_by" ] }, - "EnforceMfa": { - "type": "object", - "properties": { - "enabled": { - "type": "boolean" - } - }, - "required": [ - "enabled" - ] - }, - "CustomerLinkItem": { - "type": "object", - "properties": { - "href": { - "type": "string" - }, - "name": { - "type": "string" - }, - "description": { - "type": "string" - }, - "iconSrc": { - "type": "string" - } - }, - "required": [ - "href", - "name", - "description" - ] - }, - "CustomerSidebarSection": { - "type": "object", - "properties": { - "title": { - "type": "object" - }, - "items": { - "type": "array", - "items": { - "oneOf": [ - { - "$ref": "#/components/schemas/CustomerSidebarMenuItem" - }, - { - "$ref": "#/components/schemas/CustomerSidebarLinkItem" - } - ] - } - } - }, - "required": [ - "title", - "items" - ] - }, - "CustomerLinksConfig": { - "type": "object", - "properties": { - "home": { - "type": "array", - "items": { - "$ref": "#/components/schemas/CustomerLinkItem" - } - }, - "sidebar": { - "type": "array", - "items": { - "$ref": "#/components/schemas/CustomerSidebarSection" - } - } - } - }, - "CustomerLinksResponse": { - "type": "object", - "properties": { - "links": { - "$ref": "#/components/schemas/CustomerLinksConfig" - } - } - }, - "CustomerLinkRequest": { - "type": "object", - "properties": { - "links": { - "$ref": "#/components/schemas/CustomerLinksConfig" - } - }, - "required": [ - "links" - ] - }, "CreateShareMetadataDto": { "type": "object", "properties": { diff --git a/package-lock.json b/package-lock.json index de2fca5..098aea0 100644 --- a/package-lock.json +++ b/package-lock.json @@ -17,7 +17,7 @@ "@aws-sdk/signature-v4": "^3.370.0", "@dadosfera/dadosfera-logs": "^1.0.0-beta.4", "@dadosfera/protospack": "2.5.3", - "@dadosfera/protospack-v2": "^3.40.0-beta.3", + "@dadosfera/protospack-v2": "^3.40.0-beta.4", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", @@ -1745,9 +1745,9 @@ } }, "node_modules/@dadosfera/protospack-v2": { - "version": "3.40.0-beta.3", - "resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.40.0-beta.3.tgz", - "integrity": "sha512-m77RSqAO+hZkjUgbdrMqC+6yyyo7indH6fWP3WuygCHGYwdN6FsWTA+EGTAhXQRSgNlI8h358ZpRoILnnGzLMg==", + "version": "3.40.0-beta.4", + "resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.40.0-beta.4.tgz", + "integrity": "sha512-ftK4edb3fgGw3qLx+cT79cHYL+1gBoWXd8+SBFQXWHBrnxoRKbanBDp5Yn9PbaYZdRMSqTxSKsBwiTyJzLBunQ==", "dependencies": { "@grpc/grpc-js": "^1.9.3", "rxjs": "^7.5.5" diff --git a/package.json b/package.json index 586d74c..e4815e6 100644 --- a/package.json +++ b/package.json @@ -35,7 +35,7 @@ "@aws-sdk/signature-v4": "^3.370.0", "@dadosfera/dadosfera-logs": "^1.0.0-beta.4", "@dadosfera/protospack": "2.5.3", - "@dadosfera/protospack-v2": "^3.40.0-beta.3", + "@dadosfera/protospack-v2": "^3.40.0-beta.4", "@grpc/grpc-js": "^1.9.3", "@grpc/proto-loader": "^0.7.9", "@nestjs/cli": "^9.5.0", diff --git a/src/modules/inputs/inputs.service.ts b/src/modules/inputs/inputs.service.ts index 46ba2a4..5ead744 100644 --- a/src/modules/inputs/inputs.service.ts +++ b/src/modules/inputs/inputs.service.ts @@ -288,4 +288,8 @@ export class InputsService { }; return formatedPayload; } + + async markTableDeleted(data: { input_id: string; table_name: string; info: Info }) { + return lastValueFrom(this.inputWriteService.MarkTableDeleted(data)); + } } diff --git a/src/modules/pipelinesV2/pipelines.controller.ts b/src/modules/pipelinesV2/pipelines.controller.ts index ab39d38..5e6456e 100644 --- a/src/modules/pipelinesV2/pipelines.controller.ts +++ b/src/modules/pipelinesV2/pipelines.controller.ts @@ -48,6 +48,7 @@ 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'; @@ -64,6 +65,7 @@ export class PipelinesController { dadosferaLogger: DadosferaLogger, private pipelinesClientService: PipelinesService, private oldPipelinesService: OldPipelineService, + private platformApiService: PlatformApiService, ) { this.logger = dadosferaLogger.logger; } @@ -222,31 +224,40 @@ export class PipelinesController { language, }); - const result = await this.pipelinesClientService - .findOne({ id }, metadata) - .then(async (res) => { - //{pipeline:{tables: {tables: [], input_id: ''}}} - const parsed = JSON.parse(res.pipeline.config.tables); - const input_id = parsed?.input_id; - const tables = parsed?.tables ?? parsed; + const [pipelineRes, platformPipeline] = await Promise.all([ + this.pipelinesClientService.findOne({ id }, metadata), + this.platformApiService.proxy('GET', `/pipeline/${id.replace(/-/g, '_')}`, user).catch(() => null), + ]); - Object.assign(res.pipeline, { - transformations: res.pipeline.transformations - ? JSON.parse(res.pipeline.transformations) - : [], - config: { - cron: res.pipeline.config.cron, - tables, - input_id - }, - properties: res.pipeline.properties - ? JSON.parse(res.pipeline.properties) - : {}, - }); - return res; + //{pipeline:{tables: {tables: [], input_id: ''}}} + type PipelineTablesConfig = { input_id?: string; tables?: Array<{ name: string }> } | Array<{ name: string }>; + const parsed: PipelineTablesConfig = JSON.parse(pipelineRes.pipeline.config.tables); + const input_id = (parsed as { input_id?: string })?.input_id; + let tables: Array = (parsed as { tables?: Array })?.tables ?? (parsed as Array); + + 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; }); + } - return result; + Object.assign(pipelineRes.pipeline, { + transformations: pipelineRes.pipeline.transformations + ? JSON.parse(pipelineRes.pipeline.transformations) + : [], + config: { + cron: pipelineRes.pipeline.config.cron, + tables, + input_id, + }, + properties: pipelineRes.pipeline.properties + ? JSON.parse(pipelineRes.pipeline.properties) + : {}, + }); + + return pipelineRes; } @Patch('/:id') diff --git a/src/modules/platform-api/platform-api.controller.ts b/src/modules/platform-api/platform-api.controller.ts index a83b7bd..efe6592 100644 --- a/src/modules/platform-api/platform-api.controller.ts +++ b/src/modules/platform-api/platform-api.controller.ts @@ -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; } diff --git a/src/modules/platform-api/platform-api.module.ts b/src/modules/platform-api/platform-api.module.ts index b78e680..b2b0726 100644 --- a/src/modules/platform-api/platform-api.module.ts +++ b/src/modules/platform-api/platform-api.module.ts @@ -8,9 +8,10 @@ import { ElasticsearchModule } from '../../services/elasticsearch'; import { DynamoDBModule } from '../../services/dynamodb'; import { CustomersModule } from '../customers/customers.module'; import { CatalogModule } from '../catalog/catalog.module'; +import { InputsModule } from '../inputs/inputs.module'; @Module({ - imports: [ElasticsearchModule, DynamoDBModule, CustomersModule, CatalogModule], + imports: [ElasticsearchModule, DynamoDBModule, CustomersModule, CatalogModule, InputsModule], controllers: [PlatformApiController], providers: [PlatformApiService, DadosferaLogger], exports: [PlatformApiService],