export const RAW_DATABASES = ["raw", "raw_staging"] as const; export const RAW_S3_LOCATIONS = [ "s3_etl_tmp", "aspyn_data_lake", "data_lake", "data_event_logs", ] as const; export const RAW_FILE_FORMATS = ["json", "parquet", "aspyn_parquet", "json_gz"] as const; export type RawDatabase = (typeof RAW_DATABASES)[number]; export type RawS3Location = (typeof RAW_S3_LOCATIONS)[number]; export type RawFileFormat = (typeof RAW_FILE_FORMATS)[number]; export interface CreateRawExternalTableParams { database: RawDatabase; schema: string; table: string; s3Location?: RawS3Location; fileFormat?: RawFileFormat; customS3Path?: string; } const IDENTIFIER_PATTERN = /^[A-Za-z_][A-Za-z0-9_$]*$/; const PATH_PATTERN = /^[A-Za-z0-9_./=-]+$/; function assertIdentifier(label: string, value: string): string { const cleaned = value.trim(); if (!IDENTIFIER_PATTERN.test(cleaned)) { throw new Error(`${label} must be an unquoted Snowflake identifier containing only letters, digits, underscore, or dollar sign.`); } return cleaned; } function resolveStagePath( schema: string, table: string, s3Location: RawS3Location, customS3Path: string | undefined, ): { stage: string; path: string } { const stages: Record = { s3_etl_tmp: "S3_ETL_TMP", aspyn_data_lake: "ASPYN_DATA_LAKE", data_lake: "DATA_LAKE", data_event_logs: "DATA_EVENT_LOGS", }; if (customS3Path) { const path = customS3Path.trim().replace(/^\/+/, ""); if (!path || !PATH_PATTERN.test(path) || path.split("/").some((segment) => segment === "." || segment === "..")) { throw new Error("customS3Path may contain only letters, digits, underscore, period, slash, equals, and hyphen, and cannot contain path traversal."); } return { stage: stages[s3Location], path: path.endsWith("/") ? path : `${path}/` }; } if (s3Location === "aspyn_data_lake") { if (table.startsWith("field_service_")) { const [, , remainder] = table.split("_", 3); return { stage: stages[s3Location], path: remainder ? `field_service/${remainder}/` : `${table}/` }; } return { stage: stages[s3Location], path: `${table.replace("_", "/")}/` }; } return { stage: stages[s3Location], path: `${schema}/${table}/` }; } function columnsAndFormat(fileFormat: RawFileFormat): { columns: string; fileFormat: string } { switch (fileFormat) { case "json": return { columns: "raw_json VARIANT AS (value),\n s3_path STRING AS (metadata$filename),\n s3_row_number INTEGER AS (metadata$file_row_number),\n s3_file_last_modified TIMESTAMP AS (metadata$file_last_modified)", fileFormat: "TYPE = JSON STRIP_OUTER_ARRAY = TRUE", }; case "json_gz": return { columns: "raw_json VARIANT AS (value),\n s3_path STRING AS (metadata$filename),\n s3_row_number INTEGER AS (metadata$file_row_number),\n s3_file_last_modified TIMESTAMP AS (metadata$file_last_modified)", fileFormat: "TYPE = JSON STRIP_OUTER_ARRAY = TRUE COMPRESSION = GZIP", }; case "aspyn_parquet": return { columns: "id STRING AS (value:id::string),\n dms_timestamp TIMESTAMP AS (value:dms_timestamp::timestamp),\n all_data VARIANT AS (value),\n s3_path STRING AS (metadata$filename),\n s3_row_number INTEGER AS (metadata$file_row_number),\n s3_file_last_modified TIMESTAMP AS (metadata$file_last_modified)", fileFormat: "TYPE = PARQUET", }; case "parquet": return { columns: "value VARIANT AS (value),\n s3_path STRING AS (metadata$filename),\n s3_row_number INTEGER AS (metadata$file_row_number),\n s3_file_last_modified TIMESTAMP AS (metadata$file_last_modified)", fileFormat: "TYPE = PARQUET", }; } } /** Build the non-destructive external-table DDL used by raw dbt models. */ export function buildCreateRawExternalTableSql(params: CreateRawExternalTableParams): string { if (!RAW_DATABASES.includes(params.database)) throw new Error("database must be raw or raw_staging."); const schema = assertIdentifier("schema", params.schema); const table = assertIdentifier("table", params.table); const s3Location = params.s3Location ?? "s3_etl_tmp"; const fileFormat = params.fileFormat ?? "json"; const { stage, path } = resolveStagePath(schema, table, s3Location, params.customS3Path); const definition = columnsAndFormat(fileFormat); return `CREATE EXTERNAL TABLE IF NOT EXISTS ${params.database}.external_tables.${schema}_${table} (\n ${definition.columns}\n)\nWITH LOCATION = @${params.database}.external_tables.${stage}/${path}\nFILE_FORMAT = (${definition.fileFormat})`; }