Compare commits

..
Author SHA1 Message Date
Rafael Santana 2125884c6c Merge pull request #459 from dadosfera/force-deploy
FIX: uppercase table_name and table_schema in ES lookup
2026-03-16 10:19:02 -03:00
RafaelandClaude Opus 4.6 e3099aa2b2 FIX: uppercase table_name and table_schema in ES lookup
Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-16 10:17:56 -03:00
Rafael Santana f4c9226ef9 Merge pull request #458 from dadosfera/force-deploy
FIX: rename-tables proxy path and ES lookup
2026-03-13 18:29:53 -03:00
RafaelandClaude Opus 4.6 12c61d9b5d FIX: rename-tables proxy path and ES lookup
- Fix proxy path: /jobs/jdbc/:jobId/rename-tables → /jobs/:jobId/rename-tables
- Remove pipeline_id from ES data asset lookup, search by table_name + table_schema only

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-13 18:08:39 -03:00
Rafael Santana 209470482a Merge pull request #456 from dadosfera/force-deploy
UPDATE: force deployment of maestro
2026-03-12 18:03:19 -03:00
Rafael 99c2a9ecf5 UPDATE: force deployment of maestro 2026-03-12 18:02:49 -03:00
Rafael Santana c0f75d241f Merge pull request #449 from dadosfera/feat/rename-tables-catalog-sync
Feat/rename tables catalog sync
2026-03-12 17:57:44 -03:00
Rafael Santana 1e0fb78dff Merge branch 'beta' into feat/rename-tables-catalog-sync 2026-03-12 17:57:37 -03:00
Marcos Rodrigues Silva b70d37423d Merge pull request #454 from dadosfera/fix/header-validation
FIX: types
2026-03-11 10:27:31 -03:00
marcos-silva-rodrigues 4ddd5edcfd FIX: types 2026-03-11 10:26:33 -03:00
Marcos Rodrigues Silva 3bcbba9581 Merge pull request #453 from dadosfera/fix/header-validation
Fix/header validation
2026-03-10 17:35:31 -03:00
marcos-silva-rodrigues 67d47a9642 FIX: correct header 2026-03-10 17:35:05 -03:00
Marcos Rodrigues Silva 6f9c967c96 Merge pull request #452 from dadosfera/feat/qualify
FIX: stringify headers
2026-03-10 17:17:48 -03:00
marcos-silva-rodrigues b8bdc5beea FIX: stringify headers 2026-03-10 16:00:39 -03:00
Marcos Rodrigues Silva 269f70b309 Merge pull request #451 from dadosfera/feat/qualify
FIX: ghost commit
2026-03-10 15:22:36 -03:00
marcos-silva-rodrigues 93f452ae05 FIX: ghost commit 2026-03-10 15:19:18 -03:00
Marcos Rodrigues Silva e39378229f Merge pull request #450 from dadosfera/feat/qualify
feat: list headers
2026-03-10 14:43:43 -03:00
marcos-silva-rodrigues 61a4f724ef feat: list headers 2026-03-10 14:43:11 -03:00
RafaelandClaude Opus 4.6 38a9e21f5f feat: add rename-tables endpoint with catalog sync and rollback
Add POST /platform/jobs/:jobId/rename-tables that renames Snowflake
tables via platform-api and syncs the rename to Elasticsearch and
Nimbus (table-metadata, column-metadata, data-preview). If catalog
sync fails, all completed catalog steps are rolled back in reverse
order and the Snowflake rename is reverted.

