Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 3 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
@@ -1,8 +1,8 @@
# SQLAnvil

**SQL workflow tool for BigQuery, Postgres, and Supabase.**
**SQL workflow tool for BigQuery, Postgres, Supabase, and MySQL/MariaDB.**

SQLAnvil is an open-source fork of [Dataform OSS](https://github.com/dataform-co/dataform) (Apache 2.0), extended with first-class PostgreSQL and Supabase support. Define your data transformations in SQLX, have SQLAnvil compile them to idiomatic SQL, and run them against your warehouse.
SQLAnvil is an open-source fork of [Dataform OSS](https://github.com/dataform-co/dataform) (Apache 2.0), extended with first-class PostgreSQL, Supabase, and MySQL/MariaDB support. Define your data transformations in SQLX, have SQLAnvil compile them to idiomatic SQL, and run them against your warehouse.

> **SQLAnvil is not affiliated with or endorsed by Google.** The Dataform name and related marks are trademarks of Google LLC. See [NOTICE](NOTICE) for attribution.

Expand All @@ -13,6 +13,7 @@ SQLAnvil is an open-source fork of [Dataform OSS](https://github.com/dataform-co
- **BigQuery** — full support: partitioning, clustering, labels, materialized views, `MERGE`-based incremental upserts
- **PostgreSQL** — idiomatic DDL: native partitioning, `INSERT ... ON CONFLICT` upserts, btree/gin/gist/brin indexes, tablespaces, fillfactor
- **Supabase** — extends Postgres with RLS policies, Realtime publications, pgvector indexes, and Supabase Wrappers _(coming soon)_
- **MySQL / MariaDB** — portable MySQL DDL: CTAS tables, `CREATE OR REPLACE VIEW`, `ON DUPLICATE KEY UPDATE` incremental upserts (one adapter, validated against both engines)
- **SQLX + YAML + JS** — three authoring modes: SQL with config blocks, `actions.yaml` bulk definitions, or the JavaScript API

---
Expand Down
1 change: 1 addition & 0 deletions cli/api/BUILD
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ ts_library(
"@npm//google-sql-syntax-ts",
"@npm//js-beautify",
"@npm//js-yaml",
"@npm//mysql2",
"@npm//pg",
"@npm//promise-pool-executor",
"@npm//protobufjs",
Expand Down
7 changes: 7 additions & 0 deletions cli/api/commands/credentials.ts
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,13 @@ export function read(credentialsPath: string, warehouse: string = "bigquery"): a
// map in workflow_settings.yaml). It is not part of the write-warehouse
// connection, so exclude it from the strict warehouse-credentials validation.
const { connections, ...warehouseCredentials } = credentialsAsJson;
if (warehouse.toLowerCase() === "mysql") {
const credentials = verifyObjectMatchesProto(sqlanvil.MysqlConnection, warehouseCredentials);
if (!credentials.host) {
throw new Error(`Error reading credentials file: the host field is required`);
}
return credentials;
}
const isPostgres = warehouse.toLowerCase() === "postgres" || warehouse.toLowerCase() === "supabase";
if (isPostgres) {
const credentials = verifyObjectMatchesProto(sqlanvil.PostgresConnection, warehouseCredentials);
Expand Down
18 changes: 16 additions & 2 deletions cli/api/commands/init.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,20 @@ function postgresCredentialsTemplate(warehouse: string): string {
return `${JSON.stringify(template, null, 2)}\n`;
}

// A starter MysqlConnection (strict JSON — no comment keys). Points at a local
// MySQL/MariaDB instance with SSL disabled by default.
function mysqlCredentialsTemplate(): string {
const template = {
host: "localhost",
port: 3306,
database: "sqlanvil",
user: "root",
password: "",
sslMode: "disable"
};
return `${JSON.stringify(template, null, 2)}\n`;
}

export async function init(
projectDir: string,
projectConfig: sqlanvil.IProjectConfig
Expand Down Expand Up @@ -100,13 +114,13 @@ export async function init(
fs.writeFileSync(gitignorePath, gitIgnoreContents);
filesWritten.push(gitignorePath);

// Postgres/Supabase: scaffold a credentials template (the connection lives in a separate,
// Postgres/Supabase/MySQL: scaffold a credentials template (the connection lives in a separate,
// gitignored file — not in workflow_settings.yaml). BigQuery credentials come from gcloud / a
// BigQuery key, so no template is written for it.
if (!isBigQuery) {
fs.writeFileSync(
path.join(projectDir, CREDENTIALS_FILENAME),
postgresCredentialsTemplate(warehouse)
warehouse === "mysql" ? mysqlCredentialsTemplate() : postgresCredentialsTemplate(warehouse)
);
filesWritten.push(path.join(projectDir, CREDENTIALS_FILENAME));
}
Expand Down
3 changes: 3 additions & 0 deletions cli/api/dbadapters/execution_sql.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,5 @@
import { BigQueryExecutionSql } from "sa/cli/api/dbadapters/bigquery_execution_sql";
import { MysqlExecutionSql } from "sa/cli/api/dbadapters/mysql_execution_sql";
import { PostgresExecutionSql } from "sa/cli/api/dbadapters/postgres_execution_sql";
import { concatenateQueries, Tasks } from "sa/cli/api/dbadapters/tasks";
import { ErrorWithCause } from "sa/common/errors/errors";
Expand Down Expand Up @@ -36,6 +37,8 @@ export class ExecutionSql implements IExecutionSql {
const warehouse = (project.warehouse || "bigquery").toLowerCase();
if (warehouse === "postgres" || warehouse === "supabase") {
this.delegate = new PostgresExecutionSql(project, sqlanvilCoreVersion, uniqueIdGenerator);
} else if (warehouse === "mysql") {
this.delegate = new MysqlExecutionSql(project, sqlanvilCoreVersion, uniqueIdGenerator);
} else {
this.delegate = new BigQueryExecutionSql(project, sqlanvilCoreVersion, uniqueIdGenerator);
}
Expand Down
278 changes: 278 additions & 0 deletions cli/api/dbadapters/mysql.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,278 @@
import { collectEvaluationQueries, QueryOrAction } from "sa/cli/api/dbadapters/execution_sql";
import {
IDbAdapter,
IDbClient,
IExecutionResult,
IExecutionResultRaw,
OnCancel
} from "sa/cli/api/dbadapters/index";
import { convertFieldType, MySqlPoolExecutor } from "sa/cli/api/utils/mysql";
import { ErrorWithCause } from "sa/common/errors/errors";
import { sqlanvil } from "sa/protos/ts";

// MySQL/MariaDB has no catalog level above the database, so "schema" and
// "database" are the same thing — these are the engine-managed databases we
// never treat as user schemas.
const INTERNAL_SCHEMAS = new Set([
"information_schema",
"mysql",
"performance_schema",
"sys"
]);

export class MySqlDbAdapter implements IDbAdapter {
public static async create(
credentials: sqlanvil.IMysqlConnection,
options?: { concurrencyLimit?: number; disableSslForTestsOnly?: boolean }
): Promise<MySqlDbAdapter> {
const sslMode = (credentials.sslMode || "").toLowerCase();
const ssl =
!options?.disableSslForTestsOnly && sslMode && sslMode !== "disable"
? // Managed MySQL providers serve certs signed by their own CA; skipping
// verification is the documented path for sslmode=require. Stricter
// verification would need a CA bundle we don't ship today.
{ rejectUnauthorized: false }
: undefined;
const queryExecutor = new MySqlPoolExecutor(
{
host: credentials.host,
port: credentials.port || 3306,
user: credentials.user,
password: credentials.password,
database: credentials.database || undefined,
ssl
},
options
);
// Fail fast on a single connection before any command fans out.
try {
await queryExecutor.verifyConnection();
} catch (e) {
await queryExecutor.close().catch(() => undefined);
throw new ErrorWithCause(
`Could not connect to MySQL at ${credentials.host}:${credentials.port || 3306} ` +
`as "${credentials.user}": ${e.message}`,
e
);
}
return new MySqlDbAdapter(queryExecutor);
}

protected constructor(protected readonly queryExecutor: MySqlPoolExecutor) {}

public async execute(
statement: string,
options: {
params?: any[];
onCancel?: OnCancel;
rowLimit?: number;
byteLimit?: number;
includeQueryInError?: boolean;
} = { rowLimit: 1000, byteLimit: 1024 * 1024 }
): Promise<IExecutionResult> {
return await this.withClientLock(client => client.execute(statement, options));
}

public async executeRaw(
statement: string,
options: {
params?: any[];
rowLimit?: number;
} = { rowLimit: 1000 }
): Promise<IExecutionResultRaw> {
const result = await this.execute(statement, options);
return { ...result, schema: [] };
}

public async withClientLock<T>(callback: (client: IDbClient) => Promise<T>): Promise<T> {
return await this.queryExecutor.withClientLock(client =>
callback({
execute: async (
stmt: string,
opts: {
params?: any[];
onCancel?: OnCancel;
rowLimit?: number;
byteLimit?: number;
includeQueryInError?: boolean;
} = { rowLimit: 1000, byteLimit: 1024 * 1024 }
): Promise<IExecutionResult> => {
try {
const rows = await client.execute(stmt, { params: opts.params, rowLimit: opts.rowLimit });
return { rows, metadata: {} };
} catch (e) {
if (opts.includeQueryInError) {
throw new Error(`Error encountered while running "${stmt}": ${e.message}`);
}
throw new ErrorWithCause(`Error executing mysql query: ${e.message}`, e);
}
},
executeRaw: async (
stmt: string,
opts: { params?: { [name: string]: any }; rowLimit?: number } = { rowLimit: 1000 }
): Promise<IExecutionResultRaw> => {
const positional = opts.params ? Object.values(opts.params) : undefined;
const rows = await client.execute(stmt, { params: positional, rowLimit: opts.rowLimit });
return { rows, schema: [], metadata: {} };
}
})
);
}

public async evaluate(queryOrAction: QueryOrAction): Promise<sqlanvil.IQueryEvaluation[]> {
// EXPLAIN parses + plans without executing, catching syntax errors and
// missing tables/columns.
const validationQueries = collectEvaluationQueries(queryOrAction, false, (query: string) =>
!!query ? `explain ${query}` : ""
).map((validationQuery, index) => ({ index, validationQuery }));
const validationQueriesWithoutWrappers = collectEvaluationQueries(queryOrAction, false);

const queryEvaluations = new Array<sqlanvil.IQueryEvaluation>();
for (const { index, validationQuery } of validationQueries) {
let evaluationResponse: sqlanvil.IQueryEvaluation = {
status: sqlanvil.QueryEvaluation.QueryEvaluationStatus.SUCCESS
};
try {
await this.execute(validationQuery.query);
} catch (e) {
evaluationResponse = {
status: sqlanvil.QueryEvaluation.QueryEvaluationStatus.FAILURE,
error: sqlanvil.QueryEvaluationError.create({
message: e?.message ? String(e.message) : String(e)
})
};
}
queryEvaluations.push(
sqlanvil.QueryEvaluation.create({
...evaluationResponse,
incremental: validationQuery.incremental,
query: validationQueriesWithoutWrappers[index].query
})
);
}
return queryEvaluations;
}

public async tables(_database: string, schema?: string): Promise<sqlanvil.ITableMetadata[]> {
const params: any[] = [];
let schemaClause = "";
if (schema) {
schemaClause = "and table_schema = ?";
params.push(schema);
}
const queryResult = await this.execute(
`select table_name, table_schema
from information_schema.tables
where table_schema not in ('information_schema', 'mysql', 'performance_schema', 'sys')
${schemaClause}`,
{ params, rowLimit: 10000, includeQueryInError: true }
);
const targets = queryResult.rows.map(row => ({
schema: row.table_schema as string,
name: row.table_name as string
}));
return await Promise.all(targets.map(target => this.table(target)));
}

public async search(
searchText: string,
options: { limit: number } = { limit: 1000 }
): Promise<sqlanvil.ITableMetadata[]> {
const results = await this.execute(
`select tables.table_schema as table_schema, tables.table_name as table_name
from information_schema.tables as tables
left join information_schema.columns columns
on tables.table_schema = columns.table_schema
and tables.table_name = columns.table_name
where tables.table_schema like ?
or tables.table_name like ?
or columns.column_name like ?
group by 1, 2`,
{
params: [`%${searchText}%`, `%${searchText}%`, `%${searchText}%`],
rowLimit: options.limit
}
);
return await Promise.all(
results.rows.map(row =>
this.table({
schema: row.table_schema,
name: row.table_name
})
)
);
}

public async table(target: sqlanvil.ITarget): Promise<sqlanvil.ITableMetadata> {
const params = [target.schema, target.name];
const [tableResults, columnResults] = await Promise.all([
this.execute(
`select table_type from information_schema.tables
where table_schema = ? and table_name = ?`,
{ params, includeQueryInError: true }
),
this.execute(
`select column_name, data_type, ordinal_position
from information_schema.columns
where table_schema = ? and table_name = ?
order by ordinal_position`,
{ params, includeQueryInError: true }
)
]);

if (tableResults.rows.length === 0) {
return null;
}

// mysql2 returns information_schema column names in their canonical
// upper/lower case depending on server config; normalise via lower-cased keys.
const tableType = String(
tableResults.rows[0].table_type ?? tableResults.rows[0].TABLE_TYPE
).toUpperCase();
return sqlanvil.TableMetadata.create({
target,
type: tableType === "VIEW" ? sqlanvil.TableMetadata.Type.VIEW : sqlanvil.TableMetadata.Type.TABLE,
fields: columnResults.rows.map(row =>
sqlanvil.Field.create({
name: (row.column_name ?? row.COLUMN_NAME) as string,
primitive: convertFieldType((row.data_type ?? row.DATA_TYPE) as string)
})
)
});
}

public async deleteTable(target: sqlanvil.ITarget): Promise<void> {
const metadata = await this.table(target);
if (!metadata) {
return;
}
const kind = metadata.type === sqlanvil.TableMetadata.Type.VIEW ? "view" : "table";
await this.execute(`drop ${kind} if exists \`${target.schema}\`.\`${target.name}\``, {
includeQueryInError: true
});
}

public async schemas(_database: string): Promise<string[]> {
const result = await this.execute(`select schema_name from information_schema.schemata`, {
includeQueryInError: true
});
return result.rows
.map(row => (row.schema_name ?? row.SCHEMA_NAME) as string)
.filter(name => !INTERNAL_SCHEMAS.has(name));
}

public async createSchema(_database: string, schema: string): Promise<void> {
await this.execute(`create database if not exists \`${schema}\``, {
includeQueryInError: true
});
}

public async setMetadata(_action: sqlanvil.IExecutionAction): Promise<void> {
// Deferred for the MVP — table/column COMMENT metadata is a follow-up PR.
return;
}

public async close(): Promise<void> {
await this.queryExecutor.close();
}
}
Loading