Compare commits

...
12 changed files with 1551 additions and 1273 deletions
+1 -1
View File
@@ -9,7 +9,7 @@ maestro:
cookie_secret: "ff7bc13823edb2ae50d248e5780bddc9d4b31c36"
redis_database: "1"
platform_api_url: https://xs2hkhq07k.execute-api.us-east-1.amazonaws.com
storage_explorer_api_url: "http://storage-explorer-{customer}.data-apps.svc.cluster.local/api"
storage_explorer_api_url: "http://storage-explorer-{customer}.data-apps.svc.cluster.local:8000/api"
hostname: maestro.stg.dadosfera.ai
+1319 -1264
View File
File diff suppressed because it is too large Load Diff
+4 -4
View File
@@ -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.38.0-beta.27",
"@dadosfera/protospack-v2": "3.38.0-beta.28",
"@grpc/grpc-js": "^1.9.3",
"@grpc/proto-loader": "^0.7.9",
"@nestjs/cli": "^9.5.0",
@@ -1744,9 +1744,9 @@
}
},
"node_modules/@dadosfera/protospack-v2": {
"version": "3.38.0-beta.27",
"resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.38.0-beta.27.tgz",
"integrity": "sha512-lvV3g/SyPLFVfdtxoc/l1fCtl0o/X9Dv/KOdmDAXQyRJRjd3/81V83rwomp3UtI5IYGYuSbREgolH+X3aNYF7w==",
"version": "3.38.0-beta.28",
"resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.38.0-beta.28.tgz",
"integrity": "sha512-w3Au0qschqZJ6OSHVDKaR2KeeCF2ncuJ4k7iZQceOLo1zxUU3TpDyLYyp/rp1DoSG763uIgqm3kJUXv7RlmfNw==",
"license": "ISC",
"dependencies": {
"@grpc/grpc-js": "^1.9.3",
+1 -1
View File
@@ -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.38.0-beta.27",
"@dadosfera/protospack-v2": "3.38.0-beta.28",
"@grpc/grpc-js": "^1.9.3",
"@grpc/proto-loader": "^0.7.9",
"@nestjs/cli": "^9.5.0",
+1
View File
@@ -111,3 +111,4 @@ function configureSwagger(app: INestApplication) {
);
}
bootstrap();
+1
View File
@@ -478,6 +478,7 @@ export class AuthController {
@Get('me')
async getMe(@Req() req: Request, @Res() res: Response) {
this.logger.info('GET /auth/me ')
this.logger.info(req.headers);
// Check for API key header first
const apiKey = req.get('X-Api-key');
+1 -1
View File
@@ -200,7 +200,7 @@ export class InputsService {
}
async update(id: string, data, info: Info) {
this.validateCron({ ...data, info });
// this.validateCron({ ...data, info });
try {
const updateInputResponse: any = await this.OLD_inputClient.update({
id,
+36
View File
@@ -1,6 +1,15 @@
import { ApiProperty, ApiPropertyOptional, OmitType } from '@nestjs/swagger';
import { Info } from '@dadosfera/protospack/dist/lib/interfaces';
export class PipelineInputsDTO {
@ApiProperty()
tables: Array<{
name: string,
type: string,
}>
}
export class IPipelineV2 {
@ApiProperty()
id: string;
@@ -123,3 +132,30 @@ export class PipelineFindAllReq {
@ApiPropertyOptional()
type?: string | undefined;
}
export interface UpdateTableDTO {
name: string;
type: string;
columns: string[];
destinations: {
raw: {
table_schema: string;
table_name: string;
};
qualify: {
table_schema: string;
table_name: string;
};
};
identifier_columns: string[];
reference_column: {
name: string;
type: string;
};
memory: number;
}
export interface UpdatePlatformInputRequest {
cron: string;
tables: Array<UpdateTableDTO>;
}
@@ -43,11 +43,15 @@ import {
IPipelineV2,
IInitUploadCSVFile,
PipelineFindAllReq,
UpdatePlatformInputRequest,
} from './interfaces';
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 { ApiInternalOnlyEndpoint } from 'src/decorators/swagger.decorator';
import { TableColumns } from '../inputs/dtos/input.model';
import { UpdateInputRequest } from '../inputs/dtos/old_interfaces';
import { Info } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/entities';
@ApiTags('PipelinesV2')
@ApiHeaders([{ name: 'dadosfera-lang', enum: LanguageEnum, required: false }])
@@ -223,6 +227,7 @@ export class PipelinesController {
.then((res) => {
//{pipeline:{tables: {tables: [], input_id: ''}}}
let tables = JSON.parse(res.pipeline.config.tables);
const input_id = tables?.input_id;
if (tables?.tables) tables = tables.tables;
Object.assign(res.pipeline, {
transformations: res.pipeline.transformations
@@ -231,6 +236,7 @@ export class PipelinesController {
config: {
cron: res.pipeline.config.cron,
tables,
input_id
},
properties: res.pipeline.properties
? JSON.parse(res.pipeline.properties)
@@ -277,6 +283,45 @@ export class PipelinesController {
return response;
}
@Patch('/:pipelineId/inputs/:id')
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async updatePipelineInput(
@Language() language: LanguageEnum,
@Body() pipelineInputDTO: UpdatePlatformInputRequest,
@Param('id') inputId: string,
@Param('pipelineId') pipelineId: string,
@User() user: RequestUser,
) {
this.logger.info('PipelinesController - update', { user });
const { customer_id, customer_name, user_id, username } = user;
const info: Info = {
user_id: user.user_id,
customer: user.customer_name,
customer_id: user.customer_id,
};
const metadata = PackTheMetadata({
customer_id,
customer_name,
user_id,
username,
language,
});
const response = await this.pipelinesClientService.updatePipelineInput(
pipelineId,
inputId,
pipelineInputDTO,
info,
user,
metadata,
);
this.logger.info('PipelinesController - update: OK', { user });
return response;
}
@ApiInternalOnlyEndpoint()
@Put('/:id')
@ApiOperation({
@@ -11,6 +11,7 @@ import { PipelinesModule as OldPipelineModule } from 'src/modules/pipelines/pipe
import { ConnectorModule } from '../connector/connector.module';
import { InputsModule } from '../inputs/inputs.module';
import { TransformationsModule } from '../transformations/transformations.module';
import { PlatformApiModule } from '../platform-api/platform-api.module';
const client = new PipelinesClientConfiguration();
@@ -21,6 +22,7 @@ const client = new PipelinesClientConfiguration();
ConnectorModule,
InputsModule,
TransformationsModule,
PlatformApiModule
],
controllers: [PipelinesController],
providers: [PipelinesService, DadosferaLogger],
+139 -1
View File
@@ -1,3 +1,4 @@
/* eslint-disable no-async-promise-executor */
import {
BadRequestException,
HttpException,
@@ -16,7 +17,7 @@ import { lastValueFrom } from 'rxjs';
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
import { PipelinesClientConfiguration } from './pipelines-client';
import { ICreatePipelineV2Req } from './interfaces';
import { ICreatePipelineV2Req, UpdatePlatformInputRequest, UpdateTableDTO } from './interfaces';
import { PipelineV2CreateRequest } from '@dadosfera/protospack-v2/dist/lib/PipelineV2/interfaces/messages';
import { Metadata } from '@grpc/grpc-js';
import { ConnectorClientService } from '../connector/client.service';
@@ -26,6 +27,8 @@ import { TransformationsService } from '../transformations/transformations.servi
import { getObjValueFromPath, objHasPath } from 'src/utils/ObjValueFromPath';
import ErrorCodes from 'src/utils/errorCodes';
import ErrorBuilder from 'src/utils/ErrorBuilder';
import { PlatformApiService } from '../platform-api/platform-api.service';
import { Info } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/entities';
export class PipelinesService implements OnModuleInit {
logger: DadosferaLogger;
@@ -39,6 +42,7 @@ export class PipelinesService implements OnModuleInit {
private readonly connectorService: ConnectorClientService,
private readonly inputsService: InputsService,
private readonly transformationsService: TransformationsService,
private readonly platformAPI: PlatformApiService
) {
this.logger = dadosferaLogger.logger;
}
@@ -138,6 +142,7 @@ export class PipelinesService implements OnModuleInit {
const findOnePipelineResponse = await lastValueFrom(
this.pipelineReadService.PipelineV2FindOne(data, metadata),
);
console.log('pipeline find one response', findOnePipelineResponse);
this.logger.info('Done');
return findOnePipelineResponse;
@@ -339,4 +344,137 @@ export class PipelinesService implements OnModuleInit {
return res;
}
async updatePipelineInput(pipelineId: string, inputId: string, updateInputDTO: UpdatePlatformInputRequest, info: Info, user: RequestUser, metadata: Metadata) {
this.logger.info('InputClientService - Update');
this.logger.info('Update Dynamo Reference');
const pipelineIdFormat = pipelineId.split('-').join('_');
const updateInputResponse = await this.inputsService.update(
inputId,
updateInputDTO,
info
)
const requests = [];
this.logger.info('Dynamo Response', updateInputResponse);
for (const [index, table] of updateInputDTO.tables.entries()) {
const id = `${pipelineIdFormat}_${index}`;
this.logger.info('Updating input reference for table', table.name);
const body = {}
if (table.columns) {
body['column_include_list'] = table.columns;
}
if (table.reference_column) {
body['incremental_column_name'] = table.reference_column.name;
body['incremental_column_type'] = table.reference_column.type;
}
if (table.identifier_columns) {
body['primary_keys'] = table.identifier_columns;
}
this.logger.info('Request body', body);
const updateCollumns = this.platformAPI.proxy(
'PATCH',
`/jobs/${id}/input`,
user,
body
)
requests.push(updateCollumns);
if (table.memory) {
this.logger.info('Updating memory allocation for table', table.name);
const updateMemory = this.platformAPI.proxy(
'PUT',
`/jobs/${id}/memory`,
user,
{
amount: table.memory
}
)
requests.push(updateMemory);
}
if (table.type) {
const updateSyncMode = this.updatePipelineSyncMode(table, id, user);
requests.push(updateSyncMode);
}
}
this.logger.info('Create Platform Request for each JOB');
if (updateInputDTO.cron) {
const crnUpdatedRequest = new Promise(async (resolve, reject) => {
const response = await this.updatePipelineCron(updateInputDTO.cron, pipelineIdFormat, user);
if (response.error) {
this.logger.error('Error updating pipeline cron', response.error);
return reject(new ErrorBuilder(response.error));
}
this.logger.error('Pipeline cron updated successfully', response);
return resolve(response);
});
requests.push(crnUpdatedRequest);
}
this.logger.info('Executing all request for the platform api');
const results = await Promise.allSettled(requests);
this.logger.info('Platform api response', results);
return updateInputResponse;
}
private async updatePipelineSyncMode(table: UpdateTableDTO, pipelineId: string, user: RequestUser) {
const body = {
target_load_type: table.type
}
if (table.type === 'incremental_with_qualify') {
body['incremental_column_name'] = table.reference_column.name;
body['incremental_column_type'] = table.reference_column.type;
body['primary_keys'] = table.identifier_columns;
}
if (table.type === 'incremental') {
body['incremental_column_name'] = table.reference_column.name;
body['incremental_column_type'] = table.reference_column.type;
}
this.logger.info('Updating pipeline sync mode', {
pipelineId,
body
});
return this.platformAPI.proxy(
"POST",
`/jobs/jdbc/${pipelineId}/sync-mode`,
user,
body
)
}
private async updatePipelineCron(cron: string, pipelineId: string, user: RequestUser) {
try {
const response = await this.platformAPI.proxy(
'PATCH',
`/pipeline/${pipelineId}`,
user,
{
cron
}
);
return response
} catch (error) {
return {
error: error.message
}
}
}
}
@@ -113,7 +113,7 @@ export class StorageExplorerService {
}
// Get customer-specific storage-explorer URL
const baseUrl = STORAGE_EXPLORER_CONFIG.getUrl(user.customer_id);
const baseUrl = STORAGE_EXPLORER_CONFIG.getUrl(user.customer_name);
const url = new URL(`${baseUrl}${path}`);
// Add query params