mirror of
https://github.com/dadosfera/maestro.git
synced 2026-10-09 16:39:09 +00:00
312 lines
8.1 KiB
TypeScript
312 lines
8.1 KiB
TypeScript
import {
|
|
Body,
|
|
Controller,
|
|
Delete,
|
|
Get,
|
|
Inject,
|
|
Param,
|
|
Post,
|
|
Put,
|
|
Headers,
|
|
Query,
|
|
UseFilters,
|
|
HttpCode,
|
|
HttpStatus,
|
|
Patch,
|
|
UseInterceptors,
|
|
UploadedFile,
|
|
} from '@nestjs/common';
|
|
import {
|
|
ApiConsumes,
|
|
ApiCreatedResponse,
|
|
ApiNoContentResponse,
|
|
ApiTags,
|
|
} from '@nestjs/swagger';
|
|
import {
|
|
AuthenticateCondition,
|
|
RequireAllPermissions,
|
|
} from 'src/authentication/authentication.decorator';
|
|
import { PERMISSIONS_GROUPS } from '../../authentication/permissions.enum';
|
|
import { PipelinesService } from './pipelines.service';
|
|
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
|
|
import { Messages } from '@dadosfera/protospack-v2/dist/lib/PipelineV2';
|
|
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,
|
|
IUploadCSVFile,
|
|
} from './interfaces';
|
|
import { GrpcToHttpExceptionFilter } from 'src/error/grpc-to-http-exception.filter';
|
|
import { FileInterceptor } from '@nestjs/platform-express';
|
|
|
|
@UseFilters(new GrpcToHttpExceptionFilter())
|
|
@ApiTags('PipelinesV2')
|
|
@Controller('pipelinesV2')
|
|
@AuthenticateCondition((req, user) => {
|
|
let action;
|
|
|
|
switch (req.method) {
|
|
case 'POST':
|
|
action = 'CREATE';
|
|
break;
|
|
|
|
case 'PUT':
|
|
action = 'UPDATE';
|
|
break;
|
|
|
|
case 'PATCH':
|
|
action = 'UPDATE';
|
|
break;
|
|
|
|
default:
|
|
action = req.method;
|
|
}
|
|
|
|
return user.permissions.includes(
|
|
PERMISSIONS_GROUPS.PIPELINE.permissions[action].seqid,
|
|
);
|
|
})
|
|
export class PipelinesController {
|
|
logger: DadosferaLogger;
|
|
constructor(
|
|
@Inject(DadosferaLogger)
|
|
dadosferaLogger: DadosferaLogger,
|
|
private pipelinesClientService: PipelinesService,
|
|
private oldPipelinesService: OldPipelineService,
|
|
) {
|
|
this.logger = dadosferaLogger.logger;
|
|
}
|
|
|
|
@Post()
|
|
@ApiCreatedResponse({ type: IPipelineV2 })
|
|
async create(
|
|
@Headers('Dadosfera-Lang') language,
|
|
@User() user: RequestUser,
|
|
@Body() createPipelineDto: ICreatePipelineV2Req,
|
|
) {
|
|
this.logger.info('PipelinesController - create', { user });
|
|
if (!language) language = 'en-us';
|
|
const { customer_id, customer_name, username, user_id } = user;
|
|
const metadata = PackTheMetadata({
|
|
customer_id,
|
|
customer_name,
|
|
username,
|
|
user_id,
|
|
language,
|
|
});
|
|
|
|
const response = await this.pipelinesClientService.create(
|
|
createPipelineDto,
|
|
metadata,
|
|
);
|
|
|
|
return response;
|
|
}
|
|
|
|
@Get()
|
|
async findAll(
|
|
@User() user: RequestUser,
|
|
@Headers('Dadosfera-Lang') language,
|
|
@Query() data,
|
|
) {
|
|
this.logger.info('PipelinesController - findAll', { user });
|
|
|
|
if (!language) language = 'en-us';
|
|
|
|
const { customer_id, customer_name, user_id, username } = user;
|
|
const metadata = PackTheMetadata({
|
|
customer_id,
|
|
customer_name,
|
|
user_id,
|
|
username,
|
|
language,
|
|
});
|
|
|
|
const response = await this.pipelinesClientService.findAll(data, metadata);
|
|
|
|
return {
|
|
pipelines: [...response.pipelines],
|
|
};
|
|
}
|
|
|
|
@Get(':id/config')
|
|
async getPipelineproperties(
|
|
@Headers('Dadosfera-Lang') language,
|
|
@User() user: RequestUser,
|
|
@Param('id') id,
|
|
) {
|
|
this.logger.info('PipelinesController - getPipelineproperties', { user });
|
|
const metadata = PackTheMetadata({
|
|
...user,
|
|
language: language || 'pt-br',
|
|
});
|
|
return this.pipelinesClientService.findOneProperties(id, metadata);
|
|
}
|
|
|
|
@Get('/:id')
|
|
async findOne(
|
|
@Headers('Dadosfera-Lang') language,
|
|
@User() user: RequestUser,
|
|
@Param('id') id,
|
|
): Promise<Messages.PipelineV2FindOneResponse> {
|
|
this.logger.info('PipelinesController - findOne', { user });
|
|
|
|
if (!language) language = 'en-us';
|
|
|
|
const { customer_id, customer_name, user_id, username } = user;
|
|
|
|
const metadata = PackTheMetadata({
|
|
customer_id,
|
|
customer_name,
|
|
user_id,
|
|
username,
|
|
language,
|
|
});
|
|
|
|
const result = await this.pipelinesClientService
|
|
.findOne({ id }, metadata)
|
|
.then((res) => {
|
|
//{pipeline:{tables: {tables: [], input_id: ''}}}
|
|
let tables = JSON.parse(res.pipeline.config.tables);
|
|
if (tables?.tables) tables = tables.tables;
|
|
Object.assign(res.pipeline, {
|
|
transformations: res.pipeline.transformations
|
|
? JSON.parse(res.pipeline.transformations)
|
|
: [],
|
|
config: {
|
|
cron: res.pipeline.config.cron,
|
|
tables,
|
|
},
|
|
properties: res.pipeline.properties
|
|
? JSON.parse(res.pipeline.properties)
|
|
: {},
|
|
});
|
|
return res;
|
|
});
|
|
return result;
|
|
}
|
|
|
|
@Put('/:id')
|
|
async update(
|
|
@Headers('Dadosfera-Lang') language,
|
|
@Body() updatePipelineDto,
|
|
@Param('id') id,
|
|
@User() user: RequestUser,
|
|
) {
|
|
this.logger.info('PipelinesController - update', { user });
|
|
const { info } = updatePipelineDto;
|
|
delete updatePipelineDto.info;
|
|
|
|
const { customer_id, customer_name, user_id, username } = user;
|
|
if (!language) language = 'en-us';
|
|
|
|
const metadata = PackTheMetadata({
|
|
customer_id,
|
|
customer_name,
|
|
user_id,
|
|
username,
|
|
language,
|
|
});
|
|
|
|
const response = await this.pipelinesClientService.update(
|
|
{
|
|
...updatePipelineDto,
|
|
info,
|
|
id,
|
|
},
|
|
metadata,
|
|
);
|
|
|
|
this.logger.info('PipelinesController - update: OK', { user });
|
|
return response;
|
|
}
|
|
|
|
@Patch('/:id')
|
|
async updateByPatch(
|
|
@Headers('Dadosfera-Lang') language,
|
|
@Body() updatePipelineDto,
|
|
@Param('id') id,
|
|
@User() user: RequestUser,
|
|
) {
|
|
this.logger.info('PipelinesController - patch', { user });
|
|
const response = await this.update(language, updatePipelineDto, id, user);
|
|
this.logger.info('PipelinesController - patch: OK', { user });
|
|
return response;
|
|
}
|
|
|
|
@Delete(':id')
|
|
@ApiNoContentResponse()
|
|
@HttpCode(HttpStatus.NO_CONTENT)
|
|
async delete(@Param() params, @User() user: RequestUser) {
|
|
this.logger.info('PipelinesController - delete', { user });
|
|
const { id } = params;
|
|
const metadata = PackTheMetadata({
|
|
customer_id: user.customer_id,
|
|
customer_name: user.customer_name,
|
|
user_id: user.user_id,
|
|
});
|
|
await this.pipelinesClientService.remove({ id, metadata, user });
|
|
this.logger.info('PipelinesController - delete: OK');
|
|
}
|
|
|
|
@Post('/upload')
|
|
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.CREATE)
|
|
@ApiConsumes('multipart/form-data')
|
|
@UseInterceptors(
|
|
FileInterceptor('file', { limits: { fileSize: 50 * 1000 * 1000 + 1 } }),
|
|
)
|
|
async uploadFile(
|
|
@User() user: RequestUser,
|
|
@UploadedFile() file,
|
|
@Body() body: IUploadCSVFile,
|
|
) {
|
|
this.logger.info('/upload - Upload Connector Route');
|
|
const metadata = PackTheMetadata({ ...user });
|
|
let name = `${new Date().getTime()}_${file.originalname.split('.')[0]}`;
|
|
body.name ? (name = `${new Date().getTime()}_${body.name}`) : name;
|
|
const response = await this.pipelinesClientService.uploadFile(
|
|
{
|
|
file,
|
|
name,
|
|
},
|
|
metadata,
|
|
);
|
|
|
|
const { sep, header, encoding, description } = body;
|
|
|
|
const [file_name, file_format] = file.originalname.split('.');
|
|
|
|
const upload_pipeline = {
|
|
connection_id: process.env.UPLOAD_FILE_AGENT_CONNECTION,
|
|
connector_name: 'Amazon S3',
|
|
connector_plugin: 'aws_s3',
|
|
connector_version: '1.0.0',
|
|
image_url: 'https://assets.dadosfera.ai/images/connectors/csv.svg',
|
|
name: body.name || file_name,
|
|
description,
|
|
transformations_ids: [],
|
|
tags: [],
|
|
cron: '@once',
|
|
config: { cron: '@once', tables: [] },
|
|
properties: {
|
|
engine: 'csv',
|
|
source_bucket: process.env.BUCKET_CUSTOMER_CSV_ASSETS,
|
|
source_prefix: `${user.customer_name}/${name}.${file_format}`,
|
|
file_format_params: { sep, encoding, header: Boolean(header) },
|
|
is_a_upload_csv: true,
|
|
},
|
|
input_id: undefined,
|
|
};
|
|
|
|
const pipeline = await this.pipelinesClientService.create(
|
|
upload_pipeline,
|
|
metadata,
|
|
);
|
|
|
|
return pipeline;
|
|
}
|
|
}
|