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
15 changes: 14 additions & 1 deletion cli/api/dbadapters/postgres_execution_sql.ts
Original file line number Diff line number Diff line change
Expand Up @@ -240,10 +240,23 @@ export class PostgresExecutionSql implements IExecutionSql {
? ` include (${index.include.map(c => `"${c}"`).join(", ")})`
: "";
const where = index.where ? ` where (${index.where})` : "";
return `create ${unique}index "${index.name}" on ${target} using ${method} (${columns})${include}${where}`;
const indexName =
index.name || this.defaultIndexName(table.target.name, index.columns, !!index.unique);
return `create ${unique}index "${indexName}" on ${target} using ${method} (${columns})${include}${where}`;
});
}

// When an index config omits `name`, derive one from the table + columns
// (mirroring Postgres's own `<table>_<col>_idx` default), truncated to the
// 63-char identifier limit. Without this, an unnamed index emits
// `create index "" ...` -> "zero-length delimited identifier".
private defaultIndexName(tableName: string, columns: string[], unique: boolean): string {
const parts = [tableName, ...(columns || [])].filter(Boolean);
const base = parts.join("_");
const suffix = unique ? "_key" : "_idx";
return `${base}${suffix}`.slice(0, 63);
}

private indexMethodAsSql(method?: sqlanvil.PostgresOptions.Index.Method): string {
switch (method) {
case sqlanvil.PostgresOptions.Index.Method.HASH:
Expand Down
26 changes: 26 additions & 0 deletions cli/api/execution_sql_test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -216,6 +216,32 @@ suite("ExecutionSql with Postgres/Supabase", () => {
);
});

test("create index without a name derives one (no zero-length identifier)", () => {
const table: sqlanvil.ITable = {
...baseTable,
postgres: {
indexes: [
{ columns: ["id"] },
{ columns: ["field1", "id"], unique: true }
]
}
};
const statements = executionSql
.publishTasks(table, { fullRefresh: false })
.build()
.map(t => t.statement);
// Names are derived as <table>_<cols>_idx (or _key for unique) -- never "".
expect(statements).to.not.include.members([
'create index "" on "my_db"."public"."my_table" using btree ("id")'
]);
expect(statements[2]).to.equal(
'create index "my_table_id_idx" on "my_db"."public"."my_table" using btree ("id")'
);
expect(statements[3]).to.equal(
'create unique index "my_table_field1_id_key" on "my_db"."public"."my_table" using btree ("field1", "id")'
);
});

test("incremental fresh-create emits postgres indexes", () => {
const incTable: sqlanvil.ITable = {
...baseTable,
Expand Down
40 changes: 22 additions & 18 deletions cli/api/utils/postgres.ts
Original file line number Diff line number Diff line change
Expand Up @@ -65,18 +65,32 @@ export class PgPoolExecutor {
}) => Promise<T>
) {
const client = await this.pool.connect();
// The client can be released from several places: the client "error" handler,
// the query "error" handler (which passes the error to release() so the bad
// connection is closed, not returned to the pool), and the finally block.
// Release exactly once -- otherwise pg-pool's throwOnDoubleRelease fires and a
// confusing "Release called on client which has already been released" surfaces
// ahead of the real error on every failing statement.
let released = false;
const releaseOnce = (err?: Error) => {
if (released) {
return;
}
released = true;
try {
client.release(err);
} catch (e) {
// tslint:disable-next-line: no-console
console.error("Error thrown when releasing pg.Client", e.message, e.stack);
}
};
try {
client.on("error", err => {
// tslint:disable-next-line: no-console
console.error("pg.Client client error", err.message, err.stack);
// Errored connections cause issues when released back to the pool. Instead, close the connection
// by passing the error to release(). https://github.com/dataform-co/dataform/issues/914
try {
client.release(err);
} catch (e) {
// tslint:disable-next-line: no-console
console.error("Error thrown when releasing errored pg.Client", e.message, e.stack);
}
releaseOnce(err);
});

return await callback({
Expand Down Expand Up @@ -114,12 +128,7 @@ export class PgPoolExecutor {
// Errors don't cause "end" to fire, additionally errored connections
// cause issues when released back to the pool. Instead, close the connection
// by passing the error to release(). https://github.com/dataform-co/dataform/issues/914
try {
client.release(err);
} catch (e) {
// tslint:disable-next-line: no-console
console.error("Error thrown when releasing errored pg.Query", e.message, e.stack);
}
releaseOnce(err);
reject(err);
});
query.on("end", () => {
Expand All @@ -129,12 +138,7 @@ export class PgPoolExecutor {
}
});
} finally {
try {
client.release();
} catch (e) {
// tslint:disable-next-line: no-console
console.error("Error thrown when releasing ended pg.Client", e.message, e.stack);
}
releaseOnce();
}
}

Expand Down
50 changes: 50 additions & 0 deletions tests/integration/postgres.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,56 @@ suite("@sqlanvil/integration/postgres", { parallel: true }, ({ before, after })
);
});

test("a failing statement rejects with the real error and no double-release noise", async () => {
// Capture console.error so we can assert the pg client isn't released twice
// (which would print "Release called on client which has already been released").
// tslint:disable-next-line: no-console
const originalError = console.error;
const logged: string[] = [];
// tslint:disable-next-line: no-console
console.error = (...args: any[]) => logged.push(args.join(" "));
let err: Error | undefined;
try {
await dbadapter.execute("selct 1");
} catch (e) {
err = e;
} finally {
// tslint:disable-next-line: no-console
console.error = originalError;
}
expect(err, "a bad statement should reject").to.be.an("error");
expect(err.message.toLowerCase()).to.contain("syntax error");
expect(
logged.join("\n"),
"the pg client must be released exactly once"
).to.not.match(/already been released/i);
});

test("a table with an unnamed postgres index creates (derived index name)", async () => {
const schema = "sa_integration_test_unnamed_index";
try {
await dbadapter.execute(`drop schema if exists "${schema}" cascade`);
} catch (e) {
// ignore
}
await dbadapter.execute(`create schema "${schema}"`);
const table: sqlanvil.ITable = {
enumType: sqlanvil.TableType.TABLE,
target: { schema, name: "with_idx" },
query: "select 1 as id",
postgres: { indexes: [{ columns: ["id"], unique: true }] }
};
const adapter = new ExecutionSql({ warehouse: "postgres" }, "1.4.8");
for (const task of adapter.publishTasks(table, { fullRefresh: true }).build()) {
await dbadapter.execute(task.statement);
}
const { rows } = await dbadapter.execute(
`select indexname from pg_indexes where schemaname = '${schema}' and tablename = 'with_idx'`
);
expect(rows.map((r: any) => r.indexname)).to.include("with_idx_id_key");
await dbadapter.execute(`drop schema if exists "${schema}" cascade`);
});

test("run", { timeout: 60000 }, async () => {
const compiledGraph = await compile("tests/integration/postgres_project", "project_e2e");

Expand Down
2 changes: 1 addition & 1 deletion version.bzl
Original file line number Diff line number Diff line change
Expand Up @@ -2,5 +2,5 @@
# SemVer line). DF_VERSION is the upstream dataform-co/dataform release this fork
# is synced to — surfaced as metadata (e.g. `sqlanvil --version`), not the package
# version. Bump SQLANVIL_VERSION for sqlanvil releases; bump DF_VERSION on upstream syncs.
SQLANVIL_VERSION = "1.4.0"
SQLANVIL_VERSION = "1.4.1"
DF_VERSION = "3.0.59"