mirror of
https://github.com/dadosfera/maestro.git
synced 2026-10-10 09:39:09 +00:00
FIX: pipeline accepting properties on payload
This commit is contained in:
Generated
+7
-7
@@ -12,7 +12,7 @@
|
||||
"dependencies": {
|
||||
"@aws-sdk/client-secrets-manager": "^3.112.0",
|
||||
"@dadosfera/dadosfera-logs": "^1.0.0-beta.3",
|
||||
"@dadosfera/protospack-v2": "3.14.11",
|
||||
"@dadosfera/protospack-v2": "3.14.13",
|
||||
"@grpc/grpc-js": "^1.6.7",
|
||||
"@grpc/proto-loader": "^0.6.13",
|
||||
"@nestjs/common": "^8.4.7",
|
||||
@@ -1726,9 +1726,9 @@
|
||||
}
|
||||
},
|
||||
"node_modules/@dadosfera/protospack-v2": {
|
||||
"version": "3.14.11",
|
||||
"resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com:443/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.14.11.tgz",
|
||||
"integrity": "sha512-ieLu6h179H3isgPkj7i9BBEA6PXKCeUDFBazUSf8laYnM04DTQCXO6I+vMVdyA90vbw6ccFqz5udWSKXG/EwKQ==",
|
||||
"version": "3.14.13",
|
||||
"resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com:443/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.14.13.tgz",
|
||||
"integrity": "sha512-OVX3LFhq54+bxMgP1tX2Pw5j3rDAPI0jWtd2/XKBcj9U0Pt6YX15q6MpWK9vCrNst6w3P5AVvLLfFOUQ5wzyOg==",
|
||||
"dependencies": {
|
||||
"@grpc/grpc-js": "^1.6.7",
|
||||
"rxjs": "^7.5.5",
|
||||
@@ -12305,9 +12305,9 @@
|
||||
}
|
||||
},
|
||||
"@dadosfera/protospack-v2": {
|
||||
"version": "3.14.11",
|
||||
"resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com:443/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.14.11.tgz",
|
||||
"integrity": "sha512-ieLu6h179H3isgPkj7i9BBEA6PXKCeUDFBazUSf8laYnM04DTQCXO6I+vMVdyA90vbw6ccFqz5udWSKXG/EwKQ==",
|
||||
"version": "3.14.13",
|
||||
"resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com:443/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.14.13.tgz",
|
||||
"integrity": "sha512-OVX3LFhq54+bxMgP1tX2Pw5j3rDAPI0jWtd2/XKBcj9U0Pt6YX15q6MpWK9vCrNst6w3P5AVvLLfFOUQ5wzyOg==",
|
||||
"requires": {
|
||||
"@grpc/grpc-js": "^1.6.7",
|
||||
"rxjs": "^7.5.5",
|
||||
|
||||
+1
-1
@@ -27,7 +27,7 @@
|
||||
"dependencies": {
|
||||
"@aws-sdk/client-secrets-manager": "^3.112.0",
|
||||
"@dadosfera/dadosfera-logs": "^1.0.0-beta.3",
|
||||
"@dadosfera/protospack-v2": "3.14.11",
|
||||
"@dadosfera/protospack-v2": "3.14.13",
|
||||
"@grpc/grpc-js": "^1.6.7",
|
||||
"@grpc/proto-loader": "^0.6.13",
|
||||
"@nestjs/common": "^8.4.7",
|
||||
|
||||
-36
@@ -1,36 +0,0 @@
|
||||
import { Info } from 'protospack/dist/lib/interfaces';
|
||||
|
||||
export interface ICreatePipelineDto {
|
||||
input: IdRequest;
|
||||
transformations: IdRequest[];
|
||||
output: IdRequest;
|
||||
tags: string[];
|
||||
name: string;
|
||||
description: string;
|
||||
info: Info;
|
||||
}
|
||||
|
||||
export interface IdRequest {
|
||||
id: string;
|
||||
}
|
||||
|
||||
export interface IIdRequest {
|
||||
id: string;
|
||||
info: Info;
|
||||
}
|
||||
|
||||
export interface IUpdatePipelineRequest {
|
||||
input: IdRequest;
|
||||
transformations: IdRequest[];
|
||||
output: IdRequest;
|
||||
tags: string[];
|
||||
name: string;
|
||||
description: string;
|
||||
id: string;
|
||||
info: Info;
|
||||
}
|
||||
|
||||
export interface IGetPipelineLogsRequest {
|
||||
id: string;
|
||||
details: string;
|
||||
}
|
||||
@@ -0,0 +1,68 @@
|
||||
import { ApiProperty, ApiPropertyOptional, OmitType } from '@nestjs/swagger';
|
||||
import { Info } from 'protospack/dist/lib/interfaces';
|
||||
|
||||
export class IPipelineV2 {
|
||||
@ApiProperty()
|
||||
id: string;
|
||||
|
||||
@ApiProperty()
|
||||
name: string;
|
||||
@ApiProperty()
|
||||
description: string;
|
||||
@ApiProperty()
|
||||
connection_id: string;
|
||||
@ApiProperty()
|
||||
cron: string;
|
||||
@ApiPropertyOptional()
|
||||
input_id?: string;
|
||||
@ApiPropertyOptional()
|
||||
transformations_ids?: string[];
|
||||
@ApiPropertyOptional()
|
||||
tags?: string[];
|
||||
@ApiPropertyOptional()
|
||||
properties?: any;
|
||||
|
||||
@ApiProperty()
|
||||
connector_name: string;
|
||||
@ApiProperty()
|
||||
connector_plugin: string;
|
||||
@ApiProperty()
|
||||
connector_version: string;
|
||||
@ApiProperty()
|
||||
image_url: string;
|
||||
|
||||
@ApiProperty()
|
||||
created_at: string;
|
||||
@ApiProperty()
|
||||
updated_at: string;
|
||||
}
|
||||
export class ICreatePipelineV2Req extends OmitType(IPipelineV2, [
|
||||
'id',
|
||||
'created_at',
|
||||
'updated_at',
|
||||
]) {}
|
||||
|
||||
export interface IdRequest {
|
||||
id: string;
|
||||
}
|
||||
|
||||
export interface IIdRequest {
|
||||
id: string;
|
||||
info: Info;
|
||||
}
|
||||
|
||||
export interface IUpdatePipelineRequest {
|
||||
input: IdRequest;
|
||||
transformations: IdRequest[];
|
||||
output: IdRequest;
|
||||
tags: string[];
|
||||
name: string;
|
||||
description: string;
|
||||
id: string;
|
||||
info: Info;
|
||||
}
|
||||
|
||||
export interface IGetPipelineLogsRequest {
|
||||
id: string;
|
||||
details: string;
|
||||
}
|
||||
@@ -9,7 +9,7 @@ import {
|
||||
Put,
|
||||
Headers,
|
||||
} from '@nestjs/common';
|
||||
import { ApiTags } from '@nestjs/swagger';
|
||||
import { ApiCreatedResponse, ApiTags } from '@nestjs/swagger';
|
||||
import { AuthenticateCondition } from 'src/authentication/authentication.decorator';
|
||||
import { PERMISSIONS } from '../../authentication/permissions.enum';
|
||||
import { PipelinesService } from './pipelines.service';
|
||||
@@ -19,6 +19,7 @@ import { RequestUser, User } from 'src/authentication/user.decorator';
|
||||
import { PackTheMetadata } from 'src/utils/ PackTheMetadata';
|
||||
|
||||
import { PipelinesService as OldPipelineService } from 'src/modules/pipelines/pipelines.service';
|
||||
import { ICreatePipelineV2Req, IPipelineV2 } from './interfaces';
|
||||
@ApiTags('Pipelines')
|
||||
@Controller('pipelinesV2')
|
||||
@AuthenticateCondition((req, user) => {
|
||||
@@ -53,10 +54,11 @@ export class PipelinesController {
|
||||
}
|
||||
|
||||
@Post()
|
||||
@ApiCreatedResponse({ type: IPipelineV2 })
|
||||
async create(
|
||||
@Headers('Dadosfera-Lang') language,
|
||||
@User() user: RequestUser,
|
||||
@Body() createPipelineDto: Messages.PipelineV2CreateRequest,
|
||||
@Body() createPipelineDto: ICreatePipelineV2Req,
|
||||
) {
|
||||
this.logger.info(process.env.DEV_URL + `/pipelines ON CREATE ROUTE`);
|
||||
const { customer_id, customer_name, username, user_id } = user;
|
||||
|
||||
@@ -4,7 +4,7 @@ import {
|
||||
InternalServerErrorException,
|
||||
OnModuleInit,
|
||||
} from '@nestjs/common';
|
||||
import { ClientGrpc, Payload } from '@nestjs/microservices';
|
||||
import { ClientGrpc } from '@nestjs/microservices';
|
||||
import {
|
||||
Messages,
|
||||
WriteService,
|
||||
@@ -15,6 +15,8 @@ import { lastValueFrom } from 'rxjs';
|
||||
|
||||
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
||||
import { PipelinesClientConfiguration } from './pipelines-client';
|
||||
import { ICreatePipelineV2Req } from './interfaces';
|
||||
import { PipelineV2CreateRequest } from '@dadosfera/protospack-v2/dist/lib/PipelineV2/interfaces/messages';
|
||||
|
||||
export class PipelinesService implements OnModuleInit {
|
||||
logger: DadosferaLogger;
|
||||
@@ -42,25 +44,33 @@ export class PipelinesService implements OnModuleInit {
|
||||
}
|
||||
|
||||
async create(
|
||||
@Payload() createPipelineDto: Messages.PipelineV2CreateRequest,
|
||||
body: ICreatePipelineV2Req,
|
||||
metadata,
|
||||
): Promise<Messages.PipelineV2CreateResponse> {
|
||||
this.logger.info('PipelinesClientService - Create');
|
||||
this.logger.info('PipelinesV2ClientService - Create');
|
||||
|
||||
const pipelineV2CreateRequest: PipelineV2CreateRequest = {
|
||||
...body,
|
||||
input_id: body.input_id,
|
||||
transformations_ids: body.transformations_ids,
|
||||
tags: body.tags,
|
||||
properties: body.properties && JSON.stringify(body.properties),
|
||||
};
|
||||
|
||||
const createPipelineResponse = await lastValueFrom(
|
||||
this.pipelineWriteService.PipelineV2Create(createPipelineDto, metadata),
|
||||
)
|
||||
.then((res) => {
|
||||
this.logger.info('Done');
|
||||
return res;
|
||||
})
|
||||
.catch((error: { details: string }) => {
|
||||
if (error.details.includes('INVALID_REQUEST')) {
|
||||
// eslint-disable-next-line @typescript-eslint/no-unused-vars
|
||||
const [errorType, message] = error.details.split('|');
|
||||
throw new BadRequestException(message);
|
||||
}
|
||||
throw new InternalServerErrorException(error.details);
|
||||
});
|
||||
this.pipelineWriteService.PipelineV2Create(
|
||||
pipelineV2CreateRequest,
|
||||
metadata,
|
||||
),
|
||||
).catch((error: { details: string }) => {
|
||||
if (error.details.includes('INVALID_REQUEST')) {
|
||||
// eslint-disable-next-line @typescript-eslint/no-unused-vars
|
||||
const [errorType, message] = error.details.split('|');
|
||||
throw new BadRequestException(message);
|
||||
}
|
||||
throw new InternalServerErrorException(error.details);
|
||||
});
|
||||
this.logger.info('Done');
|
||||
|
||||
return createPipelineResponse;
|
||||
}
|
||||
|
||||
@@ -1,6 +1,6 @@
|
||||
import { Metadata } from '@grpc/grpc-js';
|
||||
|
||||
interface Info {
|
||||
interface IMetadata {
|
||||
customer_id?: string;
|
||||
customer_name?: string;
|
||||
username?: string;
|
||||
@@ -9,14 +9,11 @@ interface Info {
|
||||
details?: string;
|
||||
}
|
||||
|
||||
export const PackTheMetadata = (info: Info): Metadata => {
|
||||
export function PackTheMetadata(info: IMetadata): Metadata {
|
||||
const meta = new Metadata();
|
||||
meta.add('customer_id', info.customer_id);
|
||||
meta.add('customer_name', info.customer_name);
|
||||
meta.add('user_id', info.user_id);
|
||||
meta.add('username', info.username);
|
||||
meta.add('details', info.details);
|
||||
meta.add('language', info.language);
|
||||
for (const key in info) {
|
||||
meta.add(key, info[key]);
|
||||
}
|
||||
|
||||
return meta;
|
||||
};
|
||||
}
|
||||
|
||||
+1
-1
File diff suppressed because one or more lines are too long
Reference in New Issue
Block a user