Migração para PostgreSQL multi-driver + correções de segurança
- Camada de banco unificada (src/database.js): drivers Postgres/Firebird, tradutor de SQL, suporte a schema e pool de conexões - Conexões: novo_local (Postgres externo) e firebird_local (legado) - Tela de rotas da API redesenhada (auth, params, exemplos de body) - Correções de segurança (críticos/altos/médios/baixos): XSS no chat, escalonamento de privilégio, mídia autenticada, SQL restrito a gerente, JWT sem fallback + issuer, IDOR em conversas, CORS por allowlist, rate-limit no login, limites de corpo por rota - Deploy alinhado: install.sh grava .env com PG_*, migracoes.js driver-aware Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
This commit is contained in:
+682
@@ -0,0 +1,682 @@
|
||||
/**
|
||||
* CAMADA DE BANCO DE DADOS (arquivo único)
|
||||
* ========================================
|
||||
* Concentra tudo que diz respeito a banco de dados em UM lugar:
|
||||
*
|
||||
* 1) CONFIGURAÇÃO — aliases, drivers, conexões estáticas e customizadas
|
||||
* 2) DRIVER FIREBIRD
|
||||
* 3) DRIVER POSTGRES (pool + tradutor de SQL + schema)
|
||||
* 4) DISPATCHER — API pública usada por todo o sistema
|
||||
*
|
||||
* Suporta múltiplos drivers ('postgres' padrão, 'firebird' legado). Cada
|
||||
* conexão declara o campo `driver`. Conexões adicionadas em runtime (via API
|
||||
* /api/databases) são persistidas em databases_custom.json (ao lado deste
|
||||
* arquivo) — é um arquivo de DADOS, não de código.
|
||||
*/
|
||||
const path = require('path');
|
||||
const fs = require('fs');
|
||||
|
||||
// ============================================================
|
||||
// 1) CONFIGURAÇÃO
|
||||
// ============================================================
|
||||
const CUSTOM_DB_PATH = path.resolve(__dirname, 'databases_custom.json');
|
||||
|
||||
// Conexões dinâmicas salvas em arquivo JSON
|
||||
var customDatabases = {};
|
||||
try {
|
||||
if (fs.existsSync(CUSTOM_DB_PATH)) {
|
||||
customDatabases = JSON.parse(fs.readFileSync(CUSTOM_DB_PATH, 'utf8'));
|
||||
}
|
||||
} catch (e) {
|
||||
console.error('[Databases] Erro ao carregar databases_custom.json:', e.message);
|
||||
}
|
||||
|
||||
const DRIVERS = ['postgres', 'firebird'];
|
||||
const DEFAULT_DRIVER = (process.env.DB_DRIVER || 'postgres').toLowerCase();
|
||||
|
||||
// Valores padrão por driver
|
||||
const DRIVER_DEFAULTS = {
|
||||
postgres: {
|
||||
host: '127.0.0.1',
|
||||
port: 5432,
|
||||
user: 'postgres',
|
||||
password: 'postgres',
|
||||
schema: 'public',
|
||||
ssl: false,
|
||||
max: 10,
|
||||
idleTimeoutMillis: 30000,
|
||||
connectionTimeoutMillis: 10000,
|
||||
},
|
||||
firebird: {
|
||||
host: 'localhost',
|
||||
port: 3050,
|
||||
user: 'SYSDBA',
|
||||
password: 'masterkey',
|
||||
encoding: 'UTF-8',
|
||||
lowercase_keys: false,
|
||||
pageSize: 4096,
|
||||
wireCrypt: 1, // WIRE_CRYPT_ENABLE
|
||||
},
|
||||
};
|
||||
|
||||
/**
|
||||
* Conexões estáticas.
|
||||
* - `novo_local` → PostgreSQL (conexão principal)
|
||||
* - `firebird_local`→ Firebird (informe o CAMINHO do arquivo .FDB em `database`)
|
||||
* Sobrescrevíveis pelo .env (PG_* para Postgres, DB_* para Firebird).
|
||||
*/
|
||||
const databases = {
|
||||
novo_local: {
|
||||
driver: 'postgres',
|
||||
host: process.env.PG_HOST || '127.0.0.1',
|
||||
port: parseInt(process.env.PG_PORT, 10) || 15433,
|
||||
user: process.env.PG_USER || 'postgres',
|
||||
password: process.env.PG_PASSWORD || 'postgres',
|
||||
database: process.env.PG_DATABASE || 'novo_local',
|
||||
schema: process.env.PG_SCHEMA || 'public',
|
||||
},
|
||||
|
||||
firebird_local: {
|
||||
driver: 'firebird',
|
||||
host: process.env.DB_HOST || 'localhost',
|
||||
port: parseInt(process.env.DB_PORT, 10) || 3050,
|
||||
// CAMINHO do arquivo .FDB: defina DB_DATABASE no .env (absoluto) ou ajuste
|
||||
// o path.resolve abaixo. Ex.: path.resolve(__dirname, '../db/NOVO.FDB').
|
||||
database: process.env.DB_DATABASE || path.resolve(__dirname, '../NOVO.FDB'),
|
||||
user: process.env.DB_USER || 'SYSDBA',
|
||||
password: process.env.DB_PASSWORD || 'masterkey',
|
||||
encoding: process.env.DB_ENCODING || 'UTF-8',
|
||||
lowercase_keys: false,
|
||||
pageSize: 4096,
|
||||
wireCrypt: 1, // WIRE_CRYPT_ENABLE
|
||||
},
|
||||
};
|
||||
|
||||
function salvarCustomDatabases() {
|
||||
try {
|
||||
fs.writeFileSync(CUSTOM_DB_PATH, JSON.stringify(customDatabases, null, 2), 'utf8');
|
||||
return true;
|
||||
} catch (e) {
|
||||
console.error('[Databases] Erro ao salvar databases_custom.json:', e.message);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
/** Normaliza uma config bruta aplicando o driver e seus defaults. */
|
||||
function normalize(raw) {
|
||||
const driver = DRIVERS.includes((raw.driver || '').toLowerCase())
|
||||
? raw.driver.toLowerCase()
|
||||
: DEFAULT_DRIVER;
|
||||
const d = DRIVER_DEFAULTS[driver];
|
||||
|
||||
if (driver === 'postgres') {
|
||||
return {
|
||||
driver,
|
||||
host: raw.host || d.host,
|
||||
port: raw.port || d.port,
|
||||
database: raw.database,
|
||||
user: raw.user || d.user,
|
||||
password: raw.password || d.password,
|
||||
schema: (raw.schema || d.schema || 'public'),
|
||||
ssl: raw.ssl !== undefined ? raw.ssl : d.ssl,
|
||||
max: raw.max || d.max,
|
||||
idleTimeoutMillis: raw.idleTimeoutMillis || d.idleTimeoutMillis,
|
||||
connectionTimeoutMillis: raw.connectionTimeoutMillis || d.connectionTimeoutMillis,
|
||||
};
|
||||
}
|
||||
|
||||
// firebird
|
||||
return {
|
||||
driver,
|
||||
host: raw.host || d.host,
|
||||
port: raw.port || d.port,
|
||||
database: raw.database,
|
||||
user: raw.user || d.user,
|
||||
password: raw.password || d.password,
|
||||
encoding: raw.encoding || d.encoding,
|
||||
lowercase_keys: raw.lowercase_keys || d.lowercase_keys,
|
||||
role: raw.role || null,
|
||||
pageSize: raw.pageSize || d.pageSize,
|
||||
wireCrypt: raw.wireCrypt !== undefined ? raw.wireCrypt : d.wireCrypt,
|
||||
};
|
||||
}
|
||||
|
||||
/** Configuração normalizada (inclui `driver`) de um alias. */
|
||||
function getConfig(alias) {
|
||||
if (!alias) throw new Error('Alias do banco não informado.');
|
||||
const aliasLower = alias.toLowerCase();
|
||||
const raw = customDatabases[aliasLower] || databases[aliasLower];
|
||||
if (!raw) {
|
||||
throw new Error(`Alias "${alias}" não encontrado. Aliases disponíveis: ${listAliases().join(', ')}`);
|
||||
}
|
||||
return normalize(raw);
|
||||
}
|
||||
|
||||
/** Lista de aliases disponíveis. */
|
||||
function listAliases() {
|
||||
return Object.keys(databases).concat(
|
||||
Object.keys(customDatabases).filter((a) => !databases[a])
|
||||
);
|
||||
}
|
||||
|
||||
/** Todos os aliases com suas configs (custom + estáticas). */
|
||||
function listAllDatabases() {
|
||||
var result = {};
|
||||
Object.keys(databases).forEach(function (alias) {
|
||||
result[alias] = Object.assign({}, databases[alias], { _tipo: 'estatico' });
|
||||
});
|
||||
Object.keys(customDatabases).forEach(function (alias) {
|
||||
result[alias] = Object.assign({}, customDatabases[alias], { _tipo: 'custom' });
|
||||
});
|
||||
return result;
|
||||
}
|
||||
|
||||
/** Adiciona/atualiza uma conexão customizada. */
|
||||
function addDatabase(alias, config) {
|
||||
if (!alias || !config || !config.database) {
|
||||
throw new Error('Alias e database são obrigatórios.');
|
||||
}
|
||||
var aliasLower = alias.toLowerCase().replace(/[^a-z0-9_]/g, '_');
|
||||
var driver = DRIVERS.includes((config.driver || '').toLowerCase())
|
||||
? config.driver.toLowerCase()
|
||||
: DEFAULT_DRIVER;
|
||||
|
||||
var stored = { driver: driver };
|
||||
['host', 'port', 'database', 'schema', 'user', 'password', 'ssl',
|
||||
'encoding', 'lowercase_keys', 'role', 'pageSize', 'wireCrypt',
|
||||
'max', 'idleTimeoutMillis', 'connectionTimeoutMillis']
|
||||
.forEach(function (k) {
|
||||
if (config[k] !== undefined && config[k] !== null && config[k] !== '') {
|
||||
stored[k] = config[k];
|
||||
}
|
||||
});
|
||||
|
||||
customDatabases[aliasLower] = stored;
|
||||
return salvarCustomDatabases();
|
||||
}
|
||||
|
||||
/** Remove uma conexão customizada. */
|
||||
function removeDatabase(alias) {
|
||||
if (!alias) return false;
|
||||
var aliasLower = alias.toLowerCase();
|
||||
if (!customDatabases[aliasLower]) return false;
|
||||
delete customDatabases[aliasLower];
|
||||
return salvarCustomDatabases();
|
||||
}
|
||||
|
||||
// ============================================================
|
||||
// 2) DRIVER FIREBIRD
|
||||
// ============================================================
|
||||
const firebirdDriver = (() => {
|
||||
const Firebird = require('node-firebird');
|
||||
|
||||
function readBlob(blobFunc) {
|
||||
return new Promise(function (resolve, reject) {
|
||||
if (typeof blobFunc !== 'function') return resolve(blobFunc);
|
||||
blobFunc(function (err, name, emitter) {
|
||||
if (err) return reject(err);
|
||||
if (!emitter || typeof emitter.on !== 'function') return resolve(null);
|
||||
var chunks = [];
|
||||
var total = 0;
|
||||
emitter.on('data', function (c) { chunks.push(c); total += c.length; });
|
||||
emitter.on('end', function () { resolve(Buffer.concat(chunks, total)); });
|
||||
emitter.on('error', reject);
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
function query(config, sql, params = []) {
|
||||
return new Promise((resolve, reject) => {
|
||||
Firebird.attach(config, (err, db) => {
|
||||
if (err) return reject(new Error(`Erro ao conectar (firebird): ${err.message}`));
|
||||
var allRows = [];
|
||||
db.sequentially(sql, params, function (row) {
|
||||
var keys = Object.keys(row);
|
||||
var blobPromises = keys.map(function (key) {
|
||||
var val = row[key];
|
||||
if (typeof val === 'function') {
|
||||
return readBlob(val).then(function (data) { row[key] = data; });
|
||||
}
|
||||
return Promise.resolve();
|
||||
});
|
||||
return Promise.all(blobPromises).then(function () { allRows.push(row); });
|
||||
}, function (queryErr) {
|
||||
db.detach();
|
||||
if (queryErr) return reject(new Error(`Erro na consulta (firebird): ${queryErr.message}`));
|
||||
resolve(allRows);
|
||||
});
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
function execute(config, sql, params = []) {
|
||||
return new Promise((resolve, reject) => {
|
||||
Firebird.attach(config, (err, db) => {
|
||||
if (err) return reject(new Error(`Erro ao conectar (firebird): ${err.message}`));
|
||||
db.transaction(Firebird.ISOLATION_READ_COMMITTED, (transErr, transaction) => {
|
||||
if (transErr) {
|
||||
db.detach();
|
||||
return reject(new Error(`Erro ao iniciar transação (firebird): ${transErr.message}`));
|
||||
}
|
||||
transaction.query(sql, params, (queryErr, result) => {
|
||||
if (queryErr) {
|
||||
transaction.rollback();
|
||||
db.detach();
|
||||
return reject(new Error(`Erro na execução (firebird): ${queryErr.message}`));
|
||||
}
|
||||
transaction.commit((commitErr) => {
|
||||
db.detach();
|
||||
if (commitErr) return reject(new Error(`Erro ao commitar (firebird): ${commitErr.message}`));
|
||||
resolve({ affectedRows: result ? result.length : 0, result });
|
||||
});
|
||||
});
|
||||
});
|
||||
});
|
||||
});
|
||||
}
|
||||
|
||||
async function testConnection(config) {
|
||||
await query(config, 'SELECT 1 FROM RDB$DATABASE');
|
||||
return true;
|
||||
}
|
||||
|
||||
async function listTables(config) {
|
||||
const rows = await query(config, `
|
||||
SELECT TRIM(RDB$RELATION_NAME) AS TABLE_NAME
|
||||
FROM RDB$RELATIONS
|
||||
WHERE RDB$SYSTEM_FLAG = 0 AND RDB$RELATION_TYPE = 0
|
||||
ORDER BY RDB$RELATION_NAME
|
||||
`);
|
||||
return rows.map((r) => (r.TABLE_NAME || '').trim()).filter(Boolean);
|
||||
}
|
||||
|
||||
async function tableInfo(config, tableName) {
|
||||
const rows = await query(config, `
|
||||
SELECT
|
||||
rf.RDB$FIELD_NAME AS COLUMN_NAME,
|
||||
rf.RDB$FIELD_POSITION AS ORDINAL_POSITION,
|
||||
CASE f.RDB$FIELD_TYPE
|
||||
WHEN 7 THEN 'SMALLINT' WHEN 8 THEN 'INTEGER' WHEN 16 THEN 'BIGINT'
|
||||
WHEN 9 THEN 'QUAD' WHEN 10 THEN 'FLOAT' WHEN 27 THEN 'DOUBLE PRECISION'
|
||||
WHEN 12 THEN 'DATE' WHEN 13 THEN 'TIME' WHEN 35 THEN 'TIMESTAMP'
|
||||
WHEN 37 THEN 'VARCHAR' WHEN 40 THEN 'CSTRING' WHEN 45 THEN 'BLOB_ID'
|
||||
WHEN 261 THEN 'BLOB' WHEN 14 THEN 'CHAR' WHEN 41 THEN 'NUMERIC'
|
||||
ELSE 'UNKNOWN'
|
||||
END AS DATA_TYPE,
|
||||
f.RDB$FIELD_LENGTH AS FIELD_LENGTH,
|
||||
f.RDB$FIELD_SCALE AS FIELD_SCALE,
|
||||
f.RDB$FIELD_PRECISION AS FIELD_PRECISION,
|
||||
rf.RDB$NULL_FLAG AS NULL_FLAG
|
||||
FROM RDB$RELATION_FIELDS rf
|
||||
INNER JOIN RDB$FIELDS f ON rf.RDB$FIELD_SOURCE = f.RDB$FIELD_NAME
|
||||
WHERE rf.RDB$RELATION_NAME = ?
|
||||
ORDER BY rf.RDB$FIELD_POSITION
|
||||
`, [String(tableName).toUpperCase()]);
|
||||
|
||||
return rows.map((row) => ({
|
||||
name: (row.COLUMN_NAME || '').trim(),
|
||||
position: row.ORDINAL_POSITION,
|
||||
type: (row.DATA_TYPE || '').trim(),
|
||||
length: row.FIELD_LENGTH,
|
||||
precision: row.FIELD_PRECISION,
|
||||
scale: row.FIELD_SCALE,
|
||||
nullable: row.NULL_FLAG !== 1,
|
||||
}));
|
||||
}
|
||||
|
||||
async function tableExists(config, tableName) {
|
||||
const r = await query(config,
|
||||
"SELECT COUNT(*) AS CT FROM RDB$RELATIONS WHERE RDB$RELATION_NAME = ?",
|
||||
[String(tableName).toUpperCase()]);
|
||||
return (r[0] && r[0].CT > 0) || false;
|
||||
}
|
||||
|
||||
async function columnExists(config, tableName, columnName) {
|
||||
const r = await query(config,
|
||||
"SELECT COUNT(*) AS CT FROM RDB$RELATION_FIELDS WHERE RDB$RELATION_NAME = ? AND RDB$FIELD_NAME = ?",
|
||||
[String(tableName).toUpperCase(), String(columnName).toUpperCase()]);
|
||||
return (r[0] && r[0].CT > 0) || false;
|
||||
}
|
||||
|
||||
async function close() { /* Firebird abre conexão por consulta; nada a encerrar */ }
|
||||
|
||||
return { query, execute, testConnection, listTables, tableInfo, tableExists, columnExists, close };
|
||||
})();
|
||||
|
||||
// ============================================================
|
||||
// 3) DRIVER POSTGRES (pool + tradutor de SQL Firebird->PG + schema)
|
||||
// ============================================================
|
||||
const postgresDriver = (() => {
|
||||
const pg = require('pg');
|
||||
|
||||
// COUNT()/bigint chegam como string no pg; o código histórico usa números.
|
||||
pg.types.setTypeParser(20, (v) => (v === null ? null : parseInt(v, 10))); // int8 / bigint
|
||||
pg.types.setTypeParser(1700, (v) => (v === null ? null : parseFloat(v))); // numeric / decimal
|
||||
|
||||
const pools = new Map();
|
||||
|
||||
function safeSchema(schema) {
|
||||
const s = String(schema || 'public').trim();
|
||||
return /^[A-Za-z_][A-Za-z0-9_$]*$/.test(s) ? s : 'public';
|
||||
}
|
||||
|
||||
function poolKey(config) {
|
||||
return [config.host, config.port, config.database, config.user, safeSchema(config.schema)].join('|');
|
||||
}
|
||||
|
||||
// Schemas que o pool enxerga (configurado + public como fallback)
|
||||
function schemasFor(config) {
|
||||
const schema = safeSchema(config.schema);
|
||||
return schema === 'public' ? ['public'] : [schema, 'public'];
|
||||
}
|
||||
|
||||
function getEntry(config) {
|
||||
const key = poolKey(config);
|
||||
let entry = pools.get(key);
|
||||
if (!entry) {
|
||||
const schemas = schemasFor(config);
|
||||
const pool = new pg.Pool({
|
||||
host: config.host,
|
||||
port: config.port,
|
||||
database: config.database,
|
||||
user: config.user,
|
||||
password: config.password,
|
||||
ssl: config.ssl ? { rejectUnauthorized: false } : false,
|
||||
options: '-c search_path=' + schemas.join(','),
|
||||
max: config.max || 10,
|
||||
idleTimeoutMillis: config.idleTimeoutMillis || 30000,
|
||||
connectionTimeoutMillis: config.connectionTimeoutMillis || 10000,
|
||||
});
|
||||
pool.on('error', (err) => console.error('[postgres] Erro inesperado no pool:', err.message));
|
||||
entry = { pool, schemas, tableSet: null, loadingTables: null };
|
||||
pools.set(key, entry);
|
||||
}
|
||||
return entry;
|
||||
}
|
||||
|
||||
// Carrega (uma vez por pool) os nomes de tabela MAIÚSCULOS p/ decidir aspas
|
||||
async function ensureTableSet(entry) {
|
||||
if (entry.tableSet) return entry.tableSet;
|
||||
if (!entry.loadingTables) {
|
||||
entry.loadingTables = entry.pool
|
||||
.query('SELECT table_name FROM information_schema.tables WHERE table_schema = ANY($1)', [entry.schemas])
|
||||
.then((res) => {
|
||||
entry.tableSet = new Set(res.rows.map((r) => String(r.table_name).toUpperCase()));
|
||||
return entry.tableSet;
|
||||
})
|
||||
.catch((err) => {
|
||||
console.error('[postgres] Falha ao carregar nomes de tabela:', err.message);
|
||||
entry.tableSet = new Set();
|
||||
return entry.tableSet;
|
||||
});
|
||||
}
|
||||
return entry.loadingTables;
|
||||
}
|
||||
|
||||
// ? -> $1, $2, ... (ignora literais de string)
|
||||
function convertPlaceholders(sql) {
|
||||
let out = '', i = 0, n = 1, inStr = false;
|
||||
while (i < sql.length) {
|
||||
const ch = sql[i];
|
||||
if (inStr) {
|
||||
out += ch;
|
||||
if (ch === "'") {
|
||||
if (sql[i + 1] === "'") { out += sql[i + 1]; i += 2; continue; }
|
||||
inStr = false;
|
||||
}
|
||||
i++;
|
||||
continue;
|
||||
}
|
||||
if (ch === "'") { inStr = true; out += ch; i++; continue; }
|
||||
if (ch === '?') { out += '$' + (n++); i++; continue; }
|
||||
out += ch; i++;
|
||||
}
|
||||
return out;
|
||||
}
|
||||
|
||||
// Aspas nos nomes de tabela conhecidos (MAIÚSCULOS) após FROM/JOIN/INTO/UPDATE/ALTER/DROP
|
||||
function quoteTableNames(sql, tableSet) {
|
||||
if (!tableSet || tableSet.size === 0) return sql;
|
||||
return sql.replace(
|
||||
/(\b(?:FROM|JOIN|INTO|UPDATE|ALTER\s+TABLE|DROP\s+TABLE)\s+)("?)([A-Za-z_][A-Za-z0-9_$]*)("?)/gi,
|
||||
(match, kw, q1, name, q2) => {
|
||||
if (q1 === '"' || q2 === '"') return match;
|
||||
if (tableSet.has(name.toUpperCase())) return kw + '"' + name.toUpperCase() + '"';
|
||||
return match;
|
||||
}
|
||||
);
|
||||
}
|
||||
|
||||
// Rede de segurança: FIRST n [SKIP m] no topo -> LIMIT/OFFSET (SQL legado)
|
||||
function translateFirstSkip(sql) {
|
||||
return sql.replace(
|
||||
/^(\s*SELECT\s+)FIRST\s+(\d+)(?:\s+SKIP\s+(\d+))?\s+/i,
|
||||
(m, sel, first, skip) => sel + ' LIMITTAIL LIMIT ' + first + (skip ? ' OFFSET ' + skip : '') + ' '
|
||||
);
|
||||
}
|
||||
function applyLimitTail(sql) {
|
||||
const marker = / LIMITTAIL( LIMIT \d+(?: OFFSET \d+)?) /;
|
||||
const m = sql.match(marker);
|
||||
if (!m) return sql;
|
||||
return sql.replace(marker, ' ').trimEnd() + m[1];
|
||||
}
|
||||
|
||||
function translateSql(sql, tableSet) {
|
||||
let out = sql;
|
||||
out = out.replace(/\bCONTAINING\s+(\?|\$\d+|'(?:[^']|'')*')/gi, "ILIKE ('%' || $1 || '%')");
|
||||
out = out.replace(/\bFROM\s+RDB\$DATABASE\b/gi, '');
|
||||
out = quoteTableNames(out, tableSet);
|
||||
out = translateFirstSkip(out);
|
||||
out = applyLimitTail(out);
|
||||
out = convertPlaceholders(out);
|
||||
return out;
|
||||
}
|
||||
|
||||
function upperKeys(rows) {
|
||||
return rows.map((row) => {
|
||||
const o = {};
|
||||
for (const k in row) o[k.toUpperCase()] = row[k];
|
||||
return o;
|
||||
});
|
||||
}
|
||||
|
||||
async function query(config, sql, params = []) {
|
||||
const entry = getEntry(config);
|
||||
await ensureTableSet(entry);
|
||||
const text = translateSql(sql, entry.tableSet);
|
||||
try {
|
||||
const res = await entry.pool.query(text, params);
|
||||
return upperKeys(res.rows);
|
||||
} catch (err) {
|
||||
throw new Error(`Erro na consulta (postgres): ${err.message}`);
|
||||
}
|
||||
}
|
||||
|
||||
async function execute(config, sql, params = []) {
|
||||
const entry = getEntry(config);
|
||||
await ensureTableSet(entry);
|
||||
const text = translateSql(sql, entry.tableSet);
|
||||
try {
|
||||
const res = await entry.pool.query(text, params);
|
||||
return { affectedRows: res.rowCount || 0, result: upperKeys(res.rows || []) };
|
||||
} catch (err) {
|
||||
throw new Error(`Erro na execução (postgres): ${err.message}`);
|
||||
}
|
||||
}
|
||||
|
||||
async function testConnection(config) {
|
||||
const entry = getEntry(config);
|
||||
await entry.pool.query('SELECT 1');
|
||||
return true;
|
||||
}
|
||||
|
||||
async function listTables(config) {
|
||||
const entry = getEntry(config);
|
||||
const res = await entry.pool.query(
|
||||
`SELECT table_name FROM information_schema.tables
|
||||
WHERE table_schema = $1 AND table_type = 'BASE TABLE'
|
||||
ORDER BY table_name`,
|
||||
[safeSchema(config.schema)]
|
||||
);
|
||||
return res.rows.map((r) => r.table_name);
|
||||
}
|
||||
|
||||
async function tableInfo(config, tableName) {
|
||||
const entry = getEntry(config);
|
||||
// Resolve no primeiro schema do search_path que contém a tabela
|
||||
const res = await entry.pool.query(
|
||||
`SELECT column_name, ordinal_position, data_type,
|
||||
character_maximum_length, numeric_precision, numeric_scale, is_nullable
|
||||
FROM information_schema.columns
|
||||
WHERE UPPER(table_name) = UPPER($2)
|
||||
AND table_schema = (
|
||||
SELECT table_schema FROM information_schema.tables
|
||||
WHERE UPPER(table_name) = UPPER($2) AND table_schema = ANY($1)
|
||||
ORDER BY array_position($1, table_schema) LIMIT 1
|
||||
)
|
||||
ORDER BY ordinal_position`,
|
||||
[entry.schemas, tableName]
|
||||
);
|
||||
return res.rows.map((row) => ({
|
||||
name: row.column_name,
|
||||
position: row.ordinal_position,
|
||||
type: (row.data_type || '').toUpperCase(),
|
||||
length: row.character_maximum_length,
|
||||
precision: row.numeric_precision,
|
||||
scale: row.numeric_scale,
|
||||
nullable: row.is_nullable === 'YES',
|
||||
}));
|
||||
}
|
||||
|
||||
async function tableExists(config, tableName) {
|
||||
const entry = getEntry(config);
|
||||
await ensureTableSet(entry);
|
||||
if (entry.tableSet) return entry.tableSet.has(String(tableName).toUpperCase());
|
||||
const res = await entry.pool.query(
|
||||
`SELECT 1 FROM information_schema.tables
|
||||
WHERE table_schema = ANY($1) AND UPPER(table_name) = UPPER($2) LIMIT 1`,
|
||||
[entry.schemas, tableName]
|
||||
);
|
||||
return res.rowCount > 0;
|
||||
}
|
||||
|
||||
async function columnExists(config, tableName, columnName) {
|
||||
const entry = getEntry(config);
|
||||
const res = await entry.pool.query(
|
||||
`SELECT 1 FROM information_schema.columns
|
||||
WHERE table_schema = ANY($1) AND UPPER(table_name) = UPPER($2)
|
||||
AND UPPER(column_name) = UPPER($3) LIMIT 1`,
|
||||
[entry.schemas, tableName, columnName]
|
||||
);
|
||||
return res.rowCount > 0;
|
||||
}
|
||||
|
||||
async function close() {
|
||||
const all = Array.from(pools.values()).map((e) => e.pool.end().catch(() => {}));
|
||||
pools.clear();
|
||||
await Promise.all(all);
|
||||
}
|
||||
|
||||
return {
|
||||
query, execute, testConnection, listTables, tableInfo, tableExists, columnExists, close,
|
||||
_translateSql: translateSql, // exportado para testes
|
||||
};
|
||||
})();
|
||||
|
||||
// ============================================================
|
||||
// 4) DISPATCHER — API pública
|
||||
// ============================================================
|
||||
const drivers = { postgres: postgresDriver, firebird: firebirdDriver };
|
||||
|
||||
function getDriver(alias) {
|
||||
const config = getConfig(alias);
|
||||
const driver = drivers[config.driver];
|
||||
if (!driver) {
|
||||
throw new Error(`Driver "${config.driver}" não suportado para o alias "${alias}".`);
|
||||
}
|
||||
return { config, driver };
|
||||
}
|
||||
|
||||
function query(alias, sql, params = []) {
|
||||
let d;
|
||||
try { d = getDriver(alias); } catch (err) { return Promise.reject(err); }
|
||||
return d.driver.query(d.config, sql, params).catch((err) => { throw new Error(`[${alias}] ${err.message}`); });
|
||||
}
|
||||
|
||||
function execute(alias, sql, params = []) {
|
||||
let d;
|
||||
try { d = getDriver(alias); } catch (err) { return Promise.reject(err); }
|
||||
return d.driver.execute(d.config, sql, params).catch((err) => { throw new Error(`[${alias}] ${err.message}`); });
|
||||
}
|
||||
|
||||
async function testConnection(alias) {
|
||||
try {
|
||||
const { config, driver } = getDriver(alias);
|
||||
await driver.testConnection(config);
|
||||
return true;
|
||||
} catch (err) {
|
||||
console.error(`[${alias}] Falha na conexão:`, err.message);
|
||||
return false;
|
||||
}
|
||||
}
|
||||
|
||||
async function testAllConnections() {
|
||||
const results = {};
|
||||
for (const alias of listAliases()) {
|
||||
results[alias] = await testConnection(alias);
|
||||
}
|
||||
return results;
|
||||
}
|
||||
|
||||
function listTables(alias) {
|
||||
const { config, driver } = getDriver(alias);
|
||||
return driver.listTables(config);
|
||||
}
|
||||
|
||||
function tableInfo(alias, tableName) {
|
||||
const { config, driver } = getDriver(alias);
|
||||
return driver.tableInfo(config, tableName);
|
||||
}
|
||||
|
||||
function tableExists(alias, tableName) {
|
||||
const { config, driver } = getDriver(alias);
|
||||
return driver.tableExists(config, tableName);
|
||||
}
|
||||
|
||||
function columnExists(alias, tableName, columnName) {
|
||||
const { config, driver } = getDriver(alias);
|
||||
return driver.columnExists(config, tableName, columnName);
|
||||
}
|
||||
|
||||
/** Nome do driver de um alias (ex.: 'postgres'). */
|
||||
function driverOf(alias) {
|
||||
return getConfig(alias).driver;
|
||||
}
|
||||
|
||||
/** Encerra os pools de todos os drivers (shutdown gracioso). */
|
||||
async function closeAll() {
|
||||
await Promise.all(Object.values(drivers).map((d) => d.close && d.close()));
|
||||
}
|
||||
|
||||
module.exports = {
|
||||
// Acesso a dados
|
||||
query,
|
||||
execute,
|
||||
testConnection,
|
||||
testAllConnections,
|
||||
listTables,
|
||||
tableInfo,
|
||||
tableExists,
|
||||
columnExists,
|
||||
driverOf,
|
||||
closeAll,
|
||||
// Configuração / aliases
|
||||
databases,
|
||||
DRIVERS,
|
||||
DEFAULT_DRIVER,
|
||||
getConfig,
|
||||
listAliases,
|
||||
listAllDatabases,
|
||||
addDatabase,
|
||||
removeDatabase,
|
||||
};
|
||||
Reference in New Issue
Block a user