- Support any connector type (jdbc, singer, s3) via getJobByAnyConnectorType
- Resolve old table names from output_config (raw/qualify)
- Skip qualify sync when output_config.qualify has no table_name
- Add findDataAssetByPipelineAndTable and updateDataAsset to ElasticsearchService
- Add renameTableOnNimbus, renameColumnMetadataOnNimbus, renameDataPreviewOnNimbus to CatalogService
- Add upstream error logging to PlatformApiService

Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
2026-03-05 17:29:12 -03:00
Marcos Rodrigues Silva 1a2153d62f Merge pull request #448 from dadosfera/fix/proxy-urls
Fix/proxy urls
2026-02-20 16:03:05 -03:00
Marcos Rodrigues Silva d25bfd147c Merge pull request #447 from dadosfera/fix/proxy-urls
Fix/proxy urls
2026-02-20 14:32:50 -03:00
marcos-silva-rodrigues d4ba45fd03 FIX: storage url 2026-02-20 14:28:55 -03:00
marcos-silva-rodrigues 82f1035a5e FIX: update storage port 2026-02-20 14:28:36 -03:00
24 changed files with 352 additions and 1575 deletions
+1
View File
@@ -4,6 +4,7 @@
# Maestro
Maestro é a API principal da Dadosfera. É responsável pela comunicação do Frontend com nossos microsserviços.
```mermaid
@@ -113,12 +113,6 @@ spec:
value: {{ .Values.maestro.platform_api_url }}
- name: STORAGE_EXPLORER_API_URL
value: {{ .Values.maestro.storage_explorer_api_url | quote }}
- name: AUTODRIVE_EXTRACTOR_API_URL
value: {{ .Values.maestro.autodrive_extractor_api_url | quote }}
- name: AUTODRIVE_CORE_API_URL
value: {{ .Values.maestro.autodrive_core_api_url | quote }}
- name: AUTODRIVE_ASSISTANT_API_URL
value: {{ .Values.maestro.autodrive_assistant_api_url | quote }}
- name: JWT_PRIVATE_KEY
valueFrom:
secretKeyRef:
-3
View File
@@ -48,9 +48,6 @@ maestro:
open_group_id: 401573bb-334f-44b2-b30e-88d4cea31ae9
platform_api_url: https://oz8v2zid1e.execute-api.us-east-1.amazonaws.com
storage_explorer_api_url: "https://storage-explorer-{customer}.dadosfera.ai/api"
autodrive_extractor_api_url: "https://autodrive-extractor-api-{customer}.dadosfera.ai"
autodrive_core_api_url: "https://autodrive-api-{customer}.dadosfera.ai"
autodrive_assistant_api_url: "https://autodrive-assistant-api-{customer}.dadosfera.ai"
dedicated_proxy: ""
restricted_ip: ""
redis_host: "aaapzppmlyamkocqwstpo7zvopczyyiyuy6xzm2g6c5k4mq3a66be4a-0.redis.sa-saopaulo-1.oci.oraclecloud.com"
-6
View File
@@ -35,9 +35,6 @@ 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 { AutodriveExtractorModule } from './modules/autodrive-extractor/autodrive-extractor.module';
import { AutodriveCoreModule } from './modules/autodrive-core/autodrive-core.module';
import { AutodriveAssistantModule } from './modules/autodrive-assistant/autodrive-assistant.module';
@Module({
providers: [
@@ -80,9 +77,6 @@ import { AutodriveAssistantModule } from './modules/autodrive-assistant/autodriv
NetworkPolicyModule,
PlatformApiModule,
StorageExplorerModule,
AutodriveExtractorModule,
AutodriveCoreModule,
AutodriveAssistantModule,
//Always leave HealthModule last, so it is on the bottom of swagger
HealthModule,
],
-87
View File
@@ -697,93 +697,6 @@ export const PERMISSIONS_GROUPS = {
},
},
},
AUTODRIVE_EXTRACTOR: {
title: {
'pt-br': 'Autodrive Extractor',
'en-us': 'Autodrive Extractor',
'es-es': 'Autodrive Extractor',
},
permissions: {
READ: {
seqid: 53,
claim: 'autodrive-extractor:read',
usage: PermissionUsages.PUBLIC,
name: {
'pt-br': 'Ler dados do Autodrive Extractor',
'en-us': 'Read Autodrive Extractor data',
'es-es': 'Leer datos del Autodrive Extractor',
},
},
WRITE: {
seqid: 54,
claim: 'autodrive-extractor:write',
usage: PermissionUsages.PUBLIC,
name: {
'pt-br': 'Escrever dados no Autodrive Extractor',
'en-us': 'Write Autodrive Extractor data',
'es-es': 'Escribir datos en Autodrive Extractor',
},
},
},
},
AUTODRIVE_CORE: {
title: {
'pt-br': 'Autodrive Core',
'en-us': 'Autodrive Core',
'es-es': 'Autodrive Core',
},
permissions: {
READ: {
seqid: 55,
claim: 'autodrive-core:read',
usage: PermissionUsages.PUBLIC,
name: {
'pt-br': 'Ler dados do Autodrive Core',
'en-us': 'Read Autodrive Core data',
'es-es': 'Leer datos del Autodrive Core',
},
},
WRITE: {
seqid: 56,
claim: 'autodrive-core:write',
usage: PermissionUsages.PUBLIC,
name: {
'pt-br': 'Escrever dados no Autodrive Core',
'en-us': 'Write Autodrive Core data',
'es-es': 'Escribir datos en Autodrive Core',
},
},
},
},
AUTODRIVE_ASSISTANT: {
title: {
'pt-br': 'Autodrive Assistant',
'en-us': 'Autodrive Assistant',
'es-es': 'Autodrive Assistant',
},
permissions: {
READ: {
seqid: 57,
claim: 'autodrive-assistant:read',
usage: PermissionUsages.PUBLIC,
name: {
'pt-br': 'Ler dados do Autodrive Assistant',
'en-us': 'Read Autodrive Assistant data',
'es-es': 'Leer datos del Autodrive Assistant',
},
},
WRITE: {
seqid: 58,
claim: 'autodrive-assistant:write',
usage: PermissionUsages.PUBLIC,
name: {
'pt-br': 'Escrever dados no Autodrive Assistant',
'en-us': 'Write Autodrive Assistant data',
'es-es': 'Escribir datos en Autodrive Assistant',
},
},
},
},
};
export interface DadosferaModule {
name: string;
+1
View File
@@ -111,3 +111,4 @@ function configureSwagger(app: INestApplication) {
);
}
bootstrap();
+2 -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(JSON.stringify(req.headers));
// Check for API key header first
const apiKey = req.get('X-Api-key');
@@ -502,7 +503,7 @@ export class AuthController {
const accessToken = req.cookies['ddf-auth'];
const refreshToken = req.cookies['ddf-refresh-auth'];
const userId = req.cookies['ddf-user-id'];
const resourceHost = req.headers["host"]
const resourceHost = req.headers["x-original-url"] as string || "" ;
const hasUserSession = Boolean(accessToken) && Boolean(userId);
this.logger.info('Has User Session: ' + hasUserSession);
@@ -1,10 +0,0 @@
export const AUTODRIVE_ASSISTANT_CONFIG = {
getUrl: (customerName: string): string => {
const urlTemplate = process.env.AUTODRIVE_ASSISTANT_API_URL;
if (!urlTemplate) {
throw new Error('AUTODRIVE_ASSISTANT_API_URL environment variable is not set');
}
return urlTemplate.replace('{customer}', customerName);
},
timeout: parseInt(process.env.AUTODRIVE_ASSISTANT_TIMEOUT || '30000', 10),
};
@@ -1,201 +0,0 @@
import {
Controller,
Get,
Post,
Put,
Delete,
Param,
Body,
Query,
Inject,
} from '@nestjs/common';
import { ApiTags, ApiOperation } from '@nestjs/swagger';
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
import {
Authenticated,
RequireAllPermissions,
} from '../../decorators/authentication.decorator';
import { User, RequestUser } from '../../decorators/user.decorator';
import { AutodriveAssistantService } from './autodrive-assistant.service';
import { PERMISSIONS_GROUPS } from '../../authentication/permissions.enum';
@ApiTags('Autodrive Assistant')
@Controller('autodrive-assistant')
@Authenticated()
export class AutodriveAssistantController {
private logger: any;
constructor(
private readonly autodriveAssistantService: AutodriveAssistantService,
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
) {
this.logger = dadosferaLogger.logger;
}
// ============================================
// ASSISTANTS
// ============================================
@ApiOperation({ summary: 'List assistants' })
@Get('assistants')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_ASSISTANT.permissions.READ)
async listAssistants(
@Query('limit') limit: number,
@Query('offset') offset: number,
@User() user: RequestUser,
) {
return this.autodriveAssistantService.proxy(
'GET',
'/assistants',
user,
undefined,
{ limit, offset },
);
}
@ApiOperation({ summary: 'Create assistant' })
@Post('assistants')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_ASSISTANT.permissions.WRITE)
async createAssistant(
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveAssistantService.proxy(
'POST',
'/assistants',
user,
body,
);
}
@ApiOperation({ summary: 'Get assistant by ID' })
@Get('assistants/:assistantId')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_ASSISTANT.permissions.READ)
async getAssistant(
@Param('assistantId') assistantId: string,
@User() user: RequestUser,
) {
return this.autodriveAssistantService.proxy(
'GET',
`/assistants/${assistantId}`,
user,
);
}
@ApiOperation({ summary: 'Update assistant' })
@Put('assistants/:assistantId')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_ASSISTANT.permissions.WRITE)
async updateAssistant(
@Param('assistantId') assistantId: string,
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveAssistantService.proxy(
'PUT',
`/assistants/${assistantId}`,
user,
body,
);
}
@ApiOperation({ summary: 'Delete assistant' })
@Delete('assistants/:assistantId')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_ASSISTANT.permissions.WRITE)
async deleteAssistant(
@Param('assistantId') assistantId: string,
@User() user: RequestUser,
) {
return this.autodriveAssistantService.proxy(
'DELETE',
`/assistants/${assistantId}`,
user,
);
}
// ============================================
// KNOWLEDGE BASE
// ============================================
@ApiOperation({ summary: 'List knowledge bases linked to assistant' })
@Get('assistants/:assistantId/knowledge-bases')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_ASSISTANT.permissions.READ)
async listKnowledgeBases(
@Param('assistantId') assistantId: string,
@User() user: RequestUser,
) {
return this.autodriveAssistantService.proxy(
'GET',
`/assistants/${assistantId}/knowledge-bases`,
user,
);
}
@ApiOperation({ summary: 'Update knowledge base associations for assistant' })
@Put('assistants/:assistantId/knowledge-bases')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_ASSISTANT.permissions.WRITE)
async updateKnowledgeBases(
@Param('assistantId') assistantId: string,
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveAssistantService.proxy(
'PUT',
`/assistants/${assistantId}/knowledge-bases`,
user,
body,
);
}
// ============================================
// QUESTIONS
// ============================================
@ApiOperation({ summary: 'Ask question via assistant' })
@Post('dataset/:datasetId/ai_question')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_ASSISTANT.permissions.WRITE)
async aiQuestion(
@Param('datasetId') datasetId: string,
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveAssistantService.proxy(
'POST',
`/dataset/${datasetId}/ai_question`,
user,
body,
);
}
@ApiOperation({ summary: 'Get AI question answer' })
@Get('dataset/:datasetId/ai_question/:questionId')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_ASSISTANT.permissions.READ)
async getAiQuestionResult(
@Param('datasetId') datasetId: string,
@Param('questionId') questionId: string,
@User() user: RequestUser,
) {
return this.autodriveAssistantService.proxy(
'GET',
`/dataset/${datasetId}/ai_question/${questionId}`,
user,
);
}
// ============================================
// UTILITY
// ============================================
@ApiOperation({ summary: 'Health check' })
@Get('health')
async health(@User() user: RequestUser) {
return this.autodriveAssistantService.proxy('GET', '/health', user);
}
@ApiOperation({ summary: 'Check model availability' })
@Get('model-availability')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_ASSISTANT.permissions.READ)
async modelAvailability(@User() user: RequestUser) {
return this.autodriveAssistantService.proxy('GET', '/model-availability', user);
}
}
@@ -1,13 +0,0 @@
import { Module } from '@nestjs/common';
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
import { AutodriveAssistantController } from './autodrive-assistant.controller';
import { AutodriveAssistantService } from './autodrive-assistant.service';
@Module({
imports: [],
controllers: [AutodriveAssistantController],
providers: [AutodriveAssistantService, DadosferaLogger],
exports: [AutodriveAssistantService],
})
export class AutodriveAssistantModule {}
@@ -1,94 +0,0 @@
import { Injectable, Inject, HttpException } from '@nestjs/common';
import axios, { AxiosResponse, Method } from 'axios';
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
import { RequestUser } from '../../decorators/user.decorator';
import { AUTODRIVE_ASSISTANT_CONFIG } from './autodrive-assistant.config';
@Injectable()
export class AutodriveAssistantService {
private logger: any;
constructor(
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
) {
this.logger = dadosferaLogger.logger;
}
async proxy(
method: string,
path: string,
user: RequestUser,
body?: any,
query?: Record<string, any>
): Promise<any> {
if (!user.customer_id) {
throw new HttpException('Customer ID is required for autodrive assistant operations', 400);
}
const baseUrl = AUTODRIVE_ASSISTANT_CONFIG.getUrl(user.customer_name);
const url = new URL(`${baseUrl}${path}`);
if (query) {
Object.entries(query).forEach(([key, value]) => {
if (value !== undefined && value !== null) {
url.searchParams.set(key, String(value));
}
});
}
const headers: Record<string, string> = {
'content-type': 'application/json',
'Authorization': user.access_token,
};
this.logger.info('Proxying request to autodrive-assistant', {
method: method.toUpperCase(),
path,
customer_id: user.customer_id,
user_id: user.user_id,
});
try {
const response: AxiosResponse = await axios({
method: method as Method,
url: url.href,
headers,
data: body,
timeout: AUTODRIVE_ASSISTANT_CONFIG.timeout,
validateStatus: () => true,
});
if (response.status >= 400) {
throw new HttpException(response.data, response.status);
}
return response.data;
} catch (error) {
this.logger.error('Autodrive Assistant API proxy error', {
error: error.message,
status: error.response?.status,
path,
method: method.toUpperCase(),
});
if (error instanceof HttpException) {
throw error;
}
if (error.response) {
throw new HttpException(error.response.data, error.response.status);
}
if (error.code === 'ECONNREFUSED') {
throw new HttpException('Autodrive Assistant API service unavailable', 503);
}
if (error.code === 'ETIMEDOUT' || error.code === 'ECONNABORTED') {
throw new HttpException('Autodrive Assistant API request timeout', 504);
}
throw new HttpException('Internal server error', 500);
}
}
}
@@ -1,10 +0,0 @@
export const AUTODRIVE_CORE_CONFIG = {
getUrl: (customerName: string): string => {
const urlTemplate = process.env.AUTODRIVE_CORE_API_URL;
if (!urlTemplate) {
throw new Error('AUTODRIVE_CORE_API_URL environment variable is not set');
}
return urlTemplate.replace('{customer}', customerName);
},
timeout: parseInt(process.env.AUTODRIVE_CORE_TIMEOUT || '30000', 10),
};
@@ -1,286 +0,0 @@
import {
Controller,
Get,
Post,
Put,
Delete,
Param,
Body,
Query,
Inject,
UseInterceptors,
UploadedFiles,
} from '@nestjs/common';
import { ApiTags, ApiOperation, ApiConsumes } from '@nestjs/swagger';
import { FilesInterceptor } from '@nestjs/platform-express';
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
import FormData from 'form-data';
import {
Authenticated,
RequireAllPermissions,
} from '../../decorators/authentication.decorator';
import { User, RequestUser } from '../../decorators/user.decorator';
import { AutodriveCoreService } from './autodrive-core.service';
import { PERMISSIONS_GROUPS } from '../../authentication/permissions.enum';
@ApiTags('Autodrive Core')
@Controller('autodrive-core')
@Authenticated()
export class AutodriveCoreController {
private logger: any;
constructor(
private readonly autodriveCoreService: AutodriveCoreService,
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
) {
this.logger = dadosferaLogger.logger;
}
// ============================================
// DATASET MANAGEMENT
// ============================================
@ApiOperation({ summary: 'List datasets' })
@Get('datasets')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_CORE.permissions.READ)
async listDatasets(
@Query('name') name: string,
@Query('limit') limit: number,
@Query('offset') offset: number,
@User() user: RequestUser,
) {
return this.autodriveCoreService.proxy(
'GET',
'/datasets',
user,
undefined,
{ name, limit, offset },
);
}
@ApiOperation({ summary: 'Get dataset by ID' })
@Get('dataset/:datasetId')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_CORE.permissions.READ)
async getDataset(
@Param('datasetId') datasetId: string,
@User() user: RequestUser,
) {
return this.autodriveCoreService.proxy(
'GET',
`/dataset/${datasetId}`,
user,
);
}
@ApiOperation({ summary: 'Add documents to existing dataset' })
@Put('dataset/:datasetId')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_CORE.permissions.WRITE)
async updateDataset(
@Param('datasetId') datasetId: string,
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveCoreService.proxy(
'PUT',
`/dataset/${datasetId}`,
user,
body,
);
}
@ApiOperation({ summary: 'Delete dataset' })
@Delete('dataset/:datasetId')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_CORE.permissions.WRITE)
async deleteDataset(
@Param('datasetId') datasetId: string,
@User() user: RequestUser,
) {
return this.autodriveCoreService.proxy(
'DELETE',
`/dataset/${datasetId}`,
user,
);
}
@ApiOperation({ summary: 'Append documents from URLs to dataset' })
@Post('dataset/:datasetId/append-from-url')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_CORE.permissions.WRITE)
async appendFromUrl(
@Param('datasetId') datasetId: string,
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveCoreService.proxy(
'POST',
`/dataset/${datasetId}/append-from-url`,
user,
body,
);
}
@ApiOperation({ summary: 'Get dataset creation/update logs' })
@Get('dataset/:datasetId/logs')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_CORE.permissions.READ)
async getDatasetLogs(
@Param('datasetId') datasetId: string,
@User() user: RequestUser,
) {
return this.autodriveCoreService.proxy(
'GET',
`/dataset/${datasetId}/logs`,
user,
);
}
// ============================================
// FILE OPERATIONS
// ============================================
@ApiOperation({ summary: 'Create vector dataset from uploaded files' })
@Post('upload')
@ApiConsumes('multipart/form-data')
@UseInterceptors(FilesInterceptor('files'))
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_CORE.permissions.WRITE)
async upload(
@UploadedFiles() files: Array<Express.Multer.File>,
@Body() body: any,
@User() user: RequestUser,
) {
const formData = new FormData();
if (files && files.length > 0) {
files.forEach((file) => {
formData.append('files', file.buffer, {
filename: file.originalname,
contentType: file.mimetype,
});
});
}
// Forward additional body fields
if (body) {
Object.entries(body).forEach(([key, value]) => {
if (value !== undefined && value !== null && key !== 'files') {
formData.append(key, String(value));
}
});
}
return this.autodriveCoreService.proxyFormData(
'POST',
'/upload',
user,
formData,
);
}
@ApiOperation({ summary: 'Create vector dataset from URLs' })
@Post('upload-from-url')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_CORE.permissions.WRITE)
async uploadFromUrl(
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveCoreService.proxy(
'POST',
'/upload-from-url',
user,
body,
);
}
// ============================================
// QUERY OPERATIONS
// ============================================
@ApiOperation({ summary: 'Search dataset using semantic search' })
@Post('dataset/:datasetId/question')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_CORE.permissions.READ)
async question(
@Param('datasetId') datasetId: string,
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveCoreService.proxy(
'POST',
`/dataset/${datasetId}/question`,
user,
body,
);
}
@ApiOperation({ summary: 'Get semantic search results' })
@Get('dataset/:datasetId/question/:questionId')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_CORE.permissions.READ)
async getQuestionResult(
@Param('datasetId') datasetId: string,
@Param('questionId') questionId: string,
@User() user: RequestUser,
) {
return this.autodriveCoreService.proxy(
'GET',
`/dataset/${datasetId}/question/${questionId}`,
user,
);
}
@ApiOperation({ summary: 'Ask AI question on dataset' })
@Post('dataset/:datasetId/ai_question')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_CORE.permissions.READ)
async aiQuestion(
@Param('datasetId') datasetId: string,
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveCoreService.proxy(
'POST',
`/dataset/${datasetId}/ai_question`,
user,
body,
);
}
@ApiOperation({ summary: 'Get AI question answer' })
@Get('dataset/:datasetId/ai_question/:questionId')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_CORE.permissions.READ)
async getAiQuestionResult(
@Param('datasetId') datasetId: string,
@Param('questionId') questionId: string,
@User() user: RequestUser,
) {
return this.autodriveCoreService.proxy(
'GET',
`/dataset/${datasetId}/ai_question/${questionId}`,
user,
);
}
// ============================================
// SYSTEM
// ============================================
@ApiOperation({ summary: 'Health check' })
@Get('health')
async health(@User() user: RequestUser) {
return this.autodriveCoreService.proxy('GET', '/health', user);
}
@ApiOperation({ summary: 'Check model availability' })
@Get('model-availability')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_CORE.permissions.READ)
async modelAvailability(@User() user: RequestUser) {
return this.autodriveCoreService.proxy('GET', '/model-availability', user);
}
@ApiOperation({ summary: 'Get frontend configuration' })
@Get('config')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_CORE.permissions.READ)
async getConfig(@User() user: RequestUser) {
return this.autodriveCoreService.proxy('GET', '/config', user);
}
@ApiOperation({ summary: 'Get usage metrics' })
@Get('usage')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_CORE.permissions.READ)
async getUsage(@User() user: RequestUser) {
return this.autodriveCoreService.proxy('GET', '/usage', user);
}
}
@@ -1,13 +0,0 @@
import { Module } from '@nestjs/common';
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
import { AutodriveCoreController } from './autodrive-core.controller';
import { AutodriveCoreService } from './autodrive-core.service';
@Module({
imports: [],
controllers: [AutodriveCoreController],
providers: [AutodriveCoreService, DadosferaLogger],
exports: [AutodriveCoreService],
})
export class AutodriveCoreModule {}
@@ -1,165 +0,0 @@
import { Injectable, Inject, HttpException } from '@nestjs/common';
import axios, { AxiosResponse, Method } from 'axios';
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
import { RequestUser } from '../../decorators/user.decorator';
import { AUTODRIVE_CORE_CONFIG } from './autodrive-core.config';
@Injectable()
export class AutodriveCoreService {
private logger: any;
constructor(
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
) {
this.logger = dadosferaLogger.logger;
}
async proxy(
method: string,
path: string,
user: RequestUser,
body?: any,
query?: Record<string, any>
): Promise<any> {
if (!user.customer_id) {
throw new HttpException('Customer ID is required for autodrive core operations', 400);
}
const baseUrl = AUTODRIVE_CORE_CONFIG.getUrl(user.customer_name);
const url = new URL(`${baseUrl}${path}`);
if (query) {
Object.entries(query).forEach(([key, value]) => {
if (value !== undefined && value !== null) {
url.searchParams.set(key, String(value));
}
});
}
const headers: Record<string, string> = {
'content-type': 'application/json',
'Authorization': user.access_token,
};
this.logger.info('Proxying request to autodrive-core', {
method: method.toUpperCase(),
path,
customer_id: user.customer_id,
user_id: user.user_id,
});
try {
const response: AxiosResponse = await axios({
method: method as Method,
url: url.href,
headers,
data: body,
timeout: AUTODRIVE_CORE_CONFIG.timeout,
validateStatus: () => true,
});
if (response.status >= 400) {
throw new HttpException(response.data, response.status);
}
return response.data;
} catch (error) {
this.logger.error('Autodrive Core API proxy error', {
error: error.message,
status: error.response?.status,
path,
method: method.toUpperCase(),
});
if (error instanceof HttpException) {
throw error;
}
if (error.response) {
throw new HttpException(error.response.data, error.response.status);
}
if (error.code === 'ECONNREFUSED') {
throw new HttpException('Autodrive Core API service unavailable', 503);
}
if (error.code === 'ETIMEDOUT' || error.code === 'ECONNABORTED') {
throw new HttpException('Autodrive Core API request timeout', 504);
}
throw new HttpException('Internal server error', 500);
}
}
async proxyFormData(
method: string,
path: string,
user: RequestUser,
formData: any,
query?: Record<string, any>,
): Promise<any> {
if (!user.customer_id) {
throw new HttpException('Customer ID is required for autodrive core operations', 400);
}
const baseUrl = AUTODRIVE_CORE_CONFIG.getUrl(user.customer_name);
const url = new URL(`${baseUrl}${path}`);
if (query) {
Object.entries(query).forEach(([key, value]) => {
if (value !== undefined && value !== null) {
url.searchParams.set(key, String(value));
}
});
}
const headers: Record<string, string> = {
...formData.getHeaders?.(),
'Authorization': user.access_token,
};
this.logger.info('Proxying form data request to autodrive-core', {
method: method.toUpperCase(),
path,
customer_id: user.customer_id,
user_id: user.user_id,
});
try {
const response: AxiosResponse = await axios({
method: method as Method,
url: url.href,
headers,
data: formData,
timeout: AUTODRIVE_CORE_CONFIG.timeout,
maxContentLength: Infinity,
maxBodyLength: Infinity,
validateStatus: () => true,
});
if (response.status >= 400) {
throw new HttpException(response.data, response.status);
}
return response.data;
} catch (error) {
this.logger.error('Autodrive Core API form data proxy error', {
error: error.message,
status: error.response?.status,
path,
method: method.toUpperCase(),
});
if (error instanceof HttpException) {
throw error;
}
if (error.response) {
throw new HttpException(error.response.data, error.response.status);
}
throw new HttpException('Internal server error', 500);
}
}
}
@@ -1,10 +0,0 @@
export const AUTODRIVE_EXTRACTOR_CONFIG = {
getUrl: (customerName: string): string => {
const urlTemplate = process.env.AUTODRIVE_EXTRACTOR_API_URL;
if (!urlTemplate) {
throw new Error('AUTODRIVE_EXTRACTOR_API_URL environment variable is not set');
}
return urlTemplate.replace('{customer}', customerName);
},
timeout: parseInt(process.env.AUTODRIVE_EXTRACTOR_TIMEOUT || '30000', 10),
};
@@ -1,491 +0,0 @@
import {
Controller,
Get,
Post,
Put,
Delete,
Param,
Body,
Query,
Inject,
UseInterceptors,
UploadedFiles,
} from '@nestjs/common';
import { ApiTags, ApiOperation, ApiConsumes } from '@nestjs/swagger';
import { FilesInterceptor } from '@nestjs/platform-express';
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
import FormData from 'form-data';
import {
Authenticated,
RequireAllPermissions,
} from '../../decorators/authentication.decorator';
import { User, RequestUser } from '../../decorators/user.decorator';
import { AutodriveExtractorService } from './autodrive-extractor.service';
import { PERMISSIONS_GROUPS } from '../../authentication/permissions.enum';
@ApiTags('Autodrive Extractor')
@Controller('autodrive-extractor')
@Authenticated()
export class AutodriveExtractorController {
private logger: any;
constructor(
private readonly autodriveExtractorService: AutodriveExtractorService,
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
) {
this.logger = dadosferaLogger.logger;
}
// ============================================
// LEGACY - NPL Extraction routes (to be removed)
// ============================================
@ApiOperation({ summary: '[LEGACY] Extract features from dataset using NPL' })
@Post('npl/datasets/:datasetId/extract')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.WRITE)
async nplExtract(
@Param('datasetId') datasetId: string,
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'POST',
`/npl/datasets/${datasetId}/extract`,
user,
body,
);
}
@ApiOperation({ summary: '[LEGACY] List NPL extractions for a dataset' })
@Get('npl/datasets/:datasetId/extractions')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async nplListExtractions(
@Param('datasetId') datasetId: string,
@Query('status') status: string,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'GET',
`/npl/datasets/${datasetId}/extractions`,
user,
undefined,
{ status },
);
}
@ApiOperation({ summary: '[LEGACY] Get NPL extraction details' })
@Get('npl/datasets/:datasetId/extractions/:extractionId')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async nplGetExtraction(
@Param('datasetId') datasetId: string,
@Param('extractionId') extractionId: string,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'GET',
`/npl/datasets/${datasetId}/extractions/${extractionId}`,
user,
);
}
@ApiOperation({ summary: '[LEGACY] Export NPL extraction to CSV' })
@Get('npl/datasets/:datasetId/extractions/:extractionId/csv')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async nplExportExtractionCsv(
@Param('datasetId') datasetId: string,
@Param('extractionId') extractionId: string,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'GET',
`/npl/datasets/${datasetId}/extractions/${extractionId}/csv`,
user,
);
}
@ApiOperation({ summary: '[LEGACY] Export NPL extraction to XLSX' })
@Get('npl/datasets/:datasetId/extractions/:extractionId/xlsx')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async nplExportExtractionXlsx(
@Param('datasetId') datasetId: string,
@Param('extractionId') extractionId: string,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'GET',
`/npl/datasets/${datasetId}/extractions/${extractionId}/xlsx`,
user,
);
}
@ApiOperation({ summary: '[LEGACY] Batch export all NPL extractions to CSV' })
@Get('npl/datasets/:datasetId/extractions/export/csv')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async nplBatchExportCsv(
@Param('datasetId') datasetId: string,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'GET',
`/npl/datasets/${datasetId}/extractions/export/csv`,
user,
);
}
@ApiOperation({ summary: '[LEGACY] Batch export all NPL extractions to XLSX' })
@Get('npl/datasets/:datasetId/extractions/export/xlsx')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async nplBatchExportXlsx(
@Param('datasetId') datasetId: string,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'GET',
`/npl/datasets/${datasetId}/extractions/export/xlsx`,
user,
);
}
@ApiOperation({ summary: '[LEGACY] Get NPL extraction report' })
@Get('npl/datasets/:datasetId/extractions/:extractionId/report')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async nplGetExtractionReport(
@Param('datasetId') datasetId: string,
@Param('extractionId') extractionId: string,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'GET',
`/npl/datasets/${datasetId}/extractions/${extractionId}/report`,
user,
);
}
// ============================================
// LEGACY - NPL Configuration routes (to be removed)
// ============================================
@ApiOperation({ summary: '[LEGACY] Get NPL configuration' })
@Get('npl/config')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async nplGetConfig(@User() user: RequestUser) {
return this.autodriveExtractorService.proxy('GET', '/npl/config', user);
}
@ApiOperation({ summary: '[LEGACY] Get NPL configuration summary' })
@Get('npl/config/summary')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async nplGetConfigSummary(@User() user: RequestUser) {
return this.autodriveExtractorService.proxy('GET', '/npl/config/summary', user);
}
@ApiOperation({ summary: '[LEGACY] List NPL configuration fields' })
@Get('npl/config/fields')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async nplGetConfigFields(@User() user: RequestUser) {
return this.autodriveExtractorService.proxy('GET', '/npl/config/fields', user);
}
@ApiOperation({ summary: '[LEGACY] Get NPL configuration field by ID' })
@Get('npl/config/fields/:fieldId')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async nplGetConfigField(
@Param('fieldId') fieldId: string,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'GET',
`/npl/config/fields/${fieldId}`,
user,
);
}
@ApiOperation({ summary: '[LEGACY] Update NPL configuration field' })
@Put('npl/config/fields/:fieldId')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.WRITE)
async nplUpdateConfigField(
@Param('fieldId') fieldId: string,
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'PUT',
`/npl/config/fields/${fieldId}`,
user,
body,
);
}
@ApiOperation({ summary: '[LEGACY] Get NPL RAG queries' })
@Get('npl/config/rag-queries')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async nplGetRagQueries(@User() user: RequestUser) {
return this.autodriveExtractorService.proxy('GET', '/npl/config/rag-queries', user);
}
@ApiOperation({ summary: '[LEGACY] Get NPL system prompt' })
@Get('npl/config/system-prompt')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async nplGetSystemPrompt(@User() user: RequestUser) {
return this.autodriveExtractorService.proxy('GET', '/npl/config/system-prompt', user);
}
@ApiOperation({ summary: '[LEGACY] Update NPL system prompt' })
@Put('npl/config/system-prompt')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.WRITE)
async nplUpdateSystemPrompt(
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'PUT',
'/npl/config/system-prompt',
user,
body,
);
}
@ApiOperation({ summary: '[LEGACY] Get NPL source priorities' })
@Get('npl/config/source-priorities')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async nplGetSourcePriorities(@User() user: RequestUser) {
return this.autodriveExtractorService.proxy('GET', '/npl/config/source-priorities', user);
}
@ApiOperation({ summary: '[LEGACY] Reset NPL configuration to defaults' })
@Post('npl/config/reset')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.WRITE)
async nplResetConfig(@User() user: RequestUser) {
return this.autodriveExtractorService.proxy('POST', '/npl/config/reset', user);
}
@ApiOperation({ summary: '[LEGACY] Export NPL configuration' })
@Get('npl/config/export')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async nplExportConfig(@User() user: RequestUser) {
return this.autodriveExtractorService.proxy('GET', '/npl/config/export', user);
}
@ApiOperation({ summary: '[LEGACY] Import NPL configuration' })
@Post('npl/config/import')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.WRITE)
async nplImportConfig(
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'POST',
'/npl/config/import',
user,
body,
);
}
// ============================================
// EXTRACTION TEMPLATES
// ============================================
@ApiOperation({ summary: 'Seed NPL Brasil extraction template' })
@Post('templates/seed-npl')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.WRITE)
async seedNplTemplate(
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'POST',
'/templates/seed-npl',
user,
body,
);
}
@ApiOperation({ summary: 'List extraction templates' })
@Get('templates')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async listTemplates(@User() user: RequestUser) {
return this.autodriveExtractorService.proxy('GET', '/templates', user);
}
@ApiOperation({ summary: 'Create extraction template' })
@Post('templates')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.WRITE)
async createTemplate(
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'POST',
'/templates',
user,
body,
);
}
@ApiOperation({ summary: 'Get extraction template by ID' })
@Get('templates/:templateId')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async getTemplate(
@Param('templateId') templateId: string,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'GET',
`/templates/${templateId}`,
user,
);
}
@ApiOperation({ summary: 'Update extraction template' })
@Put('templates/:templateId')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.WRITE)
async updateTemplate(
@Param('templateId') templateId: string,
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'PUT',
`/templates/${templateId}`,
user,
body,
);
}
@ApiOperation({ summary: 'Delete extraction template' })
@Delete('templates/:templateId')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.WRITE)
async deleteTemplate(
@Param('templateId') templateId: string,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'DELETE',
`/templates/${templateId}`,
user,
);
}
@ApiOperation({ summary: 'Import extraction templates from file' })
@Post('templates/import')
@ApiConsumes('multipart/form-data')
@UseInterceptors(FilesInterceptor('file'))
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.WRITE)
async importTemplates(
@UploadedFiles() files: Array<Express.Multer.File>,
@User() user: RequestUser,
) {
const formData = new FormData();
if (files && files.length > 0) {
files.forEach((file) => {
formData.append('file', file.buffer, {
filename: file.originalname,
contentType: file.mimetype,
});
});
}
return this.autodriveExtractorService.proxyFormData(
'POST',
'/templates/import',
user,
formData,
);
}
// ============================================
// GENERIC EXTRACTION
// ============================================
@ApiOperation({ summary: 'Start template-based extraction on dataset' })
@Post('extraction/datasets/:datasetId/extract')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.WRITE)
async startExtraction(
@Param('datasetId') datasetId: string,
@Body() body: any,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'POST',
`/extraction/datasets/${datasetId}/extract`,
user,
body,
);
}
@ApiOperation({ summary: 'List extraction jobs for a dataset' })
@Get('extraction/datasets/:datasetId/jobs')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async listExtractionJobs(
@Param('datasetId') datasetId: string,
@Query('status') status: string,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'GET',
`/extraction/datasets/${datasetId}/jobs`,
user,
undefined,
{ status },
);
}
@ApiOperation({ summary: 'Get extraction job details' })
@Get('extraction/jobs/:jobId')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async getExtractionJob(
@Param('jobId') jobId: string,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'GET',
`/extraction/jobs/${jobId}`,
user,
);
}
@ApiOperation({ summary: 'Get extraction job report' })
@Get('extraction/jobs/:jobId/report')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async getExtractionJobReport(
@Param('jobId') jobId: string,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'GET',
`/extraction/jobs/${jobId}/report`,
user,
);
}
@ApiOperation({ summary: 'Export extraction job to CSV' })
@Get('extraction/jobs/:jobId/csv')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async exportExtractionJobCsv(
@Param('jobId') jobId: string,
@User() user: RequestUser,
) {
return this.autodriveExtractorService.proxy(
'GET',
`/extraction/jobs/${jobId}/csv`,
user,
);
}
// ============================================
// UTILITY
// ============================================
@ApiOperation({ summary: 'Health check' })
@Get('health')
async health(@User() user: RequestUser) {
return this.autodriveExtractorService.proxy('GET', '/health', user);
}
@ApiOperation({ summary: 'Check model availability' })
@Get('model-availability')
@RequireAllPermissions(PERMISSIONS_GROUPS.AUTODRIVE_EXTRACTOR.permissions.READ)
async modelAvailability(@User() user: RequestUser) {
return this.autodriveExtractorService.proxy('GET', '/model-availability', user);
}
}
@@ -1,13 +0,0 @@
import { Module } from '@nestjs/common';
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
import { AutodriveExtractorController } from './autodrive-extractor.controller';
import { AutodriveExtractorService } from './autodrive-extractor.service';
@Module({
imports: [],
controllers: [AutodriveExtractorController],
providers: [AutodriveExtractorService, DadosferaLogger],
exports: [AutodriveExtractorService],
})
export class AutodriveExtractorModule {}
@@ -1,165 +0,0 @@
import { Injectable, Inject, HttpException } from '@nestjs/common';
import axios, { AxiosResponse, Method } from 'axios';
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
import { RequestUser } from '../../decorators/user.decorator';
import { AUTODRIVE_EXTRACTOR_CONFIG } from './autodrive-extractor.config';
@Injectable()
export class AutodriveExtractorService {
private logger: any;
constructor(
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
) {
this.logger = dadosferaLogger.logger;
}
async proxy(
method: string,
path: string,
user: RequestUser,
body?: any,
query?: Record<string, any>
): Promise<any> {
if (!user.customer_id) {
throw new HttpException('Customer ID is required for autodrive extractor operations', 400);
}
const baseUrl = AUTODRIVE_EXTRACTOR_CONFIG.getUrl(user.customer_name);
const url = new URL(`${baseUrl}${path}`);
if (query) {
Object.entries(query).forEach(([key, value]) => {
if (value !== undefined && value !== null) {
url.searchParams.set(key, String(value));
}
});
}
const headers: Record<string, string> = {
'content-type': 'application/json',
'Authorization': user.access_token,
};
this.logger.info('Proxying request to autodrive-extractor', {
method: method.toUpperCase(),
path,
customer_id: user.customer_id,
user_id: user.user_id,
});
try {
const response: AxiosResponse = await axios({
method: method as Method,
url: url.href,
headers,
data: body,
timeout: AUTODRIVE_EXTRACTOR_CONFIG.timeout,
validateStatus: () => true,
});
if (response.status >= 400) {
throw new HttpException(response.data, response.status);
}
return response.data;
} catch (error) {
this.logger.error('Autodrive Extractor API proxy error', {
error: error.message,
status: error.response?.status,
path,
method: method.toUpperCase(),
});
if (error instanceof HttpException) {
throw error;
}
if (error.response) {
throw new HttpException(error.response.data, error.response.status);
}
if (error.code === 'ECONNREFUSED') {
throw new HttpException('Autodrive Extractor API service unavailable', 503);
}
if (error.code === 'ETIMEDOUT' || error.code === 'ECONNABORTED') {
throw new HttpException('Autodrive Extractor API request timeout', 504);
}
throw new HttpException('Internal server error', 500);
}
}
async proxyFormData(
method: string,
path: string,
user: RequestUser,
formData: any,
query?: Record<string, any>,
): Promise<any> {
if (!user.customer_id) {
throw new HttpException('Customer ID is required for autodrive extractor operations', 400);
}
const baseUrl = AUTODRIVE_EXTRACTOR_CONFIG.getUrl(user.customer_name);
const url = new URL(`${baseUrl}${path}`);
if (query) {
Object.entries(query).forEach(([key, value]) => {
if (value !== undefined && value !== null) {
url.searchParams.set(key, String(value));
}
});
}
const headers: Record<string, string> = {
...formData.getHeaders?.(),
'Authorization': user.access_token,
};
this.logger.info('Proxying form data request to autodrive-extractor', {
method: method.toUpperCase(),
path,
customer_id: user.customer_id,
user_id: user.user_id,
});
try {
const response: AxiosResponse = await axios({
method: method as Method,
url: url.href,
headers,
data: formData,
timeout: AUTODRIVE_EXTRACTOR_CONFIG.timeout,
maxContentLength: Infinity,
maxBodyLength: Infinity,
validateStatus: () => true,
});
if (response.status >= 400) {
throw new HttpException(response.data, response.status);
}
return response.data;
} catch (error) {
this.logger.error('Autodrive Extractor API form data proxy error', {
error: error.message,
status: error.response?.status,
path,
method: method.toUpperCase(),
});
if (error instanceof HttpException) {
throw error;
}
if (error.response) {
throw new HttpException(error.response.data, error.response.status);
}
throw new HttpException('Internal server error', 500);
}
}
}
+58
View File
@@ -734,6 +734,64 @@ class CatalogService implements OnModuleInit {
}
}
async renameTableOnNimbus(
nimbusUrl: string,
nimbusId: number,
changes: { table_name?: string; table_schema?: string; display_name?: string },
): Promise<void> {
const endpoint = `${nimbusUrl}/api/catalog/table-metadata/${nimbusId}`;
this.logger.info(`Renaming table-metadata ${nimbusId} on Nimbus`, { endpoint, changes });
await axios.patch(endpoint, changes);
}
async renameColumnMetadataOnNimbus(
nimbusUrl: string,
databaseName: string,
oldTableName: string,
oldTableSchema: string,
newTableName: string,
newTableSchema: string,
): Promise<void> {
const listEndpoint = `${nimbusUrl}/api/catalog/column-metadata/?database_name=${encodeURIComponent(databaseName)}&table_name=${encodeURIComponent(oldTableName)}&table_schema=${encodeURIComponent(oldTableSchema)}`;
this.logger.info(`Fetching column-metadata records to rename`, { listEndpoint });
const { data: columns } = await axios.get(listEndpoint);
const filtered = Array.isArray(columns) ? columns : [];
for (const column of filtered) {
const patchEndpoint = `${nimbusUrl}/api/catalog/column-metadata/${column.id}`;
await axios.patch(patchEndpoint, {
table_name: newTableName,
table_schema: newTableSchema,
});
}
this.logger.info(`Renamed ${filtered.length} column-metadata records on Nimbus`);
}
async renameDataPreviewOnNimbus(
nimbusUrl: string,
databaseName: string,
oldTableName: string,
oldTableSchema: string,
newTableName: string,
newTableSchema: string,
): Promise<void> {
const listEndpoint = `${nimbusUrl}/api/catalog/data-preview/?database_name=${encodeURIComponent(databaseName)}&table_name=${encodeURIComponent(oldTableName)}&table_schema=${encodeURIComponent(oldTableSchema)}`;
this.logger.info(`Fetching data-preview records to rename`, { listEndpoint });
const { data: previews } = await axios.get(listEndpoint);
const filtered = Array.isArray(previews) ? previews : [];
for (const preview of filtered) {
const patchEndpoint = `${nimbusUrl}/api/catalog/data-preview/${preview.id}`;
await axios.patch(patchEndpoint, {
table_name: newTableName,
table_schema: newTableSchema,
});
}
this.logger.info(`Renamed ${filtered.length} data-preview records on Nimbus`);
}
async catalogDatasetItem(table_metadata_id: number, metadata: Metadata) {
const customer_name_raw = metadata.get('customer_name');
@@ -10,6 +10,7 @@ import {
Query,
Inject,
BadRequestException,
HttpException,
} from '@nestjs/common';
import { ApiTags, ApiOperation } from '@nestjs/swagger';
import { DadosferaLogger } from '@dadosfera/dadosfera-logs';
@@ -25,6 +26,8 @@ import { ElasticsearchService } from '../../services/elasticsearch';
import { DynamoDBService, ReferenceColumn } from '../../services/dynamodb';
import { CustomersService } from '../customers/customers.service';
import { validateCronAgainstScheduleLimit } from '../../utils/cron-validation';
import { CatalogService } from '../catalog/catalog.service';
import { PackTheMetadata } from '../../utils/PackTheMetadata';
import { ValidationTableDTO } from './platform-api.dto';
@@ -35,6 +38,11 @@ type ValidateTablesDTO = {
}>
}
type RenameTablesBody = {
raw?: { table_name: string; table_schema: string };
qualify?: { table_name: string; table_schema: string };
}
@ApiTags('Platform API')
@Controller('platform')
export class PlatformApiController {
@@ -45,6 +53,7 @@ export class PlatformApiController {
private readonly elasticsearchService: ElasticsearchService,
private readonly dynamoDBService: DynamoDBService,
private readonly customersService: CustomersService,
private readonly catalogService: CatalogService,
@Inject(DadosferaLogger) dadosferaLogger: DadosferaLogger,
) {
this.logger = dadosferaLogger.logger;
@@ -76,6 +85,23 @@ export class PlatformApiController {
return jobId?.replace(/-/g, '_') || '';
}
private async getJobByAnyConnectorType(normalizedJobId: string, user: RequestUser): Promise<any> {
const connectorTypes = ['jdbc', 'singer', 's3'];
for (const type of connectorTypes) {
try {
const job = await this.platformApiService.proxy(
'GET',
`/jobs/${type}/${normalizedJobId}`,
user,
);
return job;
} catch (error) {
// Continue to next connector type
}
}
throw new HttpException(`Job ${normalizedJobId} not found in any connector type (jdbc, singer, s3)`, 404);
}
/**
* Extract the pipeline ID (base UUID) from a job ID.
* Job IDs have format "uuid-suffix" where suffix is the job index (e.g., "0", "1").
@@ -953,6 +979,187 @@ export class PlatformApiController {
return result;
}
@Post('jobs/:jobId/rename-tables')
@ApiOperation({ summary: 'Rename job output tables and sync to catalog' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
async renameJobTables(
@Param('jobId') jobId: string,
@Body() body: RenameTablesBody,
@User() user: RequestUser,
) {
const normalizedJobId = this.normalizeJobId(jobId);
const currentJob = await this.getJobByAnyConnectorType(normalizedJobId, user);
const result = await this.platformApiService.proxy(
'POST',
`/jobs/${normalizedJobId}/rename-tables`,
user,
body,
);
try {
await this.syncTableRenameToCatalog(jobId, body, currentJob, user);
} catch (error) {
this.logger.error('Catalog sync failed, rolling back Snowflake rename', { jobId, error: error.message });
const reverseBody = this.buildSnowflakeRollbackBody(body, currentJob.output_config || {});
if (reverseBody) {
try {
await this.platformApiService.proxy('POST', `/jobs/${normalizedJobId}/rename-tables`, user, reverseBody);
this.logger.info('Snowflake rename rolled back', { jobId });
} catch (rollbackError) {
this.logger.error('Snowflake rollback failed', { jobId, error: rollbackError.message });
}
}
throw new HttpException('Table rename failed: catalog sync error, Snowflake reverted', 500);
}
return result;
}
private buildSnowflakeRollbackBody(
body: RenameTablesBody,
outputConfig: any,
): RenameTablesBody | null {
const reverse: RenameTablesBody = {};
if (body.raw) {
const nested = outputConfig.raw;
const oldTableName = nested?.table_name || outputConfig.table_name;
const oldTableSchema = nested?.table_schema || 'PUBLIC';
if (oldTableName) reverse.raw = { table_name: oldTableName, table_schema: oldTableSchema };
}
if (body.qualify) {
const nested = outputConfig.qualify;
if (nested?.table_name) reverse.qualify = { table_name: nested.table_name, table_schema: nested.table_schema || 'STAGED' };
}
return Object.keys(reverse).length > 0 ? reverse : null;
}
/**
* Sync table rename to Elasticsearch and Nimbus.
*
* For each target (raw, qualify):
* 1. Resolve old table name from output_config
* 2. Find the ES data asset by pipeline + table + schema
* 3. Update ES, Nimbus table-metadata, column-metadata, and data-preview
* 4. If any step fails, rollback all completed steps for that target
*/
private async syncTableRenameToCatalog(
jobId: string,
body: RenameTablesBody,
currentJob: any,
user: RequestUser,
): Promise<void> {
const pipelineId = this.extractPipelineIdFromJobId(jobId);
const outputConfig = currentJob.output_config || {};
const nimbusUrl = this.catalogService._getNimbusUrl({ info: { customer: user.customer_name } });
const databaseName = `DADOSFERA_PRD_${user.customer_name.toUpperCase()}`;
const targets = this.buildRenameTargets(body, outputConfig);
for (const { key, oldTableName, oldTableSchema, newValues } of targets) {
const rollbackSteps: Array<() => Promise<void>> = [];
try {
const dataAsset = await this.elasticsearchService.findDataAssetByTable(
user.customer_name, oldTableName, oldTableSchema,
);
if (!dataAsset) {
this.logger.warn(`No data asset found for ${key}`, { jobId, pipelineId, oldTableName, oldTableSchema });
continue;
}
const { _es_id: esAssetId, nimbus_id: nimbusId } = dataAsset;
const oldValues = { table_name: oldTableName, table_schema: oldTableSchema };
// ES update
const esFields = { name: newValues.table_name, table_name: newValues.table_name, table_schema: newValues.table_schema, display_name: newValues.table_name };
await this.elasticsearchService.updateDataAsset(user.customer_name, esAssetId, esFields);
rollbackSteps.push(() => this.elasticsearchService.updateDataAsset(
user.customer_name, esAssetId,
{ name: oldTableName, table_name: oldTableName, table_schema: oldTableSchema, display_name: oldTableName },
));
// Nimbus table-metadata
if (nimbusId) {
await this.catalogService.renameTableOnNimbus(nimbusUrl, nimbusId, newValues);
rollbackSteps.push(() => this.catalogService.renameTableOnNimbus(nimbusUrl, nimbusId, oldValues));
}
// Nimbus column-metadata
await this.catalogService.renameColumnMetadataOnNimbus(
nimbusUrl, databaseName, oldTableName, oldTableSchema, newValues.table_name, newValues.table_schema,
);
rollbackSteps.push(() => this.catalogService.renameColumnMetadataOnNimbus(
nimbusUrl, databaseName, newValues.table_name, newValues.table_schema, oldTableName, oldTableSchema,
));
// Nimbus data-preview
await this.catalogService.renameDataPreviewOnNimbus(
nimbusUrl, databaseName, oldTableName, oldTableSchema, newValues.table_name, newValues.table_schema,
);
rollbackSteps.push(() => this.catalogService.renameDataPreviewOnNimbus(
nimbusUrl, databaseName, newValues.table_name, newValues.table_schema, oldTableName, oldTableSchema,
));
this.logger.info(`Synced catalog rename for ${key}`, { jobId, oldTableName, newTableName: newValues.table_name });
} catch (error) {
this.logger.error(`Catalog sync failed for ${key}, rolling back catalog`, { jobId, error: error.message });
await this.executeRollback(rollbackSteps, key, jobId);
throw error;
}
}
}
private buildRenameTargets(
body: RenameTablesBody,
outputConfig: any,
): Array<{ key: string; oldTableName: string; oldTableSchema: string; newValues: { table_name: string; table_schema: string } }> {
const DEFAULT_SCHEMAS = { raw: 'PUBLIC', qualify: 'STAGED' };
const targets: Array<{ key: string; oldTableName: string; oldTableSchema: string; newValues: { table_name: string; table_schema: string } }> = [];
for (const key of ['raw', 'qualify'] as const) {
if (!body[key]) continue;
const nested = outputConfig[key];
// qualify: only sync if output_config.qualify already exists
if (key === 'qualify' && !nested?.table_name) continue;
const oldTableName = nested?.table_name || outputConfig.table_name;
if (!oldTableName) continue;
targets.push({
key,
oldTableName,
oldTableSchema: nested?.table_schema || DEFAULT_SCHEMAS[key],
newValues: body[key],
});
}
return targets;
}
private async executeRollback(
steps: Array<() => Promise<void>>,
targetKey: string,
jobId: string,
): Promise<void> {
for (const rollback of steps.reverse()) {
try {
await rollback();
} catch (error) {
this.logger.error(`Rollback failed for ${targetKey}`, { jobId, error: error.message });
}
}
}
@Get('jobs/jdbc/configs/allowed_datatypes')
@ApiOperation({ summary: 'Get allowed datatypes for JDBC' })
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
@@ -7,9 +7,10 @@ import { PlatformApiService } from './platform-api.service';
import { ElasticsearchModule } from '../../services/elasticsearch';
import { DynamoDBModule } from '../../services/dynamodb';
import { CustomersModule } from '../customers/customers.module';
import { CatalogModule } from '../catalog/catalog.module';
@Module({
imports: [ElasticsearchModule, DynamoDBModule, CustomersModule],
imports: [ElasticsearchModule, DynamoDBModule, CustomersModule, CatalogModule],
controllers: [PlatformApiController],
providers: [PlatformApiService, DadosferaLogger],
exports: [PlatformApiService],
@@ -89,6 +89,12 @@ export class PlatformApiService {
// Propagate non-2xx responses as HttpExceptions
if (response.status >= 400) {
this.logger.error('Platform API upstream error', {
status: response.status,
data: response.data,
path,
method: method.toUpperCase(),
});
throw new HttpException(response.data, response.status);
}
@@ -358,6 +358,81 @@ export class ElasticsearchService {
}
}
private getDataAssetIndex(customerName: string): string {
return `${customerName}_data_assets_catalog`;
}
async findDataAssetByTable(
customerName: string,
tableName: string,
tableSchema: string,
): Promise<{ id: string; nimbus_id: number | null; [key: string]: any } | null> {
const index = this.getDataAssetIndex(customerName);
this.logger.info('Elasticsearch: Searching data asset', {
index,
tableName,
tableSchema,
});
try {
const response = await this.client.post(`/${index}/_search`, {
query: {
bool: {
must: [
{ term: { 'table_name.keyword': tableName.toUpperCase() } },
{ term: { 'table_schema.keyword': tableSchema.toUpperCase() } },
],
},
},
size: 1,
});
const hits = response.data.hits?.hits || [];
if (hits.length === 0) {
this.logger.warn('Elasticsearch: Data asset not found', { tableName, tableSchema, index });
return null;
}
return { ...hits[0]._source, _es_id: hits[0]._id };
} catch (error) {
this.handleError('findDataAssetByTable', error, { tableName, tableSchema, index });
throw error;
}
}
async updateDataAsset(
customerName: string,
assetId: string,
updates: Record<string, any>,
): Promise<any> {
const index = this.getDataAssetIndex(customerName);
this.logger.info('Elasticsearch: Updating data asset', {
index,
assetId,
fields: Object.keys(updates),
});
try {
const response = await this.client.post(
`/${index}/_update/${assetId}`,
{ doc: updates },
{ params: { refresh: 'wait_for' } },
);
this.logger.info('Elasticsearch: Data asset updated', {
assetId,
result: response.data.result,
});
return response.data;
} catch (error) {
this.handleError('updateDataAsset', error, { assetId, index });
throw error;
}
}
private handleError(
operation: string,
error: any,