mirror of
https://github.com/dadosfera/maestro.git
synced 2026-09-30 00:29:09 +00:00
fix: store reference_column as object with name and type
The protobuf definition expects reference_column to be an object with name and type fields, but it was being stored as just a string (column name). This caused pipeline fetching to fail with the error: ".NewTable.reference_column: object expected" Changes: - Update ReferenceColumn interface in DynamoDB service - Update extractTablesFromJobs to create reference_column object - Update syncJobInputToDynamoDB to handle reference_column object 🤖 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
fb521f53cd
commit
e99306adba
@@ -22,7 +22,7 @@ import { User, RequestUser } from '../../decorators/user.decorator';
|
||||
import { PlatformApiService } from './platform-api.service';
|
||||
import { PERMISSIONS_GROUPS } from '../../authentication/permissions.enum';
|
||||
import { ElasticsearchService } from '../../services/elasticsearch';
|
||||
import { DynamoDBService } from '../../services/dynamodb';
|
||||
import { DynamoDBService, ReferenceColumn } from '../../services/dynamodb';
|
||||
import { CustomersService } from '../customers/customers.service';
|
||||
import { validateCronAgainstScheduleLimit } from '../../utils/cron-validation';
|
||||
|
||||
@@ -155,7 +155,7 @@ export class PlatformApiController {
|
||||
* Extract and transform tables from jobs for DynamoDB input.
|
||||
* Maps connector-specific fields to a common table format.
|
||||
*
|
||||
* - JDBC: load_type, table_name, column_include_list (columns), incremental_column_name (reference_column)
|
||||
* - JDBC: load_type, table_name, column_include_list (columns), incremental_column_name/type (reference_column object)
|
||||
* - Singer: type maps replication_method (FULL_TABLE -> full_load, INCREMENTAL -> incremental), no columns
|
||||
* - S3: same mapping as Singer, no columns
|
||||
*/
|
||||
@@ -163,7 +163,7 @@ export class PlatformApiController {
|
||||
name: string;
|
||||
type: string;
|
||||
columns?: string[];
|
||||
reference_column?: string;
|
||||
reference_column?: ReferenceColumn;
|
||||
}> {
|
||||
if (!jobs || jobs.length === 0) return [];
|
||||
|
||||
@@ -171,7 +171,7 @@ export class PlatformApiController {
|
||||
name: string;
|
||||
type: string;
|
||||
columns?: string[];
|
||||
reference_column?: string;
|
||||
reference_column?: ReferenceColumn;
|
||||
}> = [];
|
||||
|
||||
for (const job of jobs) {
|
||||
@@ -179,12 +179,12 @@ export class PlatformApiController {
|
||||
if (!input) continue;
|
||||
|
||||
if (connector === 'jdbc') {
|
||||
// JDBC: table_name, load_type, column_include_list, incremental_column_name
|
||||
// JDBC: table_name, load_type, column_include_list, incremental_column_name/type
|
||||
const table: {
|
||||
name: string;
|
||||
type: string;
|
||||
columns?: string[];
|
||||
reference_column?: string;
|
||||
reference_column?: ReferenceColumn;
|
||||
} = {
|
||||
name: input.table_name || '',
|
||||
type: input.load_type || 'full_load',
|
||||
@@ -193,8 +193,11 @@ export class PlatformApiController {
|
||||
table.columns = input.column_include_list;
|
||||
}
|
||||
if (input.incremental_column_name) {
|
||||
// reference_column is stored as a string (column name)
|
||||
table.reference_column = input.incremental_column_name;
|
||||
// reference_column is stored as an object with name and type
|
||||
table.reference_column = {
|
||||
name: input.incremental_column_name,
|
||||
type: input.incremental_column_type || 'unknown',
|
||||
};
|
||||
}
|
||||
tables.push(table);
|
||||
} else if (connector === 'singer' || connector === 's3') {
|
||||
@@ -317,11 +320,11 @@ export class PlatformApiController {
|
||||
}
|
||||
|
||||
// Build changes for DynamoDB table entry
|
||||
// reference_column is stored as a string (column name), not an object
|
||||
// reference_column is stored as an object with name and type
|
||||
const changes: {
|
||||
type?: string;
|
||||
columns?: string[];
|
||||
reference_column?: string | null;
|
||||
reference_column?: ReferenceColumn | null;
|
||||
} = {};
|
||||
|
||||
|
||||
@@ -332,8 +335,15 @@ export class PlatformApiController {
|
||||
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;
|
||||
// reference_column is stored as an object with name and type
|
||||
if (body.incremental_column_name) {
|
||||
changes.reference_column = {
|
||||
name: body.incremental_column_name,
|
||||
type: body.incremental_column_type || 'unknown',
|
||||
};
|
||||
} else {
|
||||
changes.reference_column = null;
|
||||
}
|
||||
}
|
||||
|
||||
// Update DynamoDB if there are changes
|
||||
@@ -466,11 +476,13 @@ export class PlatformApiController {
|
||||
});
|
||||
}
|
||||
|
||||
const pipelineType = this.mapConnectorToDynamoType(connectorType);
|
||||
this.logger.info('Syncing pipeline to Elasticsearch', {
|
||||
customerName: user.customer_name,
|
||||
pipelineId: body.id,
|
||||
plugin,
|
||||
connector: connectorType,
|
||||
type: pipelineType,
|
||||
properties,
|
||||
inputId,
|
||||
});
|
||||
@@ -490,7 +502,7 @@ export class PlatformApiController {
|
||||
cron: body.cron,
|
||||
tables: inputId,
|
||||
properties,
|
||||
type: this.mapConnectorToDynamoType(connectorType),
|
||||
type: pipelineType,
|
||||
},
|
||||
connector,
|
||||
);
|
||||
|
||||
Reference in New Issue
Block a user