Merge Input Refactor

This commit is contained in:
Victor Radael
2022-08-12 10:43:20 -03:00
26 changed files with 397 additions and 591 deletions
+1 -11
View File
@@ -1,6 +1,5 @@
import { ApiProperty, ApiPropertyOptional } from '@nestjs/swagger';
import { Info } from 'src/clients/inputs/interfaces';
import { Method } from 'axios';
import { Info } from 'protospack/dist/lib/interfaces';
export class InputModel {
category: string;
@@ -13,15 +12,6 @@ export class InputModel {
oauth_code: string;
start_date: string;
};
options?: {
oauth?: {
get_tokens_url: string;
get_tokens_method?: Method;
get_tokens_url_params: string;
get_tokens_set_response: Record<string, string>;
content_type: string;
};
};
}
class GoogleAnalyticsClientSecrets {
+59
View File
@@ -0,0 +1,59 @@
import { Info } from 'protospack/dist/lib/interfaces';
interface Values {
jdbc_user: string;
jdbc_password: string;
database: string;
endpoint: string;
tables: string[];
port: string;
engine: string;
schema: string;
}
interface Cron {
hours: string[];
hour_interval: boolean;
hour_resourse: boolean;
hour_resourse_value: number;
month_day: string[];
month_day_interval: boolean;
month_day_resourse: boolean;
month_day_resourse_value: number;
month: string[];
month_interval: boolean;
month_resourse: boolean;
month_resourse_value: number;
week_days: string[];
week_days_interval: boolean;
week_days_resourse: boolean;
week_days_resourse_value: number;
}
export interface ICreateInputRequest {
cron: string;
name: string;
plugin: string;
values: Values;
operation: string;
info: Info;
}
export interface IIdRequest {
id: string;
info: Info;
}
export interface UpdateInputRequest {
id: string;
cron: string;
name: string;
plugin: string;
values: Values;
operation: string;
info: Info;
}
export interface ITestConnectionRequest {
plugin: string;
values: Values;
}
@@ -0,0 +1,28 @@
import { credentials } from '@grpc/grpc-js';
import {
ClientProviderOptions,
GrpcOptions,
Transport,
} from '@nestjs/microservices';
import { InputPackages, InputProtoFilePath } from 'protospack';
export class InputsGrpcClient {
private config: GrpcOptions = {
transport: Transport.GRPC,
options: {
url: process.env.INFACTORY_URL,
package: InputPackages,
credentials: process.env.LOCAL_ENV ? undefined : credentials.createSsl(),
protoPath: InputProtoFilePath,
loader: {
enums: String,
objects: true,
arrays: true,
},
},
};
providerOptions: ClientProviderOptions = {
name: 'InputsGrpcClient',
...this.config,
};
}
+2 -6
View File
@@ -10,8 +10,6 @@ import {
} from '@nestjs/common';
import { InputsService } from './inputs.service';
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
import { InputsClientService } from 'src/clients/inputs/client.service';
import { UpdateInputRequest } from 'src/clients/inputs/interfaces';
import { PERMISSIONS } from '../../authentication/permissions.enum';
import { AuthenticateCondition } from 'src/authentication/authentication.decorator';
import { ApiOkResponse, ApiTags } from '@nestjs/swagger';
@@ -25,7 +23,7 @@ import {
TestConnectionReq,
TestConnectionRes,
} from './dtos/input.model';
import { User } from 'src/authentication/user.decorator';
import { UpdateInputRequest } from './dtos/old_interfaces';
@ApiTags('inputs')
@Controller('inputs')
@@ -59,15 +57,13 @@ import { User } from 'src/authentication/user.decorator';
return user.permissions.includes(PERMISSIONS.INPUT.permissions[action].seqid);
})
export class InputsController {
inputService: InputsService;
logger: DadosferaLogger;
constructor(
@Inject(DadosferaLogger)
dadosferaLogger: DadosferaLogger,
private inputsClientService: InputsClientService,
private inputService: InputsService,
) {
this.logger = dadosferaLogger.logger;
this.inputService = new InputsService(this.inputsClientService);
}
@Get('available-entities/:plugin')
+16
View File
@@ -0,0 +1,16 @@
import { Module } from '@nestjs/common';
import { ClientsModule } from '@nestjs/microservices';
import { InputsGrpcClient } from './inputs-client.config';
import { InputsController } from './inputs.controller';
import { InputsService } from './inputs.service';
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
const client = new InputsGrpcClient();
@Module({
controllers: [InputsController],
providers: [InputsService, DadosferaLogger],
imports: [ClientsModule.register([client.providerOptions])],
exports: [InputsService],
})
export class InputsModule {}
+255 -22
View File
@@ -1,13 +1,247 @@
import { Body, HttpException, HttpStatus, Injectable } from '@nestjs/common';
import {
Body,
ForbiddenException,
HttpException,
HttpStatus,
Inject,
Injectable,
InternalServerErrorException,
NotFoundException,
} from '@nestjs/common';
import { Timeout } from '@nestjs/schedule';
import { DecodeGrpcStruct } from 'protospack';
import {
DecodeGrpcStruct,
EncodeJsonToGrpcStruct,
InputCreateS3Request,
InputNewCreateRequest,
InputService,
InputServicesNames,
TestConnectionGetColumnsRequest,
TestConnectionRequest,
} from 'protospack';
import CronParser, { CronExpression } from 'cron-parser';
import { InputsClientService } from 'src/clients/inputs/client.service';
import { IIdRequest, Info } from 'src/clients/inputs/interfaces';
import { Info } from 'protospack/dist/lib/interfaces';
import { ClientGrpc } from '@nestjs/microservices';
import DadosferaLogger from '@dadosfera/dadosfera-logs/dist';
import { lastValueFrom } from 'rxjs';
import {
objectSnakeToCamel,
objectCamelToSnake,
} from 'src/utils/CaseConverter';
import { IIdRequest, UpdateInputRequest } from './dtos/old_interfaces';
@Injectable()
//TODO refactor this class to remove OLD_inputClient
export class InputsService {
constructor(private inputClient: InputsClientService) {}
logger: DadosferaLogger;
inputService: InputService;
OLD_inputClient = {
newCreate: async (createInputDto: InputNewCreateRequest) => {
this.logger.info('InputClientService - Create');
const input = await lastValueFrom(
this.inputService.NewCreate(objectSnakeToCamel(createInputDto)),
);
objectCamelToSnake(input);
return input;
},
createGeneric: async (createInputGeneric) => {
this.logger.info('InputClientService - Create');
const { info, ...data } = createInputGeneric;
if (data.options) delete data.options;
if (data.credentials.oauth_code) delete data.credentials.oauth_code;
const grpcPayload = {
input: EncodeJsonToGrpcStruct(data),
info,
};
objectSnakeToCamel(grpcPayload);
const structReturn = await lastValueFrom(
this.inputService.Create(grpcPayload),
).catch((err: { details: string }) => {
if (err.details === 'Item Not found!') {
throw new NotFoundException('Input not found');
}
throw new InternalServerErrorException(err.details);
});
objectCamelToSnake(structReturn);
const inputCreated = DecodeGrpcStruct(structReturn.input);
return { input: inputCreated };
},
createS3Inputs: async (createInputDto: InputCreateS3Request) => {
this.logger.info('InputClientService - Create');
objectSnakeToCamel(createInputDto);
const createInputResponse = await lastValueFrom(
this.inputService.CreateS3(createInputDto),
).catch((e) => {
throw new HttpException(e.details, 500);
});
this.logger.info('Done');
objectCamelToSnake(createInputResponse);
return createInputResponse;
},
findOne: async (data: IIdRequest) => {
this.logger.info('InputClientService - FindOne');
const findOneInputResponse = await lastValueFrom(
this.inputService.FindOne(objectSnakeToCamel(data)),
).catch((e) => {
this.logger.error(e.details);
throw new HttpException(e.details, HttpStatus.INTERNAL_SERVER_ERROR);
});
objectCamelToSnake(findOneInputResponse);
return findOneInputResponse;
},
findAll: async (data) => {
this.logger.info('InputClientService - FindAll');
const findAllInputResponse = await lastValueFrom(
this.inputService.FindAll(objectSnakeToCamel(data)),
).catch((e) => {
this.logger.error(e.details);
throw new HttpException(e.details, HttpStatus.INTERNAL_SERVER_ERROR);
});
findAllInputResponse.inputs.forEach((input) => objectCamelToSnake(input));
return findAllInputResponse;
},
update: async (updateInputDTO: UpdateInputRequest) => {
this.logger.info('InputClientService - Update');
const updateInputResponse = await new Promise((resolve, reject) => {
this.inputService.Update(objectSnakeToCamel(updateInputDTO)).subscribe({
next(x) {
resolve(objectCamelToSnake(x));
},
error(err) {
reject(err);
},
complete() {
// console.log('done');
},
});
})
.then((res) => {
this.logger.info('Done');
return res;
})
.catch((err) => {
this.logger.error(err.message);
throw new Error(err);
});
return updateInputResponse;
},
remove: async (idRequest: IIdRequest) => {
this.logger.info('InputClientService - Remove');
const removeInputResponse = await new Promise((resolve, reject) => {
this.inputService.Remove(objectSnakeToCamel(idRequest)).subscribe({
next(x) {
resolve(objectCamelToSnake(x));
},
error(err) {
reject(err);
},
complete() {
// console.log('done');
},
});
})
.then((res) => {
this.logger.info('Done');
return res;
})
.catch((err) => {
this.logger.error(err.message);
throw new Error(err);
});
return removeInputResponse;
},
testConnection: async (data: TestConnectionRequest) => {
this.logger.info('InputClientService - TestConnection');
const testConnectionResponse = await new Promise((resolve, reject) => {
this.inputService
.NewTestConnection(objectSnakeToCamel(data))
.subscribe({
next(x) {
resolve(objectCamelToSnake(x));
},
error(err) {
reject(err);
},
complete() {
// console.log('done');
},
});
})
.then((res) => {
this.logger.info('Done');
return res;
})
.catch((err) => {
this.logger.error(err.message);
throw new Error(err);
});
return testConnectionResponse;
},
getColumns: async (data: TestConnectionGetColumnsRequest) => {
this.logger.info('InputClientService - TestConnection/Get-Columns');
const getColumnsResponse = await new Promise((resolve, reject) => {
this.inputService.GetColumns(objectSnakeToCamel(data)).subscribe({
next(x) {
resolve(objectCamelToSnake(x));
},
error(err) {
reject(err);
},
complete() {
// console.log('done');
},
});
})
.then((res) => {
this.logger.info('Done');
return res;
})
.catch((err) => {
this.logger.error(err.message);
throw new Error(err);
});
return getColumnsResponse;
},
getAvailableEntities: async (data) => {
this.logger.info('InputClientService - GetAvailableEntities');
objectSnakeToCamel(data);
const response = await lastValueFrom(
this.inputService.GetAvailableEntities(data),
).catch((err) => {
this.logger.info(err.details);
if (err.details && err.details.includes('400'))
throw new HttpException('Plugin inválido', HttpStatus.BAD_REQUEST);
throw new HttpException(err.details, HttpStatus.INTERNAL_SERVER_ERROR);
});
objectCamelToSnake(response);
return response;
},
};
constructor(
@Inject(DadosferaLogger)
dadosferaLogger: DadosferaLogger,
@Inject('InputsGrpcClient') private readonly grpcClient: ClientGrpc,
) {
this.logger = dadosferaLogger.logger;
}
async onModuleInit() {
this.inputService = this.grpcClient.getService<InputService>(
InputServicesNames.InputService,
);
}
secondsInADay = 60 * 60 * 24;
secondsInAnHour = 60 * 60;
@@ -37,14 +271,12 @@ export class InputsService {
const afterNextDate = interval.next().toDate();
const secondsApart = this.getDifferenceInSeconds(nextDate, afterNextDate);
if (customer_tier === 'BASIC' && secondsApart < this.secondsInADay) {
throw new HttpException(
throw new ForbiddenException(
'Intervalo de tempo não pode ser inferior a um dia.',
HttpStatus.FORBIDDEN,
);
} else if (secondsApart < this.secondsInAnHour) {
throw new HttpException(
throw new ForbiddenException(
'Intervalo de tempo não pode ser inferior a uma hora.',
HttpStatus.FORBIDDEN,
);
}
}
@@ -57,7 +289,7 @@ export class InputsService {
case 'parquet': {
const { info, ...input } = data;
const inputPayload = this.generateInputS3Payload(input);
response = await this.inputClient.createS3Inputs({
response = await this.OLD_inputClient.createS3Inputs({
input: inputPayload,
info,
});
@@ -67,10 +299,10 @@ export class InputsService {
case 'mysql':
case 'postgresql':
case 'sqlserver':
response = await this.inputClient.newCreate(data);
response = await this.OLD_inputClient.newCreate(data);
break;
default:
response = await this.inputClient.createGeneric(data);
response = await this.OLD_inputClient.createGeneric(data);
break;
}
const adjustedInput = this.adjustInputPayload(response.input);
@@ -79,16 +311,18 @@ export class InputsService {
async reCreate(data) {
const { info, cron } = data;
if (cron) this.validateCron({ info, cron });
const response = await this.inputClient.createGeneric(data);
const response = await this.OLD_inputClient.createGeneric(data);
const adjustedInput = this.adjustInputPayload(response.input);
return { ...response, input: adjustedInput };
}
async getAvailableEntities(data): Promise<{ entities: string[] }> {
return await this.inputClient.getAvailableEntities(data);
return await this.OLD_inputClient.getAvailableEntities(data);
}
async findAll(body) {
try {
const findAllInputResponse: any = await this.inputClient.findAll(body);
const findAllInputResponse: any = await this.OLD_inputClient.findAll(
body,
);
if (findAllInputResponse?.inputs?.length) {
findAllInputResponse.inputs = findAllInputResponse.inputs.map((input) =>
this.adjustInputPayload(input),
@@ -102,7 +336,7 @@ export class InputsService {
async findOne(idRequest: IIdRequest) {
try {
const findOneInputResponse: any = await this.inputClient.findOne(
const findOneInputResponse: any = await this.OLD_inputClient.findOne(
idRequest,
);
findOneInputResponse.input = this.adjustInputPayload(
@@ -117,7 +351,7 @@ export class InputsService {
async update(id: string, data, info: Info) {
this.validateCron({ ...data, info });
try {
const updateInputResponse: any = await this.inputClient.update({
const updateInputResponse: any = await this.OLD_inputClient.update({
id,
info,
...data,
@@ -134,7 +368,7 @@ export class InputsService {
async remove(idRequest: IIdRequest) {
try {
const removeInputResponse = await this.inputClient.remove(idRequest);
const removeInputResponse = await this.OLD_inputClient.remove(idRequest);
return removeInputResponse;
} catch (err) {
@@ -145,9 +379,8 @@ export class InputsService {
@Timeout(60000 * 10) // Timeout set for 10 minutes
async testConnection(data) {
try {
const testConnectionInputResponse = await this.inputClient.testConnection(
data,
);
const testConnectionInputResponse =
await this.OLD_inputClient.testConnection(data);
return testConnectionInputResponse;
} catch (err) {
@@ -158,7 +391,7 @@ export class InputsService {
async getColumns(data) {
try {
const testConnectionGetColumnsResponse =
await this.inputClient.getColumns(data);
await this.OLD_inputClient.getColumns(data);
return testConnectionGetColumnsResponse;
} catch (err) {