mirror of
https://github.com/dadosfera/maestro.git
synced 2026-10-08 07:29:09 +00:00
FEAT: add schedule limit validation and improve ES update
- Add schedule limit validation against customer's scheduleLimit from DUC - Improve ES update to only update provided fields - Mark GET /platform/pipeline/:pipelineId/pipeline_run as READY 🤖 Generated with [Claude Code](https://claude.com/claude-code) Co-Authored-By: Claude Opus 4.5 <noreply@anthropic.com>
This commit is contained in:
co-authored by
Claude Opus 4.5
parent
f7efb757bf
commit
01c1087e07
@@ -49,7 +49,52 @@ export class PlatformApiController {
|
||||
return id?.replace(/-/g, '_') || '';
|
||||
}
|
||||
|
||||
/**
|
||||
* Denormalize ID back to UUID format (replace _ with -).
|
||||
* Used when we receive a normalized ID but need the original UUID.
|
||||
*/
|
||||
private denormalizeId(id: string): string {
|
||||
return id?.replace(/_/g, '-') || '';
|
||||
}
|
||||
|
||||
/**
|
||||
* Normalize job ID to match Platform-API format.
|
||||
* Platform-API replaces '-' with '_' in job IDs.
|
||||
*
|
||||
* Example: "2ccf5481-59f5-4036-8a94-7d5f28f4f899-0" -> "2ccf5481_59f5_4036_8a94_7d5f28f4f899_0"
|
||||
*/
|
||||
private normalizeJobId(jobId: string): string {
|
||||
return jobId?.replace(/-/g, '_') || '';
|
||||
}
|
||||
|
||||
/**
|
||||
* 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").
|
||||
* Handles both hyphenated and underscored formats, always returns hyphenated UUID for ES.
|
||||
*
|
||||
* Examples:
|
||||
* - "2ccf5481-59f5-4036-8a94-7d5f28f4f899-0" -> "2ccf5481-59f5-4036-8a94-7d5f28f4f899"
|
||||
* - "2ccf5481_59f5_4036_8a94_7d5f28f4f899_0" -> "2ccf5481-59f5-4036-8a94-7d5f28f4f899"
|
||||
*/
|
||||
private extractPipelineIdFromJobId(jobId: string): string {
|
||||
if (!jobId) return '';
|
||||
|
||||
// Determine the separator used in the jobId
|
||||
const hasUnderscores = jobId.includes('_');
|
||||
const separator = hasUnderscores ? '_' : '-';
|
||||
|
||||
const parts = jobId.split(separator);
|
||||
// UUID has 5 parts (8-4-4-4-12), job suffix is the 6th part
|
||||
if (parts.length >= 6) {
|
||||
// Always return hyphenated format for Elasticsearch lookup
|
||||
return parts.slice(0, 5).join('-');
|
||||
}
|
||||
// If no suffix found, return the ID in hyphenated format
|
||||
return hasUnderscores ? jobId.replace(/_/g, '-') : jobId;
|
||||
}
|
||||
|
||||
private readonly VALID_CONNECTORS = ['jdbc', 'singer', 's3'];
|
||||
private readonly MAX_MEMORY_MB = 12000; // 12GB maximum memory per pipeline/job
|
||||
|
||||
/**
|
||||
* Validate that connector is provided and is a valid type.
|
||||
@@ -62,6 +107,17 @@ export class PlatformApiController {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Validate memory allocation against maximum limit.
|
||||
*/
|
||||
private validateMemory(memoryMb: number): void {
|
||||
if (memoryMb > this.MAX_MEMORY_MB) {
|
||||
throw new BadRequestException(
|
||||
`Memory limit exceeded. Maximum allowed: ${this.MAX_MEMORY_MB}MB (12GB)`,
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Validate cron expression against customer's schedule limit.
|
||||
* Fetches current scheduleLimit from DUC to ensure up-to-date configuration.
|
||||
@@ -107,7 +163,7 @@ export class PlatformApiController {
|
||||
name: string;
|
||||
type: string;
|
||||
columns?: string[];
|
||||
reference_column?: { name: string; type: string };
|
||||
reference_column?: string;
|
||||
}> {
|
||||
if (!jobs || jobs.length === 0) return [];
|
||||
|
||||
@@ -115,7 +171,7 @@ export class PlatformApiController {
|
||||
name: string;
|
||||
type: string;
|
||||
columns?: string[];
|
||||
reference_column?: { name: string; type: string };
|
||||
reference_column?: string;
|
||||
}> = [];
|
||||
|
||||
for (const job of jobs) {
|
||||
@@ -128,7 +184,7 @@ export class PlatformApiController {
|
||||
name: string;
|
||||
type: string;
|
||||
columns?: string[];
|
||||
reference_column?: { name: string; type: string };
|
||||
reference_column?: string;
|
||||
} = {
|
||||
name: input.table_name || '',
|
||||
type: input.load_type || 'full_load',
|
||||
@@ -137,10 +193,8 @@ export class PlatformApiController {
|
||||
table.columns = input.column_include_list;
|
||||
}
|
||||
if (input.incremental_column_name) {
|
||||
table.reference_column = {
|
||||
name: input.incremental_column_name,
|
||||
type: input.incremental_column_type || 'timestamp',
|
||||
};
|
||||
// reference_column is stored as a string (column name)
|
||||
table.reference_column = input.incremental_column_name;
|
||||
}
|
||||
tables.push(table);
|
||||
} else if (connector === 'singer' || connector === 's3') {
|
||||
@@ -205,6 +259,148 @@ export class PlatformApiController {
|
||||
return properties;
|
||||
}
|
||||
|
||||
/**
|
||||
* Sync job input changes to DynamoDB for a specific connector type.
|
||||
* Extracts pipeline ID from job ID, fetches ES document to find input ID,
|
||||
* then updates the table entry in DynamoDB.
|
||||
*
|
||||
* Job ID transformations:
|
||||
* - Raw format (from endpoint): "2ccf5481-59f5-4036-8a94-7d5f28f4f899-0"
|
||||
* - Platform API format: "2ccf5481_59f5_4036_8a94_7d5f28f4f899_0" (underscores)
|
||||
* - Elasticsearch pipeline ID: "2ccf5481-59f5-4036-8a94-7d5f28f4f899" (UUID only, hyphens)
|
||||
*
|
||||
* @param connectorType - The connector type ('jdbc', 'singer', 's3') for the Platform API endpoint
|
||||
*/
|
||||
private async syncJobInputToDynamoDB(
|
||||
jobId: string,
|
||||
body: any,
|
||||
user: RequestUser,
|
||||
connectorType: 'jdbc' | 'singer' | 's3',
|
||||
): Promise<void> {
|
||||
try {
|
||||
// Normalize job ID for Platform API GET (replace - with _)
|
||||
const normalizedJobId = this.normalizeJobId(jobId);
|
||||
|
||||
// Get job details using connector-specific endpoint to find table_name
|
||||
const jobResult = await this.platformApiService.proxy(
|
||||
'GET',
|
||||
`/jobs/${connectorType}/${normalizedJobId}`,
|
||||
user,
|
||||
);
|
||||
|
||||
// Extract the pipeline ID (base UUID) from the raw job ID for ES lookup
|
||||
const esPipelineId = this.extractPipelineIdFromJobId(jobId);
|
||||
const tableName = body.table_name || jobResult.source_config?.table_name;
|
||||
|
||||
if (!esPipelineId || !tableName) {
|
||||
this.logger.warn('Cannot sync job input: missing pipeline_id or table_name', {
|
||||
jobId,
|
||||
esPipelineId,
|
||||
tableName,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
// Get pipeline from ES to find input ID (stored in config.tables)
|
||||
const pipeline = await this.elasticsearchService.getPipeline(
|
||||
user.customer_name,
|
||||
esPipelineId,
|
||||
);
|
||||
|
||||
const inputId = pipeline?.config?.tables;
|
||||
if (!inputId) {
|
||||
this.logger.warn('Cannot sync job input: no input ID in ES', {
|
||||
jobId,
|
||||
esPipelineId,
|
||||
});
|
||||
return;
|
||||
}
|
||||
|
||||
// Build changes for DynamoDB table entry
|
||||
// reference_column is stored as a string (column name), not an object
|
||||
const changes: {
|
||||
type?: string;
|
||||
columns?: string[];
|
||||
reference_column?: string | null;
|
||||
} = {};
|
||||
|
||||
|
||||
if ('target_load_type' in body) {
|
||||
changes.type = body.target_load_type;
|
||||
}
|
||||
if ('column_include_list' in body) {
|
||||
changes.columns = body.column_include_list;
|
||||
}
|
||||
if ('incremental_column_name' in body) {
|
||||
// reference_column is just the column name as a string
|
||||
changes.reference_column = body.incremental_column_name || null;
|
||||
}
|
||||
|
||||
// Update DynamoDB if there are changes
|
||||
if (Object.keys(changes).length > 0) {
|
||||
await this.dynamoDBService.updateInputTable(
|
||||
user.customer_id,
|
||||
inputId,
|
||||
tableName,
|
||||
changes,
|
||||
);
|
||||
}
|
||||
} catch (error) {
|
||||
this.logger.error('Failed to sync job input to DynamoDB', {
|
||||
jobId,
|
||||
connectorType,
|
||||
error: error.message,
|
||||
});
|
||||
// Don't throw - Platform API update succeeded, just log the sync error
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* 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')
|
||||
@@ -446,6 +642,11 @@ export class PlatformApiController {
|
||||
@Body() body: any,
|
||||
@User() user: RequestUser,
|
||||
) {
|
||||
// Validate memory limit
|
||||
if (body.amount) {
|
||||
this.validateMemory(body.amount);
|
||||
}
|
||||
|
||||
return this.platformApiService.proxy(
|
||||
'PUT',
|
||||
`/pipeline/${pipelineId}/memory`,
|
||||
@@ -561,12 +762,23 @@ export class PlatformApiController {
|
||||
@Body() body: any,
|
||||
@User() user: RequestUser,
|
||||
) {
|
||||
return this.platformApiService.proxy(
|
||||
// Normalize job ID for Platform API (replace - with _)
|
||||
const normalizedJobId = this.normalizeJobId(jobId);
|
||||
|
||||
const result = await this.platformApiService.proxy(
|
||||
'PUT',
|
||||
`/jobs/${jobId}/input`,
|
||||
`/jobs/${normalizedJobId}/input`,
|
||||
user,
|
||||
body,
|
||||
);
|
||||
|
||||
// Sync to DynamoDB if connector type is provided
|
||||
const connectorType = body.connector as 'jdbc' | 'singer' | 's3' | undefined;
|
||||
if (connectorType && this.VALID_CONNECTORS.includes(connectorType)) {
|
||||
await this.syncJobInputToDynamoDB(jobId, body, user, connectorType);
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
@Patch('jobs/:jobId/input')
|
||||
@@ -577,12 +789,23 @@ export class PlatformApiController {
|
||||
@Body() body: any,
|
||||
@User() user: RequestUser,
|
||||
) {
|
||||
return this.platformApiService.proxy(
|
||||
// Normalize job ID for Platform API (replace - with _)
|
||||
const normalizedJobId = this.normalizeJobId(jobId);
|
||||
|
||||
const result = await this.platformApiService.proxy(
|
||||
'PATCH',
|
||||
`/jobs/${jobId}/input`,
|
||||
`/jobs/${normalizedJobId}/input`,
|
||||
user,
|
||||
body,
|
||||
);
|
||||
|
||||
// Sync to DynamoDB if connector type is provided
|
||||
const connectorType = body.connector as 'jdbc' | 'singer' | 's3' | undefined;
|
||||
if (connectorType && this.VALID_CONNECTORS.includes(connectorType)) {
|
||||
await this.syncJobInputToDynamoDB(jobId, body, user, connectorType);
|
||||
}
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
@Put('jobs/:jobId/memory')
|
||||
@@ -593,9 +816,17 @@ export class PlatformApiController {
|
||||
@Body() body: any,
|
||||
@User() user: RequestUser,
|
||||
) {
|
||||
// Validate memory limit
|
||||
if (body.amount) {
|
||||
this.validateMemory(body.amount);
|
||||
}
|
||||
|
||||
// Normalize job ID for Platform API (replace - with _)
|
||||
const normalizedJobId = this.normalizeJobId(jobId);
|
||||
|
||||
return this.platformApiService.proxy(
|
||||
'PUT',
|
||||
`/jobs/${jobId}/memory`,
|
||||
`/jobs/${normalizedJobId}/memory`,
|
||||
user,
|
||||
body,
|
||||
);
|
||||
@@ -609,9 +840,12 @@ export class PlatformApiController {
|
||||
@Body() body: any,
|
||||
@User() user: RequestUser,
|
||||
) {
|
||||
// Normalize job ID for Platform API (replace - with _)
|
||||
const normalizedJobId = this.normalizeJobId(jobId);
|
||||
|
||||
return this.platformApiService.proxy(
|
||||
'POST',
|
||||
`/jobs/${jobId}/reset-state`,
|
||||
`/jobs/${normalizedJobId}/reset-state`,
|
||||
user,
|
||||
body,
|
||||
);
|
||||
@@ -623,10 +857,12 @@ export class PlatformApiController {
|
||||
@ApiOperation({ summary: 'Get JDBC job details' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getJdbcJob(@Param('jobId') jobId: string, @User() user: RequestUser) {
|
||||
return this.platformApiService.proxy('GET', `/jobs/jdbc/${jobId}`, user);
|
||||
// Normalize job ID for Platform API (replace - with _)
|
||||
const normalizedJobId = this.normalizeJobId(jobId);
|
||||
return this.platformApiService.proxy('GET', `/jobs/jdbc/${normalizedJobId}`, user);
|
||||
}
|
||||
|
||||
@Put('jobs/jdbc/:jobId/sync-mode')
|
||||
@Post('jobs/jdbc/:jobId/sync-mode')
|
||||
@ApiOperation({ summary: 'Update JDBC job sync mode' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||
async updateJdbcSyncMode(
|
||||
@@ -634,12 +870,20 @@ export class PlatformApiController {
|
||||
@Body() body: any,
|
||||
@User() user: RequestUser,
|
||||
) {
|
||||
return this.platformApiService.proxy(
|
||||
'PUT',
|
||||
`/jobs/jdbc/${jobId}/sync-mode`,
|
||||
// Normalize job ID for Platform API (replace - with _)
|
||||
const normalizedJobId = this.normalizeJobId(jobId);
|
||||
|
||||
const result = await this.platformApiService.proxy(
|
||||
'POST',
|
||||
`/jobs/jdbc/${normalizedJobId}/sync-mode`,
|
||||
user,
|
||||
body,
|
||||
);
|
||||
|
||||
// Sync to DynamoDB (pass raw jobId for pipeline extraction)
|
||||
await this.syncJdbcSyncModeToDynamoDB(jobId, body, user);
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
@Get('jobs/jdbc/configs/allowed_datatypes')
|
||||
@@ -659,10 +903,12 @@ export class PlatformApiController {
|
||||
@ApiOperation({ summary: 'Get Singer job details' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getSingerJob(@Param('jobId') jobId: string, @User() user: RequestUser) {
|
||||
return this.platformApiService.proxy('GET', `/jobs/singer/${jobId}`, user);
|
||||
// Normalize job ID for Platform API (replace - with _)
|
||||
const normalizedJobId = this.normalizeJobId(jobId);
|
||||
return this.platformApiService.proxy('GET', `/jobs/singer/${normalizedJobId}`, user);
|
||||
}
|
||||
|
||||
@Put('jobs/singer/:jobId/sync-mode')
|
||||
@Post('jobs/singer/:jobId/sync-mode')
|
||||
@ApiOperation({ summary: 'Update Singer job sync mode' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.UPDATE)
|
||||
async updateSingerSyncMode(
|
||||
@@ -670,12 +916,20 @@ export class PlatformApiController {
|
||||
@Body() body: any,
|
||||
@User() user: RequestUser,
|
||||
) {
|
||||
return this.platformApiService.proxy(
|
||||
'PUT',
|
||||
`/jobs/singer/${jobId}/sync-mode`,
|
||||
// Normalize job ID for Platform API (replace - with _)
|
||||
const normalizedJobId = this.normalizeJobId(jobId);
|
||||
|
||||
const result = await this.platformApiService.proxy(
|
||||
'POST',
|
||||
`/jobs/singer/${normalizedJobId}/sync-mode`,
|
||||
user,
|
||||
body,
|
||||
);
|
||||
|
||||
// Sync to DynamoDB (pass raw jobId for pipeline extraction)
|
||||
await this.syncSingerSyncModeToDynamoDB(jobId, body, user);
|
||||
|
||||
return result;
|
||||
}
|
||||
|
||||
// ==================== JOBS - S3 ROUTES ====================
|
||||
@@ -684,7 +938,9 @@ export class PlatformApiController {
|
||||
@ApiOperation({ summary: 'Get S3 job details' })
|
||||
@RequireAllPermissions(PERMISSIONS_GROUPS.PIPELINE.permissions.GET)
|
||||
async getS3Job(@Param('jobId') jobId: string, @User() user: RequestUser) {
|
||||
return this.platformApiService.proxy('GET', `/jobs/s3/${jobId}`, user);
|
||||
// Normalize job ID for Platform API (replace - with _)
|
||||
const normalizedJobId = this.normalizeJobId(jobId);
|
||||
return this.platformApiService.proxy('GET', `/jobs/s3/${normalizedJobId}`, user);
|
||||
}
|
||||
|
||||
// ==================== HEALTH ROUTE ====================
|
||||
|
||||
Reference in New Issue
Block a user