Compare commits

...
Author SHA1 Message Date
marcos-silva-rodrigues 65ab16236f FEAT: update deployment to include firebase base url 2026-05-25 14:09:16 -03:00
marcos-silva-rodrigues b4cc8151d7 FIX: release note endpoint 2026-05-25 14:01:00 -03:00
Marcos Rodrigues Silva a5d78a97ae Merge pull request #486 from dadosfera/beta
Beta
2026-05-04 18:00:02 -03:00
Marcos Rodrigues Silva 2b33c22149 Merge pull request #488 from dadosfera/feature/table-schema-filter
FIX: roles
2026-04-29 18:02:04 -03:00
marcos-silva-rodrigues e32787baff FIX: roles 2026-04-29 17:06:04 -03:00
Marcos Rodrigues Silva f449d8ebd9 Merge pull request #487 from dadosfera/feature/table-schema-filter
FEAT: list assets by pipeline and multiple ids
2026-04-29 16:50:52 -03:00
marcos-silva-rodrigues 02627d023a FEAT: list assets by pipeline and multiple ids 2026-04-29 16:26:43 -03:00
Marcos Rodrigues Silva cf94c73648 Merge pull request #485 from dadosfera/feature/table-schema-filter
FIX: update package lock
2026-04-23 15:36:18 -03:00
marcos-silva-rodrigues 6b3241281b FIX: update package lock 2026-04-23 15:33:18 -03:00
Marcos Rodrigues Silva fe1caa003e Merge pull request #484 from dadosfera/feature/table-schema-filter
Feature/table schema filter
2026-04-23 15:16:11 -03:00
marcos-silva-rodrigues e980c58507 FEAT: list schema and filter assets by schema 2026-04-23 15:15:14 -03:00
Marcos Rodrigues Silva ff5f8672e2 Merge pull request #483 from dadosfera/beta
Beta
2026-04-20 10:46:37 -03:00
Marcos Rodrigues Silva d51456ecf7 Merge pull request #482 from dadosfera/hotfix/disable-edit-pipeline
FIX: block cancel first pipeline
2026-04-16 15:02:20 -03:00
marcos-silva-rodrigues 0cb9da13d0 FIX: block cancel first pipeline 2026-04-16 14:58:19 -03:00
Marcos Rodrigues Silva f198449c16 Merge pull request #481 from dadosfera/hotfix/disable-edit-pipeline
Hotfix/disable edit pipeline
2026-04-13 14:29:58 -03:00
marcos-silva-rodrigues 0d627b0451 FEAT: guard to prevent pipeline update when pipeline is running 2026-04-13 14:23:44 -03:00
24 changed files with 502 additions and 320 deletions
@@ -113,6 +113,8 @@ spec:
value: {{ .Values.maestro.platform_api_url }}
- name: STORAGE_EXPLORER_API_URL
value: {{ .Values.maestro.storage_explorer_api_url | quote }}
- name: FIREBASE_BASE_URL
value: {{ .Values.maestro.firebase_base_url }}
- name: JWT_PRIVATE_KEY
valueFrom:
secretKeyRef:
+1
View File
@@ -10,6 +10,7 @@ maestro:
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:8000/api"
firebase_base_url: https://feature-flag-25bf6-default-rtdb.firebaseio.com/stg
hostname: maestro.stg.dadosfera.ai
+1
View File
@@ -55,6 +55,7 @@ maestro:
redis_database: "0"
redis_tls: "true"
cookie_secret: "13cc5e136d3074bcc05bec8697092ec1f5f376bf"
firebase_base_url: https://feature-flag-25bf6-default-rtdb.firebaseio.com/prd
autoscaling:
enabled: false
minReplicas: 1
+135 -230
View File
@@ -3541,6 +3541,64 @@
]
}
},
"/pipelinesV2/{id}/data-assets": {
"get": {
"operationId": "PipelinesController_findAllDataAssetByPipeline",
"parameters": [
{
"name": "dadosfera-lang",
"in": "header",
"required": false,
"schema": {
"enum": [
"pt-br",
"en-us"
],
"type": "string"
}
},
{
"name": "id",
"required": true,
"in": "path",
"schema": {
"type": "string"
}
},
{
"name": "object",
"required": true,
"in": "query",
"schema": {
"type": "string"
}
}
],
"responses": {
"200": {
"description": "",
"content": {
"application/json": {
"schema": {
"type": "array",
"items": {
"type": "object"
}
}
}
}
}
},
"tags": [
"PipelinesV2"
],
"security": [
{
"access-token": []
}
]
}
},
"/pipelinesV2/{pipelineId}/inputs/{id}": {
"patch": {
"operationId": "PipelinesController_updatePipelineInput",
@@ -4001,7 +4059,7 @@
]
}
},
"/platform/pipeline": {
"/platform/pipelines": {
"post": {
"operationId": "PlatformApiController_createPipeline",
"summary": "Create a new pipeline",
@@ -4026,9 +4084,7 @@
"access-token": []
}
]
}
},
"/platform/pipelines": {
},
"get": {
"operationId": "PlatformApiController_getPipelines",
"summary": "List all pipelines for customer",
@@ -4055,7 +4111,7 @@
]
}
},
"/platform/pipeline/{pipelineId}": {
"/platform/pipelines/{pipelineId}": {
"get": {
"operationId": "PlatformApiController_getPipeline",
"summary": "Get pipeline by ID",
@@ -4159,7 +4215,7 @@
]
}
},
"/platform/pipeline/execute": {
"/platform/pipelines/execute": {
"post": {
"operationId": "PlatformApiController_executePipeline",
"summary": "Execute a pipeline",
@@ -4186,7 +4242,7 @@
]
}
},
"/platform/pipeline/pause": {
"/platform/pipelines/pause": {
"post": {
"operationId": "PlatformApiController_pausePipeline",
"summary": "Pause a pipeline",
@@ -4213,7 +4269,7 @@
]
}
},
"/platform/pipeline/unpause": {
"/platform/pipelines/unpause": {
"post": {
"operationId": "PlatformApiController_unpausePipeline",
"summary": "Unpause a pipeline",
@@ -4240,7 +4296,7 @@
]
}
},
"/platform/pipeline/{pipelineId}/memory": {
"/platform/pipelines/{pipelineId}/memory": {
"put": {
"operationId": "PlatformApiController_updatePipelineMemory",
"summary": "Update pipeline memory configuration",
@@ -4276,7 +4332,7 @@
]
}
},
"/platform/pipeline/{pipelineId}/metadata": {
"/platform/pipelines/{pipelineId}/metadata": {
"put": {
"operationId": "PlatformApiController_updatePipelineMetadata",
"summary": "Update pipeline metadata",
@@ -4403,7 +4459,7 @@
]
}
},
"/platform/pipeline/{pipelineId}/pipeline_run": {
"/platform/pipelines/{pipelineId}/pipeline_run": {
"get": {
"operationId": "PlatformApiController_getPipelineRuns",
"summary": "Get pipeline runs for a pipeline",
@@ -4439,7 +4495,7 @@
]
}
},
"/platform/pipeline/{pipelineId}/pipeline_run/{runId}": {
"/platform/pipelines/{pipelineId}/pipeline_run/{runId}": {
"get": {
"operationId": "PlatformApiController_getPipelineRun",
"summary": "Get specific pipeline run",
@@ -4483,7 +4539,7 @@
]
}
},
"/platform/pipeline/pipeline_run/{runId}/logs": {
"/platform/pipelines/pipeline_run/{runId}/logs": {
"get": {
"operationId": "PlatformApiController_getPipelineRunLogs",
"summary": "Get pipeline run logs",
@@ -4519,7 +4575,7 @@
]
}
},
"/platform/pipeline/{pipelineId}/pipeline_run/{runId}/cancel": {
"/platform/pipelines/{pipelineId}/pipeline_run/{runId}/cancel": {
"post": {
"operationId": "PlatformApiController_cancelPipelineRun",
"summary": "Cancel a running pipeline run",
@@ -4742,114 +4798,6 @@
]
}
},
"/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/{jobId}/rename-tables": {
"post": {
"operationId": "PlatformApiController_renameJobTables",
"summary": "Rename job output tables and sync to catalog",
"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",
@@ -4877,114 +4825,6 @@
]
}
},
"/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",
@@ -5649,6 +5489,48 @@
]
}
},
"/catalog/schemas": {
"get": {
"operationId": "CatalogController_findSchemas",
"parameters": [
{
"name": "dadosfera-lang",
"in": "header",
"required": false,
"schema": {
"enum": [
"pt-br",
"en-us"
],
"type": "string"
}
}
],
"responses": {
"200": {
"description": "",
"content": {
"application/json": {
"schema": {
"type": "object"
}
}
}
}
},
"tags": [
"Catalog"
],
"security": [
{
"access-token": []
},
{
"access-token": []
}
]
}
},
"/catalog/data-asset/{id}": {
"get": {
"operationId": "CatalogController_getDataAsset",
@@ -8714,6 +8596,29 @@
"Health"
]
}
},
"/release_note": {
"get": {
"operationId": "ReleaseNoteController_getLatestReleaseNote",
"parameters": [],
"responses": {
"200": {
"description": "",
"content": {
"application/json": {
"schema": {
"type": "object"
}
}
}
}
},
"security": [
{
"access-token": []
}
]
}
}
},
"info": {
+5 -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.40.0-beta.5",
"@dadosfera/protospack-v2": "3.40.0-beta.8",
"@grpc/grpc-js": "^1.9.3",
"@grpc/proto-loader": "^0.7.9",
"@nestjs/cli": "^9.5.0",
@@ -1745,9 +1745,10 @@
}
},
"node_modules/@dadosfera/protospack-v2": {
"version": "3.40.0-beta.5",
"resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.40.0-beta.5.tgz",
"integrity": "sha512-iocKv/XXp2jKAasO5ONgm31cKfLgNsU4pEKZMr6YnR7nQaH11WcW7rnuagNxWoik++wLUqbYyf0bZWRDzMlCPA==",
"version": "3.40.0-beta.8",
"resolved": "https://dadosfera-611330257153.d.codeartifact.us-east-1.amazonaws.com/npm/dadosfera-npm/@dadosfera/protospack-v2/-/protospack-v2-3.40.0-beta.8.tgz",
"integrity": "sha512-JE5qMjqB3UOM+tCUxB1EwYLQW0PecsaQIa1KDpKEaG3lrzG/H13z8iJi3WH/DuVav2EI94i9VcJWJ1Y0F7ribw==",
"license": "ISC",
"dependencies": {
"@grpc/grpc-js": "^1.9.3",
"rxjs": "^7.5.5"
+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.40.0-beta.5",
"@dadosfera/protospack-v2": "3.40.0-beta.8",
"@grpc/grpc-js": "^1.9.3",
"@grpc/proto-loader": "^0.7.9",
"@nestjs/cli": "^9.5.0",
+3
View File
@@ -35,6 +35,8 @@ import { ShareMetadataModule } from './modules/share-metadata/share-metadata.mod
import { ApiKeyModule } from './modules/api-key/api-key.module';
import { PlatformApiModule } from './modules/platform-api/platform-api.module';
import { StorageExplorerModule } from './modules/storage-explorer/storage-explorer.module';
import { ReleaseNoteModule } from './modules/release_note/release_note.module';
@Module({
providers: [
@@ -79,6 +81,7 @@ import { StorageExplorerModule } from './modules/storage-explorer/storage-explor
StorageExplorerModule,
//Always leave HealthModule last, so it is on the bottom of swagger
HealthModule,
ReleaseNoteModule,
],
})
export class AppModule {}
@@ -153,6 +153,7 @@ export class AuthenticationGuard
user_id: accessTokenPayload.user_id,
username: accessTokenPayload.username,
permissions: accessTokenPayload.permissions,
roles: accessTokenPayload.roles,
customer_id: accessTokenPayload.customer_id,
customer_name: accessTokenPayload.customer_name,
customer_tier: accessTokenPayload.customer_tier,
+1
View File
@@ -11,6 +11,7 @@ export function extractUserFrom(aRawJwt: string) {
user_id: payload.user_id,
username: payload.username,
permissions: payload.permissions,
roles: payload.roles,
customer_id: payload.customer_id,
customer_name: payload.customer_name,
customer_tier: payload.customer_tier,
+1
View File
@@ -12,6 +12,7 @@ export interface RequestUser {
customer_tier: string;
access_token: string;
customer_modules: string[];
roles: string[];
}
export const User: (options?: { required?: boolean }) => ParameterDecorator =
+65
View File
@@ -0,0 +1,65 @@
import {
BadRequestException,
CanActivate,
ExecutionContext,
Inject,
Injectable,
OnModuleInit,
} from '@nestjs/common';
import { ClientGrpc } from '@nestjs/microservices';
import { map, Observable } from 'rxjs';
import { PackTheMetadata } from 'src/utils/PackTheMetadata';
import {
ReadService,
ProtoServices,
} from '@dadosfera/protospack-v2/dist/lib/PipelineV2';
import { PipelinesClientConfiguration } from 'src/modules/pipelinesV2/pipelines-client';
import { PlatformApiService } from 'src/modules/platform-api/platform-api.service';
import DadosferaLogger from '@dadosfera/dadosfera-logs';
@Injectable()
export class PipelineExecutionGuard implements CanActivate {
logger: DadosferaLogger;
constructor(
@Inject(DadosferaLogger)
dadosferaLogger: DadosferaLogger,
private readonly platformApiService: PlatformApiService,
) {
this.logger = dadosferaLogger.logger;
}
async canActivate(context: ExecutionContext): Promise<boolean> {
try {
this.logger.info(
'PipelineExecutionGuard: Checking if pipeline can be executed...',
);
const request = context.switchToHttp().getRequest();
const pipelineId = request.params.pipelineId;
const user = request.user;
const idRegex = /[^0-9a-zA-Z_$]+/g;
const convertedId = pipelineId.replace(idRegex, '_');
const status = await this.platformApiService.proxy(
'GET',
`/pipeline/${convertedId}/pipeline_run`,
user,
);
const currentStatus = status[status.length - 1]
this.logger.info('Pipeline current status response:' + JSON.stringify(currentStatus));
if (currentStatus.last_status.toLowerCase() === 'running') {
this.logger.error('Pipeline is running, cannot update input now');
throw new BadRequestException('Pipeline is running, cannot update input now');
} else {
return true;
}
} catch (error) {
this.logger.error('Error in PipelineExecutionGuard: ' + error.message);
throw new BadRequestException('Error checking pipeline status: ' + error.message);
}
}
}
+28
View File
@@ -241,6 +241,34 @@ export class CatalogController {
return res;
}
@Get('schemas')
@RequireSomePermission(
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
)
async findSchemas(@User() user: RequestUser) {
const { username, user_id, customer_id, customer_name } = user;
this.logger.info(`/catalog - ON FIND SCHEMAS ROUTE`, {
username,
customer_name,
});
const metadata = PackTheMetadata({
username,
user_id,
customer_id,
customer_name,
});
try {
const res = await this.catalogService.findSchemas(metadata);
return res;
} catch (error) {
throw new HttpException(error.message, HttpStatus.NOT_FOUND);
}
}
@Get('data-asset/:id')
@RequireSomePermission(
PERMISSIONS_GROUPS.CATALOG.permissions.GET,
+16 -18
View File
@@ -214,15 +214,6 @@ class CatalogService implements OnModuleInit {
this.logger.debug('Extracted filters:', { filters });
console.log('MAESTRO VAI CHAMAR PI-FACTORY COM (ANTES AJUSTE):', {
search,
page,
size,
sort_by,
order,
filters,
});
if (
filters.manually !== undefined &&
filters.manually !== null &&
@@ -233,15 +224,6 @@ class CatalogService implements OnModuleInit {
delete filters.manually;
}
console.log('MAESTRO VAI CHAMAR PI-FACTORY COM (DEPOIS AJUSTE):', {
search,
page,
size,
sort_by,
order,
filters,
});
if (filters.owner) {
const { users: customer_users } =
await this.userService.findAllUsersByCustomerId(customer_id);
@@ -517,6 +499,22 @@ class CatalogService implements OnModuleInit {
return response;
}
async findSchemas(metadata: Metadata) {
this.logger.info('CatalogService - findSchemas');
try {
const response = await lastValueFrom(
this.catalogReadService.GetSchemas({}, metadata),
);
return response;
} catch (error) {
this.logger.error('Error fetching schemas:', error);
throw error;
}
}
async getAssetsUsersAndRoles(data_assets: Array<any>, customer_id: string) {
const { users: customer_users } =
await this.userService.findAllUsersByCustomerId(customer_id);
@@ -262,6 +262,7 @@ export class ShareService implements OnModuleInit {
user_id: accessTokenPayload.user_id,
username: accessTokenPayload.username,
permissions: accessTokenPayload.permissions,
roles: accessTokenPayload.roles,
customer_id: accessTokenPayload.customer_id,
customer_name: accessTokenPayload.customer_name,
customer_tier: accessTokenPayload.customer_tier,
@@ -14,7 +14,7 @@ import {
Patch,
HttpException,
BadRequestException,
CacheTTL,
UseGuards,
} from '@nestjs/common';
import {
ApiCreatedResponse,
@@ -24,7 +24,6 @@ import {
ApiTags,
} from '@nestjs/swagger';
import {
AuthenticateCondition,
RequireAllPermissions,
RequireSomePermission,
} from 'src/decorators/authentication.decorator';
@@ -49,9 +48,8 @@ import { GrpcToHttpExceptionFilter } from 'src/error/grpc-to-http-exception.filt
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';
import { PipelineExecutionGuard } from 'src/guards/pipeline-execution.guard';
type PipelineTable = { name: string; job_id?: string; is_deleted?: boolean; [key: string]: any };
type PipelineTablesConfig = { input_id?: string; tables: PipelineTable[] };
@@ -195,6 +193,7 @@ export class PipelinesController {
@Get(':id/status')
@RequireSomePermission(PERMISSIONS_GROUPS.IMPORT_FILES.permissions.VIEW, PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
async getPipelineStatus(@Body() body, @Param('id') id: string) {
body.id = id;
this.logger.info(`/pipeline/${id} - ON GET PIPELINE STATUS ROUTE`, {
@@ -249,6 +248,31 @@ export class PipelinesController {
return pipelineRes;
}
@Get("/:id/data-assets")
@RequireSomePermission(
PERMISSIONS_GROUPS.PIPELINE.permissions.GET,
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER,
PERMISSIONS_GROUPS.CATALOG.permissions.GET
)
async findAllDataAssetByPipeline(
@Language() language: LanguageEnum,
@Param('id') id: string,
@User() user: RequestUser,
@Query('object') object: string
) {
const payload = {
pipeline: id,
object: object,
};
this.logger.info(`GET pipelinesV2/:id/data-assets` + JSON.stringify(payload));
const result =
await this.pipelinesClientService.findAllDataAssetByPipeline(payload, user);
return result;
}
@Patch('/:id')
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async update(
@@ -286,6 +310,7 @@ export class PipelinesController {
@Patch('/:pipelineId/inputs/:id')
@RequireSomePermission(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
@UseGuards(PipelineExecutionGuard)
async updatePipelineInput(
@Language() language: LanguageEnum,
@Body() pipelineInputDTO: UpdatePlatformInputRequest,
+3 -1
View File
@@ -14,6 +14,7 @@ import { TransformationsModule } from '../transformations/transformations.module
import { PlatformApiModule } from '../platform-api/platform-api.module';
import { NimbusServicesModule } from 'src/services/nimbus/nimbus.module';
import { NimbusService } from 'src/services/nimbus/nimbus.service';
import { CatalogModule } from '../catalog/catalog.module';
const client = new PipelinesClientConfiguration();
@@ -25,7 +26,8 @@ const client = new PipelinesClientConfiguration();
InputsModule,
TransformationsModule,
PlatformApiModule,
NimbusServicesModule
NimbusServicesModule,
CatalogModule
],
controllers: [PipelinesController],
providers: [PipelinesService, DadosferaLogger, NimbusService],
+47 -1
View File
@@ -32,6 +32,10 @@ import { Info } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/entitie
import { TableUpdate } from '@dadosfera/protospack-v2/dist/lib/Input/interfaces/messages';
import { AxiosError } from 'axios';
import { NimbusService } from 'src/services/nimbus/nimbus.service';
import { PERMISSIONS_GROUPS } from 'src/authentication/permissions.enum';
import { PackTheMetadata } from 'src/utils/PackTheMetadata';
import { IDataAsset } from '../catalog/dtos';
import { CatalogService } from '../catalog/catalog.service';
type RollbackPromise = () => Promise<any>;
@@ -48,7 +52,8 @@ export class PipelinesService implements OnModuleInit {
private readonly inputsService: InputsService,
private readonly transformationsService: TransformationsService,
private readonly platformAPI: PlatformApiService,
private readonly nimbusService: NimbusService
private readonly nimbusService: NimbusService,
private readonly catalogService: CatalogService
) {
this.logger = dadosferaLogger.logger;
}
@@ -577,4 +582,45 @@ export class PipelinesService implements OnModuleInit {
this.logger.info('Platform api response: ' + JSON.stringify(response));
}
async findAllDataAssetByPipeline(data: {
pipeline: string,
object?: string
}, user: RequestUser) {
const metadata = PackTheMetadata(user);
const isDataAdmin = user.permissions.includes(
PERMISSIONS_GROUPS.CATALOG.permissions.DATA_MANAGER.seqid,
);
let has_permission = false;
const {
data_assets: resultString
} = await lastValueFrom(
this.pipelineReadService.FindAllDataAssetByPipeline(data, metadata)
);
const result = JSON.parse(resultString) as any;
const data_assets: IDataAsset[] = []
result.forEach(data_asset => {
if (data_asset?.owner === user.username) has_permission = true;
for (const role of user.roles) {
if (data_asset.roles.includes(role)) has_permission = true;
}
if (data_asset.users.includes(user.user_id)) has_permission = true;
if (isDataAdmin || has_permission) {
delete data_asset.p_roles;
delete data_asset.p_users;
data_assets.push(data_asset as IDataAsset);
}
});
const assets = await this.catalogService.getAssetsUsersAndRoles(data_assets, user.customer_id);
return assets;
}
}
@@ -12,6 +12,7 @@ import {
BadRequestException,
HttpException,
NotFoundException,
UseGuards,
} from '@nestjs/common';
import { ApiTags, ApiOperation } from '@nestjs/swagger';
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
@@ -31,6 +32,7 @@ import { CatalogService } from '../catalog/catalog.service';
import { PackTheMetadata } from '../../utils/PackTheMetadata';
import { ValidationTableDTO } from './platform-api.dto';
import { InputsService } from '../inputs/inputs.service';
import { PipelineExecutionGuard } from 'src/guards/pipeline-execution.guard';
type ValidateTablesDTO = {
@@ -132,7 +134,7 @@ export class PlatformApiController {
}
private readonly VALID_CONNECTORS = ['jdbc', 'singer', 's3'];
private readonly MAX_MEMORY_MB = 12000; // 12GB maximum memory per pipeline/job
private readonly MAX_MEMORY_MB = 12000; // 12GB maximum memory per pipelines/job
/**
* Validate that connector is provided and is a valid type.
@@ -403,55 +405,9 @@ export class PlatformApiController {
}
}
/**
* Sync sync-mode changes to DynamoDB for JDBC connectors.
* Always passes both target_load_type and incremental_column_name to ensure proper sync.
*/
private async syncJdbcSyncModeToDynamoDB(
jobId: string,
body: any,
user: RequestUser,
): Promise<void> {
// JDBC sync mode uses target_load_type field
const changes: any = {};
if ('target_load_type' in body) {
changes.target_load_type = body.target_load_type;
}
// Handle incremental_column_name:
// - If provided in body, use that value
// - If changing to full_load, explicitly clear it
if ('incremental_column_name' in body) {
changes.incremental_column_name = body.incremental_column_name;
changes.incremental_column_type = body.incremental_column_type;
} else if (body.target_load_type === 'full_load') {
// Changing to full_load without specifying incremental_column - clear it
changes.incremental_column_name = null;
}
await this.syncJobInputToDynamoDB(jobId, changes, user, 'jdbc');
}
/**
* Sync sync-mode changes to DynamoDB for Singer connectors.
*/
private async syncSingerSyncModeToDynamoDB(
jobId: string,
body: any,
user: RequestUser,
): Promise<void> {
// Singer sync mode uses replication_method field
// Map to DynamoDB type: FULL_TABLE -> full_load, INCREMENTAL -> incremental
if ('replication_method' in body) {
const type = body.replication_method === 'INCREMENTAL' ? 'incremental' : 'full_load';
await this.syncJobInputToDynamoDB(jobId, { load_type: type }, user, 'singer');
}
}
// ==================== PIPELINE ROUTES ====================
@Post('pipeline')
@Post('pipelines')
@ApiOperation({ summary: 'Create a new pipeline' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.CREATE)
async createPipeline(@Body() body: any, @User() user: RequestUser) {
@@ -572,7 +528,7 @@ export class PlatformApiController {
);
}
@Get('pipeline/:pipelineId')
@Get('pipelines/:pipelineId')
@ApiOperation({ summary: 'Get pipeline by ID' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
async getPipeline(
@@ -583,7 +539,7 @@ export class PlatformApiController {
return this.platformApiService.proxy('GET', `/pipeline/${normalizedId}`, user);
}
@Patch('pipeline/:pipelineId')
@Patch('pipelines/:pipelineId')
@ApiOperation({ summary: 'Update pipeline by ID' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async updatePipeline(
@@ -634,7 +590,7 @@ export class PlatformApiController {
return result;
}
@Delete('pipeline/:pipelineId')
@Delete('pipelines/:pipelineId')
@ApiOperation({ summary: 'Delete pipeline by ID' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE)
async deletePipeline(
@@ -664,7 +620,7 @@ export class PlatformApiController {
return result;
}
@Post('pipeline/execute')
@Post('pipelines/execute')
@ApiOperation({ summary: 'Execute a pipeline' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async executePipeline(@Body() body: any, @User() user: RequestUser) {
@@ -674,10 +630,10 @@ export class PlatformApiController {
...body,
customer_id: user.customer_name,
};
return this.platformApiService.proxy('POST', '/pipeline/execute', user, enrichedBody);
return this.platformApiService.proxy('POST', '/pipelines/execute', user, enrichedBody);
}
@Post('pipeline/pause')
@Post('pipelines/pause')
@ApiOperation({ summary: 'Pause a pipeline' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async pausePipeline(@Body() body: any, @User() user: RequestUser) {
@@ -690,7 +646,7 @@ export class PlatformApiController {
return this.platformApiService.proxy('POST', '/pipeline/pause', user, enrichedBody);
}
@Post('pipeline/unpause')
@Post('pipelines/unpause')
@ApiOperation({ summary: 'Unpause a pipeline' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async unpausePipeline(@Body() body: any, @User() user: RequestUser) {
@@ -703,7 +659,7 @@ export class PlatformApiController {
return this.platformApiService.proxy('POST', '/pipeline/unpause', user, enrichedBody);
}
@Put('pipeline/:pipelineId/memory')
@Put('pipelines/:pipelineId/memory')
@ApiOperation({ summary: 'Update pipeline memory configuration' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async updatePipelineMemory(
@@ -726,7 +682,7 @@ export class PlatformApiController {
// ==================== PIPELINE METADATA ROUTES ====================
@Put('pipeline/:pipelineId/metadata')
@Put('pipelines/:pipelineId/metadata')
@ApiOperation({ summary: 'Update pipeline metadata' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async updatePipelineMetadata(
@@ -795,7 +751,7 @@ export class PlatformApiController {
// ==================== PIPELINE RUN ROUTES ====================
@Get('pipeline/:pipelineId/pipeline_run')
@Get('pipelines/:pipelineId/pipeline_run')
@ApiOperation({ summary: 'Get pipeline runs for a pipeline' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
async getPipelineRuns(
@@ -813,7 +769,7 @@ export class PlatformApiController {
);
}
@Get('pipeline/:pipelineId/pipeline_run/:runId')
@Get('pipelines/:pipelineId/pipeline_run/:runId')
@ApiOperation({ summary: 'Get specific pipeline run' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
async getPipelineRun(
@@ -830,7 +786,7 @@ export class PlatformApiController {
);
}
@Get('pipeline/pipeline_run/:runId/logs')
@Get('pipelines/pipeline_run/:runId/logs')
@ApiOperation({ summary: 'Get pipeline run logs' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
async getPipelineRunLogs(
@@ -848,7 +804,7 @@ export class PlatformApiController {
);
}
@Post('pipeline/:pipelineId/pipeline_run/:runId/cancel')
@Post('pipelines/:pipelineId/pipeline_run/:runId/cancel')
@ApiOperation({ summary: 'Cancel a running pipeline run' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async cancelPipelineRun(
@@ -859,6 +815,16 @@ export class PlatformApiController {
const normalizedPipelineId = this.normalizePipelineId(pipelineId);
const normalizedRunId = this.normalizePipelineId(runId);
const status = await this.platformApiService.proxy(
'GET',
`/pipeline/${normalizedPipelineId}/pipeline_run`,
user,
);
if (status.length === 1) {
throw new BadRequestException('The first pipeline cannot be canceled');
}
return this.platformApiService.proxy(
'POST',
`/pipeline/${normalizedPipelineId}/pipeline_run/${normalizedRunId}/cancel`,
@@ -968,6 +934,7 @@ export class PlatformApiController {
@Delete('pipelines/:pipelineId/inputs/:inputId')
@ApiOperation({ summary: 'Mark a table as deleted and delete its associated job via platform-api' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.DELETE)
@UseGuards(PipelineExecutionGuard)
async deleteTable(
@Param('pipelineId') pipelineId: string,
@Param('inputId') inputId: string,
@@ -0,0 +1,14 @@
export type ReleaseNoteDTO = {
id: string;
date: string;
tag: string;
title: string;
visible: boolean;
expiryDate: string;
content: string;
showEmojis: boolean;
image?: string;
link?: string;
linkText?: string;
};
@@ -0,0 +1,20 @@
import { Test, TestingModule } from '@nestjs/testing';
import { ReleaseNoteController } from './release_note.controller';
import { ReleaseNoteService } from './release_note.service';
describe('ReleaseNoteController', () => {
let controller: ReleaseNoteController;
beforeEach(async () => {
const module: TestingModule = await Test.createTestingModule({
controllers: [ReleaseNoteController],
providers: [ReleaseNoteService],
}).compile();
controller = module.get<ReleaseNoteController>(ReleaseNoteController);
});
it('should be defined', () => {
expect(controller).toBeDefined();
});
});
@@ -0,0 +1,26 @@
import { Controller, Get, Inject } from '@nestjs/common';
import { ReleaseNoteService } from './release_note.service';
import { Authenticated } from 'src/decorators/authentication.decorator';
import { Language } from 'src/decorators/language.decorator';
import { LanguageEnum } from 'src/utils/languages.enum';
import DadosferaLogger from '@dadosfera/dadosfera-logs';
@Controller('release_note')
@Authenticated()
export class ReleaseNoteController {
logger: DadosferaLogger;
constructor(
@Inject(DadosferaLogger)
dadosferaLogger: DadosferaLogger,
private readonly releaseNoteService: ReleaseNoteService,
) {
this.logger = dadosferaLogger.logger;
}
@Get()
async getLatestReleaseNote(@Language() language: LanguageEnum) {
this.logger.info(`Fetching latest release note for language: ${language}`);
return await this.releaseNoteService.getLatestReleaseNote(language);
}
}
@@ -0,0 +1,10 @@
import { Module } from '@nestjs/common';
import { ReleaseNoteService } from './release_note.service';
import { ReleaseNoteController } from './release_note.controller';
import DadosferaLogger from '@dadosfera/dadosfera-logs';
@Module({
controllers: [ReleaseNoteController],
providers: [ReleaseNoteService, DadosferaLogger]
})
export class ReleaseNoteModule {}
@@ -0,0 +1,18 @@
import { Test, TestingModule } from '@nestjs/testing';
import { ReleaseNoteService } from './release_note.service';
describe('ReleaseNoteService', () => {
let service: ReleaseNoteService;
beforeEach(async () => {
const module: TestingModule = await Test.createTestingModule({
providers: [ReleaseNoteService],
}).compile();
service = module.get<ReleaseNoteService>(ReleaseNoteService);
});
it('should be defined', () => {
expect(service).toBeDefined();
});
});
@@ -0,0 +1,46 @@
import { Inject, Injectable } from '@nestjs/common';
import axios, { AxiosInstance } from 'axios';
import { LanguageEnum } from 'src/utils/languages.enum';
import { ReleaseNoteDTO } from './dto/release_note.dto';
import DadosferaLogger from '@dadosfera/dadosfera-logs';
@Injectable()
export class ReleaseNoteService {
client: AxiosInstance;
logger: DadosferaLogger;
constructor(
@Inject(DadosferaLogger)
dadosferaLogger: DadosferaLogger,
) {
this.logger = dadosferaLogger.logger;
this.client = axios.create({
baseURL: process.env.FIREBASE_BASE_URL,
});
}
async getLatestReleaseNote(lang: LanguageEnum) {
try {
const lng = lang.split('-');
const language = lng[0] + '-' + lng[1].toUpperCase();
const endpoint = `/release_note/${language}.json`;
const {
data,
status,
config
} = await this.client.get<ReleaseNoteDTO>(endpoint)
this.logger.info(`Fetched release note for language: ${lang} with status: ${status}`);
this.logger.info(`Request URL: ${config.baseURL}/${config.url}`);
return data;
} catch (error) {
this.logger.error(`Error fetching release note: ${error.message}`);
if (axios.isAxiosError(error)) {
this.logger.error(`Axios error details: ${error.toJSON()}`);
}
}
}
}