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>
This commit is contained in:
Rafael
2026-03-05 17:29:12 -03:00
co-authored by Claude Opus 4.6
parent 1a2153d62f
commit 38a9e21f5f
5 changed files with 351 additions and 1 deletions
@@ -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';
type ValidateTablesDTO = {
@@ -34,6 +37,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 {
@@ -44,6 +52,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;
@@ -75,6 +84,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").
@@ -968,6 +994,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.findDataAssetByPipelineAndTable(
user.customer_name, pipelineId, 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)