This commit is contained in:
Benjamin Toby
2025-07-05 14:59:30 +01:00
parent 6e334c2525
commit 7e8bb37c09
526 changed files with 17560 additions and 11386 deletions
@@ -8,7 +8,7 @@ import {
import _ from "lodash";
import EJSON from "../../utils/ejson";
import generateTypeDefinition from "./generate-type-definitions";
import path from "path";
import { AppNames } from "../../dict/app-names";
type Params = {
dbSchema?: DSQL_DatabaseSchemaType;
@@ -16,16 +16,12 @@ type Params = {
export default function dbSchemaToType(params?: Params): string[] | undefined {
let datasquirelSchema;
const defaultTableFieldsJSONFilePath = path.resolve(
__dirname,
"../../data/defaultFields.json"
);
const { mainShemaJSONFilePath, defaultTableFieldsJSONFilePath } =
grabDirNames();
if (params?.dbSchema) {
datasquirelSchema = params.dbSchema;
} else {
const { mainShemaJSONFilePath } = grabDirNames();
const mainSchema = EJSON.parse(
fs.readFileSync(mainShemaJSONFilePath, "utf-8")
) as DSQL_DatabaseSchemaType[];
@@ -49,7 +45,7 @@ export default function dbSchemaToType(params?: Params): string[] | undefined {
let newDefaultFields = _.cloneDeep(defaultFields);
return {
...tblSchm,
fields: params?.dbSchema
fields: tblSchm.fields.find((fld) => fld.fieldName == "id")
? tblSchm.fields
: [
newDefaultFields.shift(),
@@ -62,7 +58,10 @@ export default function dbSchemaToType(params?: Params): string[] | undefined {
const defDbName = (
datasquirelSchema.dbName ||
datasquirelSchema.dbFullName?.replace(/datasquirel_user_\d+_/, "")
datasquirelSchema.dbFullName?.replace(
new RegExp(`${AppNames["DsqlDbPrefix"]}\\d+_`),
""
)
)
?.toUpperCase()
.replace(/ /g, "_");
+1
View File
@@ -52,6 +52,7 @@ export default function decrypt({
return decrypted;
} catch (error: any) {
console.log("Error in decrypting =>", error.message);
console.log("encryptedString =>", encryptedString);
global.ERROR_CALLBACK?.(`Error Decrypting data`, error as Error);
return encryptedString;
}
@@ -1,16 +1,9 @@
// @ts-check
/**
* Check for user in local storage
* Regular expression to match default fields
*
* @description Preventdefault, declare variables
* @description Regular expression to match default fields
*/
const defaultFieldsRegexp =
/^id$|^uuid$|^date_created$|^date_created_code$|^date_created_timestamp$|^date_updated$|^date_updated_code$|^date_updated_timestamp$/;
////////////////////////////////////////
////////////////////////////////////////
////////////////////////////////////////
/^id$|^uuid$|^uid$|^date_created$|^date_created_code$|^date_created_timestamp$|^date_updated$|^date_updated_code$|^date_updated_timestamp$/;
export default defaultFieldsRegexp;
@@ -1,4 +1,4 @@
import { DSQL_TableSchemaType } from "../../types";
import { DSQL_FieldSchemaType, DSQL_TableSchemaType } from "../../types";
import defaultFieldsRegexp from "./default-fields-regexp";
type Param = {
@@ -8,6 +8,7 @@ type Param = {
typeDefName?: string;
allValuesOptional?: boolean;
addExport?: boolean;
dbName?: string;
};
export default function generateTypeDefinition({
@@ -17,21 +18,34 @@ export default function generateTypeDefinition({
typeDefName,
allValuesOptional,
addExport,
dbName,
}: Param): string | null {
let typeDefinition: string | null = ``;
try {
const tdName =
typeDefName ||
`DSQL_${query.single}_${query.single_table}`.toUpperCase();
const tdName = typeDefName
? typeDefName
: dbName
? `DSQL_${dbName}_${table.tableName}`.toUpperCase()
: `DSQL_${query.single}_${query.single_table}`.toUpperCase();
const fields = table.fields;
function typeMap(type: string) {
if (type?.match(/int/i)) {
function typeMap(schemaType: DSQL_FieldSchemaType) {
if (schemaType.options && schemaType.options.length > 0) {
return schemaType.options
.map((opt) =>
schemaType.dataType?.match(/int/i) ||
typeof opt == "number"
? `${opt}`
: `"${opt}"`
)
.join(" | ");
}
if (schemaType.dataType?.match(/int/i)) {
return "number";
}
if (type?.match(/text|varchar|timestamp/i)) {
if (schemaType.dataType?.match(/text|varchar|timestamp/i)) {
return "string";
}
@@ -48,21 +62,19 @@ export default function generateTypeDefinition({
fields.forEach((field) => {
const nullValue = allValuesOptional
? "?"
: field.nullValue
? "?"
: field.fieldName?.match(defaultFieldsRegexp)
? "?"
: "";
: field.notNullValue
? ""
: "?";
typesArrayTypeScript.push(
` ${field.fieldName}${nullValue}: ${typeMap(
field.dataType || ""
)};`
` ${field.fieldName}${nullValue}: ${typeMap(field)};`
);
typesArrayJavascript.push(
` * @property {${typeMap(field.dataType || "")}${nullValue}} ${
` * @property {${typeMap(field)}${nullValue}} ${
field.fieldName
}`
);
@@ -1,3 +1,6 @@
import { SQLDeleteGeneratorParams } from "../../../types";
import sqlEqualityParser from "../../../utils/sql-equality-parser";
interface SQLDeleteGenReturn {
query: string;
values: string[];
@@ -8,13 +11,10 @@ interface SQLDeleteGenReturn {
*/
export default function sqlDeleteGenerator({
tableName,
data,
deleteKeyValues,
dbFullName,
}: {
data: any;
tableName: string;
dbFullName?: string;
}): SQLDeleteGenReturn | undefined {
data,
}: SQLDeleteGeneratorParams): SQLDeleteGenReturn | undefined {
const finalDbName = dbFullName ? `${dbFullName}.` : "";
try {
@@ -23,17 +23,44 @@ export default function sqlDeleteGenerator({
let deleteBatch: string[] = [];
let queryArr: string[] = [];
Object.keys(data).forEach((ky) => {
deleteBatch.push(`${ky}=?`);
queryArr.push(data[ky]);
});
if (data) {
Object.keys(data).forEach((ky) => {
let value = data[ky] as string | number | null | undefined;
const parsedValue =
typeof value == "number" ? String(value) : value;
if (!parsedValue) return;
if (parsedValue.match(/%/)) {
deleteBatch.push(`${ky} LIKE ?`);
queryArr.push(parsedValue);
} else {
deleteBatch.push(`${ky}=?`);
queryArr.push(parsedValue);
}
});
} else if (deleteKeyValues) {
deleteKeyValues.forEach((ky) => {
let value = ky.value as string | number | null | undefined;
const parsedValue =
typeof value == "number" ? String(value) : value;
if (!parsedValue) return;
const operator = sqlEqualityParser(ky.operator || "EQUAL");
deleteBatch.push(`${ky.key} ${operator} ?`);
queryArr.push(parsedValue);
});
}
queryStr += ` WHERE ${deleteBatch.join(" AND ")}`;
return {
query: queryStr,
values: queryArr,
};
} catch (/** @type {any} */ error: any) {
} catch (error: any) {
console.log(`SQL delete gen ERROR: ${error.message}`);
return undefined;
}
@@ -0,0 +1,50 @@
import sqlEqualityParser from "../../../utils/sql-equality-parser";
import { ServerQueryEqualities } from "../../../types";
type Params = {
fieldName: string;
value?: string;
equality: (typeof ServerQueryEqualities)[number];
};
/**
* # SQL Gen Operator Gen
* @description Generates an SQL operator for node module `mysql` or `serverless-mysql`
*/
export default function sqlGenOperatorGen({
fieldName,
value,
equality,
}: Params): string {
if (value) {
if (equality == "LIKE") {
return `LOWER(${fieldName}) LIKE LOWER('%${value}%')`;
} else if (equality == "LIKE_RAW") {
return `LOWER(${fieldName}) LIKE LOWER('${value}')`;
} else if (equality == "NOT LIKE") {
return `LOWER(${fieldName}) NOT LIKE LOWER('%${value}%')`;
} else if (equality == "NOT LIKE_RAW") {
return `LOWER(${fieldName}) NOT LIKE LOWER('${value}')`;
} else if (equality == "REGEXP") {
return `LOWER(${fieldName}) REGEXP LOWER('${value}')`;
} else if (equality == "FULLTEXT") {
return `MATCH(${fieldName}) AGAINST('${value}' IN BOOLEAN MODE)`;
} else if (equality == "NOT EQUAL") {
return `${fieldName} != ${value}`;
} else if (equality) {
return `${fieldName} ${sqlEqualityParser(equality)} ${value}`;
} else {
return `${fieldName} = ${value}`;
}
} else {
if (equality == "IS NULL") {
return `${fieldName} IS NULL`;
} else if (equality == "IS NOT NULL") {
return `${fieldName} IS NOT NULL`;
} else if (equality) {
return `${fieldName} ${sqlEqualityParser(equality)} ?`;
} else {
return `${fieldName} = ?`;
}
}
}
@@ -1,3 +1,4 @@
import sqlEqualityParser from "../../../utils/sql-equality-parser";
import {
ServerQueryParam,
ServerQueryParamsJoin,
@@ -64,16 +65,30 @@ export default function sqlGenerator<
typeof queryObj.value == "number"
) {
const valueParsed = String(queryObj.value);
const operator = sqlEqualityParser(queryObj.equality || "EQUAL");
if (queryObj.equality == "LIKE") {
str = `LOWER(${finalFieldName}) LIKE LOWER('%${valueParsed}%')`;
} else if (queryObj.equality == "LIKE_RAW") {
str = `LOWER(${finalFieldName}) LIKE LOWER(?)`;
sqlSearhValues.push(valueParsed);
} else if (queryObj.equality == "NOT LIKE") {
str = `LOWER(${finalFieldName}) NOT LIKE LOWER('%${valueParsed}%')`;
} else if (queryObj.equality == "NOT LIKE_RAW") {
str = `LOWER(${finalFieldName}) NOT LIKE LOWER(?)`;
sqlSearhValues.push(valueParsed);
} else if (queryObj.equality == "REGEXP") {
str = `${finalFieldName} REGEXP '${valueParsed}'`;
str = `LOWER(${finalFieldName}) REGEXP LOWER(?)`;
sqlSearhValues.push(valueParsed);
} else if (queryObj.equality == "FULLTEXT") {
str = `MATCH(${finalFieldName}) AGAINST('${valueParsed}' IN BOOLEAN MODE)`;
str = `MATCH(${finalFieldName}) AGAINST(? IN BOOLEAN MODE)`;
sqlSearhValues.push(valueParsed);
} else if (queryObj.equality == "NOT EQUAL") {
str = `${finalFieldName} != ?`;
sqlSearhValues.push(valueParsed);
} else if (queryObj.equality) {
str = `${finalFieldName} ${operator} ?`;
sqlSearhValues.push(valueParsed);
} else {
sqlSearhValues.push(valueParsed);
}
@@ -173,7 +188,7 @@ export default function sqlGenerator<
} else if (genObject?.selectFields?.[0]) {
if (genObject.join) {
str += ` ${genObject.selectFields
?.map((fld) => `${finalDbName}${tableName}.${fld}`)
?.map((fld) => `${finalDbName}${tableName}.${String(fld)}`)
.join(",")}`;
} else {
str += ` ${genObject.selectFields?.join(",")}`;
@@ -272,12 +287,19 @@ export default function sqlGenerator<
queryString += ` WHERE ${sqlSearhString.join(` ${stringOperator} `)}`;
}
if (genObject?.order && !count)
if (genObject?.group?.[0]) {
queryString += ` GROUP BY ${genObject.group
.map((g) => `\`${g.toString()}\``)
.join(",")}`;
}
if (genObject?.order && !count) {
queryString += ` ORDER BY ${
genObject.join
? `${finalDbName}${tableName}.${String(genObject.order.field)}`
: String(genObject.order.field)
} ${genObject.order.strategy}`;
}
if (genObject?.limit && !count) queryString += ` LIMIT ${genObject.limit}`;
if (genObject?.offset && !count)
@@ -0,0 +1,117 @@
import mysql, { Connection } from "mysql";
import { exec } from "child_process";
import { promisify } from "util";
// Configuration interface
interface DatabaseConfig {
host: string;
user: string;
password: string;
database?: string; // Optional for global connection
}
// Master status interface
interface MasterStatus {
File: string;
Position: number;
Binlog_Do_DB?: string;
Binlog_Ignore_DB?: string;
}
function getConnection(config: DatabaseConfig): Connection {
return mysql.createConnection(config);
}
function getMasterStatus(config: DatabaseConfig): Promise<MasterStatus> {
return new Promise((resolve, reject) => {
const connection = getConnection(config);
connection.query("SHOW MASTER STATUS", (error, results) => {
connection.end();
if (error) reject(error);
else resolve(results[0] as MasterStatus);
});
});
}
async function syncDatabases() {
const config: DatabaseConfig = {
host: "localhost",
user: "root",
password: "your_password",
};
let lastPosition: number | null = null; // Track last synced position
while (true) {
try {
// Get current master status
const { File, Position } = await getMasterStatus(config);
// Determine start position (use lastPosition or 4 if first run)
const startPosition = lastPosition !== null ? lastPosition + 1 : 4;
if (startPosition >= Position) {
await new Promise((resolve) => setTimeout(resolve, 5000)); // Wait 5 seconds if no new changes
continue;
}
// Execute mysqlbinlog to get changes
const execPromise = promisify(exec);
const { stdout } = await execPromise(
`mysqlbinlog --database=db_master ${File} --start-position=${startPosition} --stop-position=${Position}`
);
if (stdout) {
const connection = getConnection({
...config,
database: "db_slave",
});
return new Promise((resolve, reject) => {
connection.query(stdout, (error) => {
connection.end();
if (error) reject(error);
else {
lastPosition = Position;
console.log(
`Synced up to position ${Position} at ${new Date().toISOString()}`
);
resolve(null);
}
});
});
}
} catch (error) {
console.error("Sync error:", error);
}
await new Promise((resolve) => setTimeout(resolve, 5000)); // Check every 5 seconds
}
}
// Initialize db_slave with db_master data
async function initializeSlave() {
const config: DatabaseConfig = {
host: "localhost",
user: "root",
password: "your_password",
};
try {
await promisify(exec)(
`mysqldump -u ${config.user} -p${config.password} db_master > db_master_backup.sql`
);
await promisify(exec)(
`mysql -u ${config.user} -p${config.password} db_slave < db_master_backup.sql`
);
console.log("Slave initialized with master data");
} catch (error) {
console.error("Initialization error:", error);
}
}
// Run the sync process
async function main() {
await initializeSlave();
await syncDatabases();
}
main().catch(console.error);
@@ -0,0 +1,3 @@
type Params = {};
function createDuplicateTablesTriggers({}: Params) {}
@@ -0,0 +1,106 @@
```sql
DELIMITER //
CREATE PROCEDURE dsql_replicate_databases(IN source_db VARCHAR(64), IN target_db VARCHAR(64))
BEGIN
-- Declare variables
DECLARE done INT DEFAULT FALSE;
DECLARE table_name VARCHAR(64);
DECLARE column_list TEXT;
DECLARE trigger_sql TEXT;
-- Cursor to iterate over tables in source_db
DECLARE cur CURSOR FOR
SELECT TABLE_NAME
FROM INFORMATION_SCHEMA.TABLES
WHERE TABLE_SCHEMA = source_db;
-- Handler for end of cursor
DECLARE CONTINUE HANDLER FOR NOT FOUND SET done = TRUE;
-- Start transaction to ensure consistency
START TRANSACTION;
-- Open cursor
OPEN cur;
read_loop: LOOP
FETCH cur INTO table_name;
IF done THEN
LEAVE read_loop;
END IF;
-- Dynamically get column names for the table
SELECT GROUP_CONCAT(CONCAT('NEW.', COLUMN_NAME))
INTO column_list
FROM INFORMATION_SCHEMA.COLUMNS
WHERE TABLE_SCHEMA = source_db
AND TABLE_NAME = table_name;
-- Drop existing triggers if they exist
SET @drop_trigger_insert = CONCAT('DROP TRIGGER IF EXISTS after_insert_', table_name);
SET @drop_trigger_update = CONCAT('DROP TRIGGER IF EXISTS after_update_', table_name);
SET @drop_trigger_delete = CONCAT('DROP TRIGGER IF EXISTS after_delete_', table_name);
PREPARE stmt_drop_insert FROM @drop_trigger_insert;
EXECUTE stmt_drop_insert;
DEALLOCATE PREPARE stmt_drop_insert;
PREPARE stmt_drop_update FROM @drop_trigger_update;
EXECUTE stmt_drop_update;
DEALLOCATE PREPARE stmt_drop_update;
PREPARE stmt_drop_delete FROM @drop_trigger_delete;
EXECUTE stmt_drop_delete;
DEALLOCATE PREPARE stmt_drop_delete;
-- Create INSERT trigger
SET @trigger_sql = CONCAT(
'CREATE TRIGGER after_insert_', table_name,
' AFTER INSERT ON ', source_db, '.', table_name, ' FOR EACH ROW ',
'BEGIN ',
'INSERT INTO ', target_db, '.', table_name, ' (',
(SELECT GROUP_CONCAT(COLUMN_NAME) FROM INFORMATION_SCHEMA.COLUMNS WHERE TABLE_SCHEMA = source_db AND TABLE_NAME = table_name), ') ',
'VALUES (', column_list, '); ',
'END;'
);
PREPARE stmt FROM @trigger_sql;
EXECUTE stmt;
DEALLOCATE PREPARE stmt;
-- Create UPDATE trigger
SET @trigger_sql = CONCAT(
'CREATE TRIGGER after_update_', table_name,
' AFTER UPDATE ON ', source_db, '.', table_name, ' FOR EACH ROW ',
'BEGIN ',
'UPDATE ', target_db, '.', table_name, ' SET ',
(SELECT GROUP_CONCAT(CONCAT(COLUMN_NAME, '=NEW.', COLUMN_NAME))
FROM INFORMATION_SCHEMA.COLUMNS WHERE TABLE_SCHEMA = source_db AND TABLE_NAME = table_name),
' WHERE ',
(SELECT CONCAT('id=NEW.id')
FROM INFORMATION_SCHEMA.COLUMNS
WHERE TABLE_SCHEMA = source_db AND TABLE_NAME = table_name AND COLUMN_NAME = 'id' LIMIT 1), '; ',
'END;'
);
PREPARE stmt FROM @trigger_sql;
EXECUTE stmt;
DEALLOCATE PREPARE stmt;
-- Create DELETE trigger
SET @trigger_sql = CONCAT(
'CREATE TRIGGER after_delete_', table_name,
' AFTER DELETE ON ', source_db, '.', table_name, ' FOR EACH ROW ',
'BEGIN ',
'DELETE FROM ', target_db, '.', table_name, ' WHERE id=OLD.id; ',
'END;'
);
PREPARE stmt FROM @trigger_sql;
EXECUTE stmt;
DEALLOCATE PREPARE stmt;
END LOOP;
CLOSE cur;
COMMIT;
END //
DELIMITER ;
```
@@ -0,0 +1,212 @@
DELIMITER / / CREATE PROCEDURE replicate_databases(
IN source_db VARCHAR(64),
IN target_db VARCHAR(64)
) BEGIN -- Declare variables
DECLARE done INT DEFAULT FALSE;
DECLARE table_name VARCHAR(64);
DECLARE column_list TEXT;
DECLARE trigger_sql TEXT;
-- Cursor to iterate over tables in source_db
DECLARE cur CURSOR FOR
SELECT
TABLE_NAME
FROM
INFORMATION_SCHEMA.TABLES
WHERE
TABLE_SCHEMA = source_db;
-- Handler for end of cursor
DECLARE CONTINUE HANDLER FOR NOT FOUND
SET
done = TRUE;
-- Start transaction to ensure consistency
START TRANSACTION;
-- Open cursor
OPEN cur;
read_loop: LOOP FETCH cur INTO table_name;
IF done THEN LEAVE read_loop;
END IF;
-- Dynamically get column names for the table
SELECT
GROUP_CONCAT(CONCAT('NEW.', COLUMN_NAME)) INTO column_list
FROM
INFORMATION_SCHEMA.COLUMNS
WHERE
TABLE_SCHEMA = source_db
AND TABLE_NAME = table_name;
-- Drop existing triggers if they exist
SET
@drop_trigger_insert = CONCAT(
'DROP TRIGGER IF EXISTS after_insert_',
table_name
);
SET
@drop_trigger_update = CONCAT(
'DROP TRIGGER IF EXISTS after_update_',
table_name
);
SET
@drop_trigger_delete = CONCAT(
'DROP TRIGGER IF EXISTS after_delete_',
table_name
);
PREPARE stmt_drop_insert
FROM
@drop_trigger_insert;
EXECUTE stmt_drop_insert;
DEALLOCATE PREPARE stmt_drop_insert;
PREPARE stmt_drop_update
FROM
@drop_trigger_update;
EXECUTE stmt_drop_update;
DEALLOCATE PREPARE stmt_drop_update;
PREPARE stmt_drop_delete
FROM
@drop_trigger_delete;
EXECUTE stmt_drop_delete;
DEALLOCATE PREPARE stmt_drop_delete;
-- Create INSERT trigger
SET
@trigger_sql = CONCAT(
'CREATE TRIGGER after_insert_',
table_name,
' AFTER INSERT ON ',
source_db,
'.',
table_name,
' FOR EACH ROW ',
'BEGIN ',
'INSERT INTO ',
target_db,
'.',
table_name,
' (',
(
SELECT
GROUP_CONCAT(COLUMN_NAME)
FROM
INFORMATION_SCHEMA.COLUMNS
WHERE
TABLE_SCHEMA = source_db
AND TABLE_NAME = table_name
),
') ',
'VALUES (',
column_list,
'); ',
'END;'
);
PREPARE stmt
FROM
@trigger_sql;
EXECUTE stmt;
DEALLOCATE PREPARE stmt;
-- Create UPDATE trigger
SET
@trigger_sql = CONCAT(
'CREATE TRIGGER after_update_',
table_name,
' AFTER UPDATE ON ',
source_db,
'.',
table_name,
' FOR EACH ROW ',
'BEGIN ',
'UPDATE ',
target_db,
'.',
table_name,
' SET ',
(
SELECT
GROUP_CONCAT(CONCAT(COLUMN_NAME, '=NEW.', COLUMN_NAME))
FROM
INFORMATION_SCHEMA.COLUMNS
WHERE
TABLE_SCHEMA = source_db
AND TABLE_NAME = table_name
),
' WHERE ',
(
SELECT
CONCAT('id=NEW.id')
FROM
INFORMATION_SCHEMA.COLUMNS
WHERE
TABLE_SCHEMA = source_db
AND TABLE_NAME = table_name
AND COLUMN_NAME = 'id'
LIMIT
1
), '; ', 'END;'
);
PREPARE stmt
FROM
@trigger_sql;
EXECUTE stmt;
DEALLOCATE PREPARE stmt;
-- Create DELETE trigger
SET
@trigger_sql = CONCAT(
'CREATE TRIGGER after_delete_',
table_name,
' AFTER DELETE ON ',
source_db,
'.',
table_name,
' FOR EACH ROW ',
'BEGIN ',
'DELETE FROM ',
target_db,
'.',
table_name,
' WHERE id=OLD.id; ',
'END;'
);
PREPARE stmt
FROM
@trigger_sql;
EXECUTE stmt;
DEALLOCATE PREPARE stmt;
END LOOP;
CLOSE cur;
COMMIT;
END / / DELIMITER;
@@ -0,0 +1,23 @@
export const TriggerParadigms = ["sync_tables", "sync_dbs"] as const;
type Params = {
userId?: string | number;
paradigm: (typeof TriggerParadigms)[number];
dbId?: string | number;
tableName?: string;
};
export default function grabTriggerName({
userId,
paradigm,
dbId,
tableName,
}: Params) {
let triggerName = `dsql_trig_${paradigm}`;
if (userId) triggerName += `_${userId}`;
if (dbId) triggerName += `_${dbId}`;
if (tableName) triggerName += `_${tableName}`;
return triggerName;
}
@@ -0,0 +1,44 @@
import { DSQL_DatabaseSchemaType, DSQL_TableSchemaType } from "../../../types";
const TriggerTypes = [
{
name: "after_insert",
value: "INSERT",
},
{
name: "after_update",
value: "UPDATE",
},
{
name: "after_delete",
value: "DELETE",
},
] as const;
export type TriggerSQLGenParams = {
type: (typeof TriggerTypes)[number];
srcDbSchema: DSQL_DatabaseSchemaType;
srcTableSchema: DSQL_TableSchemaType;
content: string;
proceedureName: string;
};
export default function triggerSQLGen({
type,
srcDbSchema,
srcTableSchema,
content,
proceedureName,
}: TriggerSQLGenParams) {
let sql = `DELIMITER //\n`;
sql += `CREATE PROCEDURE ${proceedureName}`;
sql += `\nBEGIN`;
sql += ` ${content}`;
sql += `\nEND //`;
sql += `\nDELIMITER\n`;
return sql;
}
@@ -0,0 +1,48 @@
import { DSQL_DatabaseSchemaType, DSQL_TableSchemaType } from "../../../types";
import triggerSQLGen, { TriggerSQLGenParams } from "./trigger-sql-gen";
type Params = TriggerSQLGenParams & {
dstDbSchema: DSQL_DatabaseSchemaType;
dstTableSchema: DSQL_TableSchemaType;
};
export default function tableReplicationTriggerSQLGen({
type,
dstDbSchema,
dstTableSchema,
srcDbSchema,
srcTableSchema,
userId,
paradigm,
}: Params) {
let sql = `CREATE TRIGGER`;
const srcColumns = srcTableSchema.fields
.map((fld) => fld.fieldName)
.filter((fld) => typeof fld == "string");
const dstColumns = dstTableSchema.fields
.map((fld) => fld.fieldName)
.filter((fld) => typeof fld == "string");
if (type.name == "after_insert") {
sql += ` INSERT INTO ${dstDbSchema.dbFullName}.${dstTableSchema.tableName}`;
sql += ` (${dstColumns.join(",")})`;
sql += ` VALUES (${dstColumns.map((c) => `NEW.${c}`).join(",")})`;
} else if (type.name == "after_update") {
sql += ` UPDATE ${dstDbSchema.dbFullName}.${dstTableSchema.tableName}`;
sql += ` SET ${dstColumns.map((c) => `${c}=NEW.${c}`).join(",")}`;
sql += ` WHERE id = NEW.id`;
} else if (type.name == "after_delete") {
sql += ` DELETE FROM ${dstDbSchema.dbFullName}.${dstTableSchema.tableName}`;
sql += ` WHERE id = OLD.id`;
}
return triggerSQLGen({
content: sql,
srcDbSchema,
srcTableSchema,
type,
paradigm,
userId,
});
}
@@ -0,0 +1,92 @@
```sql
DELIMITER //
CREATE PROCEDURE dsql_replicate_two_tables(
IN source_db VARCHAR(64),
IN target_db VARCHAR(64),
IN source_table VARCHAR(64),
IN target_table VARCHAR(64)
)
BEGIN
-- Declare variables
DECLARE column_list TEXT;
DECLARE set_clause TEXT;
DECLARE trigger_sql TEXT;
-- Start transaction to ensure consistency
START TRANSACTION;
-- Dynamically get column names for the source table
SELECT GROUP_CONCAT(CONCAT('NEW.', COLUMN_NAME))
INTO column_list
FROM INFORMATION_SCHEMA.COLUMNS
WHERE TABLE_SCHEMA = source_db
AND TABLE_NAME = source_table;
SELECT GROUP_CONCAT(CONCAT(COLUMN_NAME, '=NEW.', COLUMN_NAME))
INTO set_clause
FROM INFORMATION_SCHEMA.COLUMNS
WHERE TABLE_SCHEMA = source_db
AND TABLE_NAME = source_table;
-- Drop existing triggers if they exist
SET @drop_trigger_insert = CONCAT('DROP TRIGGER IF EXISTS after_insert_', source_table);
SET @drop_trigger_update = CONCAT('DROP TRIGGER IF EXISTS after_update_', source_table);
SET @drop_trigger_delete = CONCAT('DROP TRIGGER IF EXISTS after_delete_', source_table);
PREPARE stmt_drop_insert FROM @drop_trigger_insert;
EXECUTE stmt_drop_insert;
DEALLOCATE PREPARE stmt_drop_insert;
PREPARE stmt_drop_update FROM @drop_trigger_update;
EXECUTE stmt_drop_update;
DEALLOCATE PREPARE stmt_drop_update;
PREPARE stmt_drop_delete FROM @drop_trigger_delete;
EXECUTE stmt_drop_delete;
DEALLOCATE PREPARE stmt_drop_delete;
-- Create INSERT trigger
SET @trigger_sql = CONCAT(
'CREATE TRIGGER after_insert_', source_table,
' AFTER INSERT ON ', source_db, '.', source_table, ' FOR EACH ROW ',
'BEGIN ',
'INSERT INTO ', target_db, '.', target_table, ' (',
(SELECT GROUP_CONCAT(COLUMN_NAME) FROM INFORMATION_SCHEMA.COLUMNS WHERE TABLE_SCHEMA = source_db AND TABLE_NAME = source_table), ') ',
'VALUES (', column_list, '); ',
'END;'
);
PREPARE stmt FROM @trigger_sql;
EXECUTE stmt;
DEALLOCATE PREPARE stmt;
-- Create UPDATE trigger
-- Assume 'id' as the primary key; adjust if different
SET @trigger_sql = CONCAT(
'CREATE TRIGGER after_update_', source_table,
' AFTER UPDATE ON ', source_db, '.', source_table, ' FOR EACH ROW ',
'BEGIN ',
'UPDATE ', target_db, '.', target_table, ' SET ',
set_clause,
' WHERE id = NEW.id; ',
'END;'
);
PREPARE stmt FROM @trigger_sql;
EXECUTE stmt;
DEALLOCATE PREPARE stmt;
-- Create DELETE trigger
SET @trigger_sql = CONCAT(
'CREATE TRIGGER after_delete_', source_table,
' AFTER DELETE ON ', source_db, '.', source_table, ' FOR EACH ROW ',
'BEGIN ',
'DELETE FROM ', target_db, '.', target_table, ' WHERE id = OLD.id; ',
'END;'
);
PREPARE stmt FROM @trigger_sql;
EXECUTE stmt;
DEALLOCATE PREPARE stmt;
COMMIT;
END //
DELIMITER ;
```
@@ -0,0 +1,53 @@
import { DSQL_DatabaseSchemaType, DSQL_TableSchemaType } from "../../../types";
import grabTriggerName, { TriggerParadigms } from "./grab-trigger-name";
const TriggerTypes = [
{
name: "after_insert",
value: "INSERT",
},
{
name: "after_update",
value: "UPDATE",
},
{
name: "after_delete",
value: "DELETE",
},
] as const;
export type TriggerSQLGenParams = {
type: (typeof TriggerTypes)[number];
srcDbSchema: DSQL_DatabaseSchemaType;
srcTableSchema: DSQL_TableSchemaType;
content: string;
userId?: string | number;
paradigm: (typeof TriggerParadigms)[number];
};
export default function triggerSQLGen({
type,
srcDbSchema,
srcTableSchema,
content,
userId,
paradigm,
}: TriggerSQLGenParams) {
let sql = `CREATE TRIGGER`;
let triggerName = grabTriggerName({
paradigm,
dbId: srcDbSchema.id,
tableName: srcTableSchema.tableName,
userId,
});
sql += ` ${triggerName}`;
sql += ` AFTER ${type.value} ON ${srcTableSchema.tableName}`;
sql += ` FOR EACH ROW BEGIN`;
sql += ` ${content}`;
sql += ` END`;
return sql;
}