180 lines
7.2 KiB
JavaScript
180 lines
7.2 KiB
JavaScript
const { Pool } = require('pg');
|
|
const deasync = require('deasync');
|
|
const SqliteDatabase = require('better-sqlite3');
|
|
|
|
const pools = new Map();
|
|
|
|
function isPostgresEnabled() {
|
|
return String(process.env.DURIIN_DB_BACKEND || '').toLowerCase() === 'postgres' || Boolean((process.env.DURIIN_POSTGRES_URL || process.env.DATABASE_URL) && process.env.DURIIN_USE_POSTGRES === 'true');
|
|
}
|
|
|
|
function poolFor(schema) {
|
|
const connectionString = process.env.DURIIN_POSTGRES_URL || process.env.DATABASE_URL;
|
|
if (!connectionString) throw new Error('DURIIN_POSTGRES_URL is required for postgres runtime');
|
|
const key = `${connectionString}|${schema}`;
|
|
if (!pools.has(key)) {
|
|
pools.set(key, new Pool({
|
|
connectionString,
|
|
max: Math.max(1, Number(process.env.POSTGRES_RUNTIME_POOL_SIZE) || 4),
|
|
options: `-c search_path=${schema},public`,
|
|
}));
|
|
}
|
|
return pools.get(key);
|
|
}
|
|
|
|
function querySync(pool, sql, params = []) {
|
|
let done = false;
|
|
let result;
|
|
let error;
|
|
pool.query(sql, params).then((value) => { result = value; done = true; }).catch((err) => { error = err; done = true; });
|
|
deasync.loopWhile(() => !done);
|
|
if (error) throw error;
|
|
return result;
|
|
}
|
|
|
|
function connectSync(pool) {
|
|
let done = false;
|
|
let client;
|
|
let error;
|
|
pool.connect().then((value) => { client = value; done = true; }).catch((err) => { error = err; done = true; });
|
|
deasync.loopWhile(() => !done);
|
|
if (error) throw error;
|
|
return client;
|
|
}
|
|
|
|
function normalizeParams(params) {
|
|
if (params.length === 1 && params[0] && typeof params[0] === 'object' && !Array.isArray(params[0]) && !Buffer.isBuffer(params[0])) {
|
|
return params[0];
|
|
}
|
|
return params.flat();
|
|
}
|
|
|
|
function rewritePlaceholders(sql, params) {
|
|
if (params && !Array.isArray(params)) {
|
|
const values = [];
|
|
const text = sql.replace(/@([A-Za-z_][A-Za-z0-9_]*)/g, (_, name) => {
|
|
values.push(params[name]);
|
|
return `$${values.length}`;
|
|
});
|
|
return { sql: text, params: values };
|
|
}
|
|
let index = 0;
|
|
return { sql: sql.replace(/\?/g, () => `$${++index}`), params: params || [] };
|
|
}
|
|
|
|
function rewriteSql(sql, params) {
|
|
let text = String(sql).trim();
|
|
const pragmaTable = text.match(/^PRAGMA\s+table_info\((?:"([^"]+)"|'([^']+)'|([^)]+))\)$/i);
|
|
if (pragmaTable) {
|
|
const table = String(pragmaTable[1] || pragmaTable[2] || pragmaTable[3] || '').trim();
|
|
return {
|
|
sql: `
|
|
SELECT ordinal_position - 1 AS cid,
|
|
column_name AS name,
|
|
data_type AS type,
|
|
CASE WHEN is_nullable = 'NO' THEN 1 ELSE 0 END AS notnull,
|
|
column_default AS dflt_value,
|
|
0 AS pk
|
|
FROM information_schema.columns
|
|
WHERE table_schema = current_schema() AND table_name = $1
|
|
ORDER BY ordinal_position
|
|
`,
|
|
params: [table],
|
|
};
|
|
}
|
|
text = text.replace(/INSERT\s+OR\s+IGNORE\s+INTO/gi, 'INSERT INTO');
|
|
text = text.replace(/INSERT\s+OR\s+REPLACE\s+INTO\s+autonomy_outcomes\s*\(([^)]+)\)\s*VALUES\s*\(([^)]+)\)/i,
|
|
(match, columns, values) => {
|
|
const names = columns.split(',').map((item) => item.trim().replace(/"/g, ''));
|
|
const updates = names.filter((name) => name !== 'prediction_id').map((name) => `${name}=EXCLUDED.${name}`).join(', ');
|
|
return `INSERT INTO autonomy_outcomes (${columns}) VALUES (${values}) ON CONFLICT (prediction_id) DO UPDATE SET ${updates}`;
|
|
});
|
|
text = text.replace(/AUTOINCREMENT/gi, 'GENERATED BY DEFAULT AS IDENTITY');
|
|
text = text.replace(/INTEGER\s+PRIMARY\s+KEY\s+GENERATED BY DEFAULT AS IDENTITY/gi, 'BIGINT PRIMARY KEY GENERATED BY DEFAULT AS IDENTITY');
|
|
text = text.replace(/INTEGER\s+PRIMARY\s+KEY\s+AUTOINCREMENT/gi, 'BIGINT PRIMARY KEY GENERATED BY DEFAULT AS IDENTITY');
|
|
text = text.replace(/ingested_at\s*>=\s*datetime\('now',\s*'-48 hours'\)/gi, "ingested_at >= to_char(CURRENT_TIMESTAMP - interval '48 hours', 'YYYY-MM-DD HH24:MI:SS')");
|
|
text = text.replace(/datetime\('now',\s*\?\)/gi, 'CURRENT_TIMESTAMP + (?::interval)');
|
|
text = text.replace(/datetime\('now',\s*'\+60 seconds'\)/gi, "CURRENT_TIMESTAMP + interval '60 seconds'");
|
|
text = text.replace(/datetime\('now',\s*'([^']+)'\)/gi, "CURRENT_TIMESTAMP + interval '$1'");
|
|
text = text.replace(/datetime\('now'\)/gi, 'CURRENT_TIMESTAMP');
|
|
text = text.replace(/date\('now'\)/gi, 'CURRENT_DATE');
|
|
text = text.replace(/datetime\(COALESCE\(([^)]+)\)\)/gi, 'COALESCE($1)::timestamp');
|
|
text = text.replace(/datetime\((p\.information_cutoff),\s*'\+'\s*\|\|\s*(p\.horizon_days)\s*\|\|\s*' days'\)/gi, "($1::timestamp + ($2 || ' days')::interval)");
|
|
text = text.replace(/datetime\(([^)]+)\)/gi, '($1)::timestamp');
|
|
|
|
const rewritten = rewritePlaceholders(text, params);
|
|
let finalSql = rewritten.sql;
|
|
if (/^INSERT\s+INTO\s+autonomy_jobs\b/i.test(finalSql) && !/ON\s+CONFLICT/i.test(finalSql)) finalSql += ' ON CONFLICT DO NOTHING';
|
|
if (/^INSERT\s+INTO\s+autonomy_order_intents\b/i.test(finalSql) && !/ON\s+CONFLICT/i.test(finalSql)) finalSql += ' ON CONFLICT DO NOTHING';
|
|
if (/^INSERT\s+INTO\s+autonomy_schema\b/i.test(finalSql) && !/ON\s+CONFLICT/i.test(finalSql)) finalSql += ' ON CONFLICT DO NOTHING';
|
|
if (/^INSERT\s+INTO\s+autonomy_proposals\b/i.test(finalSql) && !/RETURNING\s+id/i.test(finalSql)) finalSql += ' RETURNING id';
|
|
return { sql: finalSql, params: rewritten.params };
|
|
}
|
|
|
|
class PgCompatDb {
|
|
constructor(schema) {
|
|
this.schema = schema;
|
|
this.dialect = 'postgres';
|
|
this.pool = poolFor(schema);
|
|
}
|
|
|
|
pragma() { return undefined; }
|
|
|
|
prepare(sql) {
|
|
const db = this;
|
|
const target = () => db.activeClient || db.pool;
|
|
return {
|
|
get(...rawParams) {
|
|
const { sql: text, params } = rewriteSql(sql, normalizeParams(rawParams));
|
|
return querySync(target(), text, params).rows[0];
|
|
},
|
|
all(...rawParams) {
|
|
const { sql: text, params } = rewriteSql(sql, normalizeParams(rawParams));
|
|
return querySync(target(), text, params).rows;
|
|
},
|
|
run(...rawParams) {
|
|
const { sql: text, params } = rewriteSql(sql, normalizeParams(rawParams));
|
|
const result = querySync(target(), text, params);
|
|
return { changes: result.rowCount || 0, lastInsertRowid: result.rows?.[0]?.id ?? null };
|
|
},
|
|
};
|
|
}
|
|
|
|
exec(sql) {
|
|
const statements = String(sql).split(';').map((statement) => statement.trim()).filter(Boolean);
|
|
for (const statement of statements) {
|
|
const { sql: text, params } = rewriteSql(statement, []);
|
|
querySync(this.activeClient || this.pool, text, params);
|
|
}
|
|
}
|
|
|
|
transaction(fn) {
|
|
const db = this;
|
|
const run = (...args) => {
|
|
const client = connectSync(db.pool);
|
|
try {
|
|
db.activeClient = client;
|
|
querySync(client, 'BEGIN', []);
|
|
const result = fn(...args);
|
|
querySync(client, 'COMMIT', []);
|
|
return result;
|
|
} catch (error) {
|
|
try { querySync(client, 'ROLLBACK', []); } catch (_) {}
|
|
throw error;
|
|
} finally {
|
|
db.activeClient = null;
|
|
client.release();
|
|
}
|
|
};
|
|
run.immediate = run;
|
|
return run;
|
|
}
|
|
}
|
|
|
|
function openRuntimeDb(path, { schema = 'intelligence', readonly = false } = {}) {
|
|
if (isPostgresEnabled()) return new PgCompatDb(schema);
|
|
return new SqliteDatabase(path, readonly ? { readonly: true } : undefined);
|
|
}
|
|
|
|
module.exports = { openRuntimeDb, isPostgresEnabled, PgCompatDb, rewriteSql };
|