import type { PlaySheetContract } from '../plays/static-pipeline'; import { legacyPostgresTruncatedPlayColumnName, sqlSafePlayColumnName, } from './column-names'; export type PhysicalSheetColumnProjection = { sqlName: string; fieldName: string; }; export type LegacyPhysicalSheetColumnBackfill = { sqlName: string; legacySqlName: string; }; const RUNTIME_SHEET_SYSTEM_FIELDS = new Set([ '_key', '_status', '_run_id', '_error', '_stage', '_provider', '_input_index', '_created_at', '_updated_at', '_cell_meta', ]); export function isDatasetPayloadField(field: string): boolean { return ( field.length > 0 && !field.startsWith('__deepline') && !RUNTIME_SHEET_SYSTEM_FIELDS.has(field) ); } export function physicalSheetColumnProjections( sheetContract: PlaySheetContract | null | undefined, ): PhysicalSheetColumnProjection[] { if (!sheetContract) { return []; } const seen = new Set(); const columns: PhysicalSheetColumnProjection[] = []; for (const column of sheetContract.columns) { const sqlName = column.sqlName.trim(); const fieldName = column.source === 'input' || column.source === 'datasetColumn' ? column.field?.trim() || column.id.trim() || sqlName : sqlName; if (!sqlName || sqlName.startsWith('_') || seen.has(sqlName)) { continue; } seen.add(sqlName); columns.push({ sqlName, fieldName }); } return columns; } export function physicalSheetColumnNames( sheetContract: PlaySheetContract | null | undefined, ): string[] { return physicalSheetColumnProjections(sheetContract).map( (column) => column.sqlName, ); } export function legacyPhysicalSheetColumnBackfills( sheetContract: PlaySheetContract | null | undefined, ): LegacyPhysicalSheetColumnBackfill[] { if (!sheetContract) { return []; } const seen = new Set(); const backfills: LegacyPhysicalSheetColumnBackfill[] = []; for (const column of sheetContract.columns) { const sqlName = column.sqlName.trim(); if (!sqlName || sqlName.startsWith('_')) { continue; } const candidates = [ typeof column.field === 'string' ? column.field : null, column.id, ]; for (const candidate of candidates) { if (!candidate) continue; const legacySqlName = legacyPostgresTruncatedPlayColumnName(candidate); const key = legacySqlName ? `${sqlName}:${legacySqlName}` : ''; if (!legacySqlName || legacySqlName === sqlName || seen.has(key)) { continue; } seen.add(key); backfills.push({ sqlName, legacySqlName }); } } return backfills; } export function physicalSheetColumnFieldLookup( sheetContract: PlaySheetContract | null | undefined, ): Map { return new Map( physicalSheetColumnProjections(sheetContract).map((column) => [ column.sqlName, column.fieldName, ]), ); } export function outputFieldsFromSheetContract( sheetContract: PlaySheetContract | null | undefined, ): string[] { const fields = new Set(); for (const column of sheetContract?.columns ?? []) { if (column.source === 'input' || typeof column.field !== 'string') { continue; } fields.add(column.field); } return [...fields]; } export function outputPhysicalSheetColumnProjections( sheetContract: PlaySheetContract | null | undefined, ): PhysicalSheetColumnProjection[] { if (!sheetContract) { return []; } const seen = new Set(); const columns: PhysicalSheetColumnProjection[] = []; for (const column of sheetContract.columns) { const sqlName = column.sqlName.trim(); if ( column.source === 'input' || typeof column.field !== 'string' || !sqlName || sqlName.startsWith('_') || seen.has(sqlName) ) { continue; } seen.add(sqlName); columns.push({ sqlName, fieldName: column.field }); } return columns; } export function outputPhysicalSheetColumnNames( sheetContract: PlaySheetContract | null | undefined, ): string[] { return outputPhysicalSheetColumnProjections(sheetContract).map( (column) => column.sqlName, ); } export function augmentSheetContractWithDatasetFields(input: { contract: PlaySheetContract; rows: readonly Record[]; outputFields?: readonly string[]; }): PlaySheetContract { const outputFields = new Set(input.outputFields ?? []); const candidateFields = new Set(); for (const row of input.rows) { for (const field of Object.keys(row)) { if (isDatasetPayloadField(field)) { candidateFields.add(field); } } } for (const field of outputFields) { if (isDatasetPayloadField(field)) { candidateFields.add(field); } } const existingDatasetPayloadFields = new Set(); const existingSqlNames = new Set(); const inputColumns: PlaySheetContract['columns'] = []; const outputColumns: PlaySheetContract['columns'] = []; const appendColumn = ( target: PlaySheetContract['columns'], column: PlaySheetContract['columns'][number], ) => { const field = typeof column.field === 'string' ? column.field : column.id; const sqlName = column.sqlName.trim(); const isDatasetPayloadColumn = column.source === 'input' || column.source === 'datasetColumn'; if ( !field || !sqlName || (isDatasetPayloadColumn && existingDatasetPayloadFields.has(field)) || existingSqlNames.has(sqlName) ) { return; } if (isDatasetPayloadColumn) { existingDatasetPayloadFields.add(field); } existingSqlNames.add(sqlName); target.push(column); }; for (const column of input.contract.columns) { const field = typeof column.field === 'string' ? column.field : column.id; if ( column.source === 'input' && ((field === input.contract.tableNamespace && !candidateFields.has(field)) || outputFields.has(field)) ) { continue; } appendColumn( column.source === 'input' ? inputColumns : outputColumns, column, ); } for (const field of candidateFields) { if (existingDatasetPayloadFields.has(field)) continue; const sqlName = sqlSafePlayColumnName(field); if (existingSqlNames.has(sqlName)) continue; appendColumn(outputFields.has(field) ? outputColumns : inputColumns, { id: `runtime:${input.contract.tableNamespace}:${field}`, sqlName, source: outputFields.has(field) ? 'datasetColumn' : 'input', field, }); } return { ...input.contract, columns: [...inputColumns, ...outputColumns] }; }