forked from activepieces/activepieces
-
Notifications
You must be signed in to change notification settings - Fork 0
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
feat: add multi-query transaction support to snowflake piece
- Loading branch information
1 parent
4d1dd96
commit bef5b0d
Showing
3 changed files
with
170 additions
and
3 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -1,4 +1,4 @@ | ||
{ | ||
"name": "@activepieces/piece-snowflake", | ||
"version": "0.0.9" | ||
"version": "0.0.10" | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
166 changes: 166 additions & 0 deletions
166
packages/pieces/community/snowflake/src/lib/actions/run-multiple-queries.ts
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,166 @@ | ||
import { createAction, Property } from '@activepieces/pieces-framework'; | ||
import snowflake, { Statement, SnowflakeError } from 'snowflake-sdk'; | ||
import { snowflakeAuth } from '../../index'; | ||
|
||
type QueryResult = unknown[] | undefined; | ||
type QueryResults = { query: string; result: QueryResult }[]; | ||
|
||
const DEFAULT_APPLICATION_NAME = 'ActivePieces'; | ||
const DEFAULT_QUERY_TIMEOUT = 30000; | ||
|
||
export const runMultipleQueries = createAction({ | ||
name: 'runMultipleQueries', | ||
displayName: 'Run Multiple Queries', | ||
description: 'Run Multiple Queries', | ||
auth: snowflakeAuth, | ||
props: { | ||
sqlTexts: Property.Array({ | ||
displayName: 'SQL queries', | ||
description: | ||
'Array of SQL queries to execute in order, in the same transaction. Use :1, :2… placeholders to use binding parameters. ' + | ||
'Avoid using "?" to avoid unexpected behaviors when having multiple queries.', | ||
required: true, | ||
}), | ||
binds: Property.Array({ | ||
displayName: 'Parameters', | ||
description: | ||
'Binding parameters shared across all queries to prevent SQL injection attacks. ' + | ||
'Use :1, :2, etc. to reference parameters in order. ' + | ||
'Avoid using "?" to avoid unexpected behaviors when having multiple queries. ' + | ||
'Unused parameters are allowed.', | ||
required: false, | ||
}), | ||
useTransaction: Property.Checkbox({ | ||
displayName: 'Use Transaction', | ||
description: | ||
'When enabled, all queries will be executed in a single transaction. If any query fails, all changes will be rolled back.', | ||
required: false, | ||
defaultValue: false, | ||
}), | ||
timeout: Property.Number({ | ||
displayName: 'Query timeout (ms)', | ||
description: | ||
'An integer indicating the maximum number of milliseconds to wait for a query to complete before timing out.', | ||
required: false, | ||
defaultValue: DEFAULT_QUERY_TIMEOUT, | ||
}), | ||
application: Property.ShortText({ | ||
displayName: 'Application name', | ||
description: | ||
'A string indicating the name of the client application connecting to the server.', | ||
required: false, | ||
defaultValue: DEFAULT_APPLICATION_NAME, | ||
}), | ||
}, | ||
|
||
async run(context) { | ||
const { username, password, role, database, warehouse, account } = | ||
context.auth; | ||
|
||
const connection = snowflake.createConnection({ | ||
application: context.propsValue.application, | ||
timeout: context.propsValue.timeout, | ||
username, | ||
password, | ||
role, | ||
database, | ||
warehouse, | ||
account, | ||
}); | ||
|
||
return new Promise<QueryResults>((resolve, reject) => { | ||
connection.connect(async function (err: SnowflakeError | undefined) { | ||
if (err) { | ||
reject(err); | ||
return; | ||
} | ||
|
||
const { sqlTexts, binds, useTransaction } = context.propsValue; | ||
const queryResults: QueryResults = []; | ||
|
||
function handleError(err: SnowflakeError) { | ||
if (useTransaction) { | ||
connection.execute({ | ||
sqlText: 'ROLLBACK', | ||
complete: () => { | ||
connection.destroy(() => { | ||
reject(err); | ||
}); | ||
}, | ||
}); | ||
} else { | ||
connection.destroy(() => { | ||
reject(err); | ||
}); | ||
} | ||
} | ||
|
||
async function executeQueriesSequentially() { | ||
try { | ||
if (useTransaction) { | ||
await new Promise<void>((resolveBegin, rejectBegin) => { | ||
connection.execute({ | ||
sqlText: 'BEGIN', | ||
complete: (err: SnowflakeError | undefined) => { | ||
if (err) rejectBegin(err); | ||
else resolveBegin(); | ||
}, | ||
}); | ||
}); | ||
} | ||
for (const sqlText of sqlTexts) { | ||
const result = await new Promise<QueryResult>( | ||
(resolveQuery, rejectQuery) => { | ||
connection.execute({ | ||
sqlText: sqlText as string, | ||
binds: binds as snowflake.Binds, | ||
complete: ( | ||
err: SnowflakeError | undefined, | ||
stmt: Statement, | ||
rows: QueryResult | ||
) => { | ||
if (err) { | ||
rejectQuery(err); | ||
return; | ||
} | ||
resolveQuery(rows); | ||
}, | ||
}); | ||
} | ||
); | ||
|
||
queryResults.push({ | ||
query: sqlText as string, | ||
result, | ||
}); | ||
} | ||
|
||
if (useTransaction) { | ||
await new Promise<void>((resolveCommit, rejectCommit) => { | ||
connection.execute({ | ||
sqlText: 'COMMIT', | ||
complete: (err: SnowflakeError | undefined) => { | ||
if (err) rejectCommit(err); | ||
else resolveCommit(); | ||
}, | ||
}); | ||
}); | ||
} | ||
|
||
connection.destroy((err: SnowflakeError | undefined) => { | ||
if (err) { | ||
reject(err); | ||
return; | ||
} | ||
resolve(queryResults); | ||
}); | ||
} catch (err) { | ||
handleError(err as SnowflakeError); // Reject with the original error! | ||
} | ||
} | ||
|
||
executeQueriesSequentially(); | ||
}); | ||
}); | ||
}, | ||
}); |