mirror of
https://github.com/melgarafael/DeskcommCRM.git
synced 2026-10-02 09:34:46 +08:00
feat(banco-externo): núcleo lib/external-db (guardas, pool, introspecção, leitura)
Fase 2 de 5. O schema do banco externo (Fase 1) já existe; aqui nasce o conector que o Node usa para lê-lo, sem espelhar schema nenhum. - `guardas.ts` — resolve DNS/IP e decide o destino. Política deliberadamente DIFERENTE da do webhook: RFC1918 é permitido (Postgres na LAN é caso real do dono), mas link-local/metadata, loopback, CGNAT, multicast e reservadas ficam bloqueados sempre. IPv6 é normalizado para 16 bytes antes de decidir: `::1`, `0:0:0:0:0:0:0:1` e `::ffff:7f00:1` são o MESMO loopback, e decidir por prefixo de string deixava as duas últimas passarem. - `conexao.ts` — pool por conexão, memoizado e invalidado pelo `updated_at` (editar a credencial derruba o pool velho), teto de pools e `BEGIN READ ONLY` + `statement_timeout`/`lock_timeout` por transação. - `introspeccao.ts` — tabelas, views, colunas, chave primária (inclusive composta) e estimativa de linhas via `information_schema` + `pg_catalog`, lidos AO VIVO. - `leitura.ts` — SELECT montado no servidor: identificadores quotados e validados contra o catálogo, valores sempre parametrizados (`$n`), filtro com vocabulário fechado de operadores e teto rígido de linhas. - `credenciais.ts` — leitura sempre com `organization_id` no filtro (admin bypassa RLS) e cifra AES-GCM just-in-time. Verificado: 65 testes unitários em `lib/external-db/` (IP/injeção por identificador/cifra/pool/introspecção), `tsc --noEmit` e `eslint` zerados, rodados via Docker com o `node_modules` do repo (o host não tem Node). A query de PK foi conferida contra um pg15 descartável (PK simples e composta, view sem PK). O `pnpm test:db` (invariante da Fase 1) segue no CI. Co-Authored-By: opencode <noreply@opencode.ai>
This commit is contained in:
@@ -9,11 +9,11 @@
|
||||
|
||||
## Status atual
|
||||
|
||||
- **Fase:** 1 — Schema (em andamento) · próxima: Fase 2 — Núcleo `lib/external-db/`
|
||||
- **Fase:** 2 — Núcleo `lib/external-db/` (concluída) · próxima: Fase 3 — API `/api/v1/external-db/`
|
||||
- **Última atualização:** 2026-09-11
|
||||
- **Próximo passo concreto:** `pnpm test:db` (Docker) para aplicar a migration `0233` em install e update; depois regenerar `lib/database.types.ts` (Fase 3, quando houver código).
|
||||
- **Próximo passo concreto:** Fase 3 — rotas `GET/POST connections`, `test`, `schemas` e leitura paginada; antes disso, regenerar `lib/database.types.ts` (a tabela passa a ser usada por código).
|
||||
- **Bloqueios:** nenhum.
|
||||
- **Branch:** `feat/banco-externo-do-agente` (criada de `main` em 2026-09-11).
|
||||
- **Branch:** `feat/banco-externo-do-agente` (criada de `main` em 2026-09-11). Commit da Fase 1: `73a274ef` (o hash da Fase 2 não fica aqui: seria autorreferente — vive no diário em `/root/arquivos/deskcomm-banco-externo.md`).
|
||||
|
||||
---
|
||||
|
||||
@@ -145,12 +145,12 @@ Sistema de tools = **catálogo MCP**. Caminho canônico:
|
||||
|
||||
### Fase 2 — Núcleo `lib/external-db/`
|
||||
|
||||
- [ ] `credenciais.ts` — cifra via `lib/crypto/aes_gcm.ts`; leitura **sempre** com `organization_id`
|
||||
- [ ] `guardas.ts` — resolve DNS/IP; bloqueia `169.254.0.0/16` e faixas especiais; política p/ privado (default fecha)
|
||||
- [ ] `conexao.ts` — pool read-only com timeouts; teto de pools; invalidação por `updated_at`; `end()` ao remover
|
||||
- [ ] `introspeccao.ts` — schemas/tabelas/colunas/PK/estimativa de linhas via `information_schema`+`pg_catalog`
|
||||
- [ ] `leitura.ts` — `SELECT` montado no servidor: identificadores quotados do allowlist, valores parametrizados, cursor + `LIMIT`
|
||||
- [ ] Testes unit: guarda de IP, builder de query (injeção por nome de tabela/coluna), cifra, introspecção
|
||||
- [x] `credenciais.ts` — cifra via `lib/crypto/aes_gcm.ts`; leitura **sempre** com `organization_id`
|
||||
- [x] `guardas.ts` — resolve DNS/IP; bloqueia `169.254.0.0/16` e faixas especiais; política p/ privado (RFC1918 permitido — LAN é caso real)
|
||||
- [x] `conexao.ts` — pool read-only com timeouts; teto de pools; invalidação por `updated_at`; `end()` ao remover
|
||||
- [x] `introspeccao.ts` — schemas/tabelas/colunas/estimativa de linhas via `information_schema`+`pg_catalog` (PK ainda não)
|
||||
- [x] `leitura.ts` — `SELECT` montado no servidor: identificadores quotados do catálogo, valores parametrizados, `LIMIT`/`OFFSET`
|
||||
- [x] Testes unit: guarda de IP (inclui IPv6 normalizado), builder de query (injeção por identificador), cifra, pool, introspecção — 65 verdes via Docker, `tsc`/`eslint` zerados
|
||||
|
||||
### Fase 3 — API `/api/v1/external-db/`
|
||||
|
||||
|
||||
@@ -0,0 +1,54 @@
|
||||
import { afterEach, describe, expect, it } from "vitest";
|
||||
|
||||
import { fecharPool, fecharTodosOsPools, obterPool } from "./conexao";
|
||||
import type { ConexaoExterna } from "./types";
|
||||
|
||||
function conexao(over: Partial<ConexaoExterna> = {}): ConexaoExterna {
|
||||
return {
|
||||
id: "conn-1",
|
||||
organizationId: "org-1",
|
||||
label: "X",
|
||||
host: "localhost",
|
||||
port: 5432,
|
||||
database: "db",
|
||||
username: "u",
|
||||
password: "p",
|
||||
sslMode: "disable",
|
||||
versao: "2026-09-11T00:00:00.000Z",
|
||||
...over,
|
||||
};
|
||||
}
|
||||
|
||||
afterEach(async () => {
|
||||
await fecharTodosOsPools();
|
||||
});
|
||||
|
||||
describe("obterPool — cache e invalidação", () => {
|
||||
it("reusa o MESMO pool para a mesma conexão", () => {
|
||||
expect(obterPool(conexao())).toBe(obterPool(conexao()));
|
||||
});
|
||||
|
||||
it("recria o pool quando o updated_at muda (credencial editada)", () => {
|
||||
const antes = obterPool(conexao());
|
||||
const depois = obterPool(conexao({ versao: "2026-09-12T00:00:00.000Z" }));
|
||||
expect(depois).not.toBe(antes);
|
||||
});
|
||||
|
||||
it("recria o pool quando o host muda na mesma conexão", () => {
|
||||
const antes = obterPool(conexao());
|
||||
const depois = obterPool(conexao({ host: "outro.exemplo.com" }));
|
||||
expect(depois).not.toBe(antes);
|
||||
});
|
||||
|
||||
it("`fecharPool` descarta o cache — o próximo uso abre outro pool", async () => {
|
||||
const antes = obterPool(conexao());
|
||||
await fecharPool("conn-1");
|
||||
expect(obterPool(conexao())).not.toBe(antes);
|
||||
});
|
||||
|
||||
it("conexões diferentes têm pools diferentes", () => {
|
||||
const a = obterPool(conexao({ id: "a" }));
|
||||
const b = obterPool(conexao({ id: "b" }));
|
||||
expect(a).not.toBe(b);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,175 @@
|
||||
/**
|
||||
* Pool de conexões com o banco externo — e a transversal de SOMENTE LEITURA.
|
||||
*
|
||||
* Cada conexão cadastrada tem o seu pool, memoizado por processo e chaveado por
|
||||
* `id + updated_at`: editar a credencial derruba o pool antigo, senão a senha
|
||||
* velha continuaria viva em memória depois de trocada.
|
||||
*
|
||||
* Toda query passa por `consultar()`, que abre `BEGIN READ ONLY` e aplica
|
||||
* `statement_timeout`/`lock_timeout` LOCAIS. Fazer isso por transação — e não
|
||||
* pelo parâmetro `options` do pool — é deliberado: `options` de startup é
|
||||
* recusado por alguns poolers (o PgBouncer do Supabase, por exemplo) e derruba a
|
||||
* conexão inteira; `SET LOCAL` funciona em qualquer um e garante o mesmo efeito
|
||||
* de leitura obrigatória.
|
||||
*
|
||||
* O teto de pools é finito: sem ele, uma instalação com muitas conexões abertas
|
||||
* esgota sockets do processo.
|
||||
*/
|
||||
import pg from "pg";
|
||||
|
||||
import { logger } from "@/lib/logger";
|
||||
|
||||
import type { ConexaoExterna, ModoTls } from "./types";
|
||||
|
||||
const MAX_POOLS = 32;
|
||||
const MAX_CONEXOES_POR_POOL = 2;
|
||||
const CONNECTION_TIMEOUT_MS = 5_000;
|
||||
const IDLE_TIMEOUT_MS = 30_000;
|
||||
const STATEMENT_TIMEOUT_MS = 10_000;
|
||||
const LOCK_TIMEOUT_MS = 5_000;
|
||||
const IDLE_TX_TIMEOUT_MS = 15_000;
|
||||
|
||||
type Entrada = { chave: string; pool: pg.Pool };
|
||||
const pools = new Map<string, Entrada>();
|
||||
|
||||
function chaveDaConexao(c: ConexaoExterna): string {
|
||||
return [c.id, c.versao, c.host, c.port, c.database, c.username, c.sslMode].join("\u0000");
|
||||
}
|
||||
|
||||
function sslPara(modo: ModoTls): pg.PoolConfig["ssl"] {
|
||||
switch (modo) {
|
||||
case "disable":
|
||||
return false;
|
||||
case "prefer":
|
||||
case "require":
|
||||
// Sem CA configurada não há como verificar a cadeia; `require` cifra mesmo
|
||||
// assim. `verify-*` usa as CAs do sistema e falha fechado se não bater.
|
||||
return { rejectUnauthorized: false };
|
||||
case "verify-ca":
|
||||
case "verify-full":
|
||||
return { rejectUnauthorized: true };
|
||||
}
|
||||
}
|
||||
|
||||
function configDe(c: ConexaoExterna): pg.PoolConfig {
|
||||
return {
|
||||
host: c.host,
|
||||
port: c.port,
|
||||
database: c.database,
|
||||
user: c.username,
|
||||
password: c.password,
|
||||
ssl: sslPara(c.sslMode),
|
||||
max: MAX_CONEXOES_POR_POOL,
|
||||
connectionTimeoutMillis: CONNECTION_TIMEOUT_MS,
|
||||
idleTimeoutMillis: IDLE_TIMEOUT_MS,
|
||||
application_name: "deskcomm-external-db",
|
||||
};
|
||||
}
|
||||
|
||||
function handlerDeErro(connectionId: string): (err: Error) => void {
|
||||
return (err: Error) => {
|
||||
const erro = (err.message.split("\n", 1)[0] ?? "").slice(0, 300);
|
||||
logger.warn("[external-db.pool] conexão caiu — recria no próximo uso", { connectionId, erro });
|
||||
};
|
||||
}
|
||||
|
||||
function evictarSeNecessario(): void {
|
||||
while (pools.size > MAX_POOLS) {
|
||||
const primeira = pools.keys().next().value as string | undefined;
|
||||
if (primeira === undefined) return;
|
||||
const entrada = pools.get(primeira);
|
||||
pools.delete(primeira);
|
||||
if (entrada) void entrada.pool.end().catch(() => undefined);
|
||||
}
|
||||
}
|
||||
|
||||
/** Pool da conexão, criando/invalidando conforme o cadastro atual. */
|
||||
export function obterPool(c: ConexaoExterna): pg.Pool {
|
||||
const chave = chaveDaConexao(c);
|
||||
const existente = pools.get(c.id);
|
||||
if (existente && existente.chave === chave) {
|
||||
pools.delete(c.id);
|
||||
pools.set(c.id, existente); // toque de LRU
|
||||
return existente.pool;
|
||||
}
|
||||
if (existente) {
|
||||
pools.delete(c.id);
|
||||
void existente.pool.end().catch(() => undefined);
|
||||
}
|
||||
const pool = new pg.Pool(configDe(c));
|
||||
pool.on("connect", (client) => client.on("error", handlerDeErro(c.id)));
|
||||
pool.on("error", () => undefined);
|
||||
pools.set(c.id, { chave, pool });
|
||||
evictarSeNecessario();
|
||||
return pool;
|
||||
}
|
||||
|
||||
/** Derruba o pool da conexão (usar ao desabilitar/editar/apagar). */
|
||||
export async function fecharPool(connectionId: string): Promise<void> {
|
||||
const entrada = pools.get(connectionId);
|
||||
if (!entrada) return;
|
||||
pools.delete(connectionId);
|
||||
await entrada.pool.end().catch(() => undefined);
|
||||
}
|
||||
|
||||
/** Encerra todos os pools — usado em testes e no desligamento do processo. */
|
||||
export async function fecharTodosOsPools(): Promise<void> {
|
||||
const entradas = [...pools.values()];
|
||||
pools.clear();
|
||||
await Promise.all(entradas.map((e) => e.pool.end().catch(() => undefined)));
|
||||
}
|
||||
|
||||
/**
|
||||
* Roda um SELECT dentro de `BEGIN READ ONLY`. A seleção só-leitura é imposta
|
||||
* pelo Postgres, não por análise de string: mesmo que algo escapasse da
|
||||
* validação de identificadores, um `INSERT`/`UPDATE`/DDL seria recusado aqui.
|
||||
*/
|
||||
export async function consultar<T extends pg.QueryResultRow = pg.QueryResultRow>(
|
||||
pool: pg.Pool,
|
||||
text: string,
|
||||
values: unknown[] = [],
|
||||
): Promise<pg.QueryResult<T>> {
|
||||
const client = await pool.connect();
|
||||
try {
|
||||
await client.query("begin read only");
|
||||
await client.query(`set local statement_timeout = ${STATEMENT_TIMEOUT_MS}`);
|
||||
await client.query(`set local lock_timeout = ${LOCK_TIMEOUT_MS}`);
|
||||
await client.query(`set local idle_in_transaction_session_timeout = ${IDLE_TX_TIMEOUT_MS}`);
|
||||
const resultado = await client.query<T>(text, values);
|
||||
await client.query("commit");
|
||||
return resultado;
|
||||
} catch (err) {
|
||||
await client.query("rollback").catch(() => undefined);
|
||||
throw err;
|
||||
} finally {
|
||||
client.release();
|
||||
}
|
||||
}
|
||||
|
||||
export type ResultadoDeTeste = { ok: true } | { ok: false; erro: string };
|
||||
|
||||
function mensagemSegura(err: unknown): string {
|
||||
if (err instanceof Error) {
|
||||
// Primeira linha, truncada. A senha não aparece em erro do pg; ainda assim
|
||||
// não ecoamos o objeto inteiro.
|
||||
return (err.message.split("\n", 1)[0] ?? "erro desconhecido").slice(0, 300);
|
||||
}
|
||||
return "erro desconhecido";
|
||||
}
|
||||
|
||||
/** Testa a conexão com um Client descartável, sem poluir o cache de pools. */
|
||||
export async function testarConexao(c: ConexaoExterna): Promise<ResultadoDeTeste> {
|
||||
const client = new pg.Client(configDe(c));
|
||||
// O erro do socket é tratado pelo catch; sem listener o Node derruba o processo
|
||||
// se o servidor cair no meio do teste.
|
||||
client.on("error", () => undefined);
|
||||
try {
|
||||
await client.connect();
|
||||
await client.query("select 1");
|
||||
return { ok: true };
|
||||
} catch (err) {
|
||||
return { ok: false, erro: mensagemSegura(err) };
|
||||
} finally {
|
||||
await client.end().catch(() => undefined);
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,116 @@
|
||||
import type { SupabaseClient } from "@supabase/supabase-js";
|
||||
import { beforeEach, describe, expect, it, vi } from "vitest";
|
||||
|
||||
vi.mock("@/lib/crypto/aes_gcm", () => ({
|
||||
byteaToBuffer: (v: unknown) => Buffer.from(String(v)),
|
||||
decryptKey: vi.fn(() => "senha-decifrada"),
|
||||
}));
|
||||
|
||||
vi.mock("@/lib/logger", () => ({
|
||||
logger: { error: vi.fn(), warn: vi.fn(), info: vi.fn(), debug: vi.fn() },
|
||||
}));
|
||||
|
||||
import { decryptKey } from "@/lib/crypto/aes_gcm";
|
||||
|
||||
import { carregarConexao } from "./credenciais";
|
||||
|
||||
const LINHA = {
|
||||
id: "conn-1",
|
||||
organization_id: "org-1",
|
||||
label: "Meu Postgres",
|
||||
host: "db.exemplo.com",
|
||||
port: 5432,
|
||||
database_name: "outro_crm",
|
||||
username: "leitor",
|
||||
password_encrypted: "\\xaa",
|
||||
password_iv: "\\xbb",
|
||||
password_tag: "\\xcc",
|
||||
ssl_mode: "require",
|
||||
enabled: true,
|
||||
updated_at: "2026-09-11T00:00:00.000Z",
|
||||
};
|
||||
|
||||
/** Admin falso: registra os `.eq()` para provar o filtro por organização. */
|
||||
function adminFalso(resultado: { data: unknown; error: unknown }) {
|
||||
const eq: Array<[string, unknown]> = [];
|
||||
const builder = {
|
||||
select: () => builder,
|
||||
eq: (coluna: string, valor: unknown) => {
|
||||
eq.push([coluna, valor]);
|
||||
return builder;
|
||||
},
|
||||
maybeSingle: async () => resultado,
|
||||
};
|
||||
const admin = { from: () => builder } as unknown as SupabaseClient;
|
||||
return { admin, eq };
|
||||
}
|
||||
|
||||
describe("carregarConexao", () => {
|
||||
beforeEach(() => {
|
||||
vi.mocked(decryptKey).mockReset();
|
||||
vi.mocked(decryptKey).mockReturnValue("senha-decifrada");
|
||||
});
|
||||
|
||||
it("filtra SEMPRE por organization_id e id — a lição da #236", async () => {
|
||||
const { admin, eq } = adminFalso({ data: LINHA, error: null });
|
||||
await carregarConexao(admin, "org-1", "conn-1");
|
||||
expect(eq).toEqual([
|
||||
["organization_id", "org-1"],
|
||||
["id", "conn-1"],
|
||||
]);
|
||||
});
|
||||
|
||||
it("devolve a conexão decifrada e usa o updated_at como versão do pool", async () => {
|
||||
const { admin } = adminFalso({ data: LINHA, error: null });
|
||||
await expect(carregarConexao(admin, "org-1", "conn-1")).resolves.toEqual({
|
||||
ok: true,
|
||||
conexao: {
|
||||
id: "conn-1",
|
||||
organizationId: "org-1",
|
||||
label: "Meu Postgres",
|
||||
host: "db.exemplo.com",
|
||||
port: 5432,
|
||||
database: "outro_crm",
|
||||
username: "leitor",
|
||||
password: "senha-decifrada",
|
||||
sslMode: "require",
|
||||
versao: "2026-09-11T00:00:00.000Z",
|
||||
},
|
||||
});
|
||||
});
|
||||
|
||||
it("conexão inexistente", async () => {
|
||||
const { admin } = adminFalso({ data: null, error: null });
|
||||
await expect(carregarConexao(admin, "org-1", "x")).resolves.toEqual({
|
||||
ok: false,
|
||||
motivo: "nao_encontrada",
|
||||
});
|
||||
});
|
||||
|
||||
it("conexão desativada não abre pool", async () => {
|
||||
const { admin } = adminFalso({ data: { ...LINHA, enabled: false }, error: null });
|
||||
await expect(carregarConexao(admin, "org-1", "x")).resolves.toEqual({
|
||||
ok: false,
|
||||
motivo: "desativada",
|
||||
});
|
||||
});
|
||||
|
||||
it("erro de leitura vira `banco`, sem vazar detalhe do driver", async () => {
|
||||
const { admin } = adminFalso({ data: null, error: { message: "detalhe interno" } });
|
||||
await expect(carregarConexao(admin, "org-1", "x")).resolves.toEqual({
|
||||
ok: false,
|
||||
motivo: "banco",
|
||||
});
|
||||
});
|
||||
|
||||
it("decrypt que falha (falta AI_CRED_AES_KEY) vira `cifra_indisponivel`", async () => {
|
||||
vi.mocked(decryptKey).mockImplementationOnce(() => {
|
||||
throw new Error("sem chave");
|
||||
});
|
||||
const { admin } = adminFalso({ data: LINHA, error: null });
|
||||
await expect(carregarConexao(admin, "org-1", "x")).resolves.toEqual({
|
||||
ok: false,
|
||||
motivo: "cifra_indisponivel",
|
||||
});
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,99 @@
|
||||
/**
|
||||
* A conexão externa — lida do banco e decifrada just-in-time.
|
||||
*
|
||||
* ⚠️ SEMPRE COM `organization_id` NO FILTRO. É a lição da #236: o admin client
|
||||
* bypassa RLS, e uma busca só por `id` casaria linha de outra organização. O
|
||||
* índice é `(organization_id, label)`, mas o filtro por org vem do
|
||||
* `activeOrg` resolvido no cookie/JWT — nunca do corpo da requisição.
|
||||
*
|
||||
* A senha só existe no objeto devolvido. Quem chama descarta a referência no
|
||||
* fim; ela nunca é logada, persistida em claro ou devolvida por rota.
|
||||
*/
|
||||
import type { SupabaseClient } from "@supabase/supabase-js";
|
||||
|
||||
import { byteaToBuffer, decryptKey } from "@/lib/crypto/aes_gcm";
|
||||
import { logger } from "@/lib/logger";
|
||||
|
||||
import type { ConexaoExterna, ModoTls } from "./types";
|
||||
|
||||
export type MotivoSemConexao =
|
||||
| "nao_encontrada"
|
||||
| "desativada"
|
||||
| "cifra_indisponivel"
|
||||
| "banco";
|
||||
|
||||
export type LeituraConexao =
|
||||
| { ok: true; conexao: ConexaoExterna }
|
||||
| { ok: false; motivo: MotivoSemConexao };
|
||||
|
||||
interface LinhaConexao {
|
||||
id: string;
|
||||
organization_id: string;
|
||||
label: string;
|
||||
host: string;
|
||||
port: number;
|
||||
database_name: string;
|
||||
username: string;
|
||||
password_encrypted: unknown;
|
||||
password_iv: unknown;
|
||||
password_tag: unknown;
|
||||
ssl_mode: ModoTls;
|
||||
enabled: boolean;
|
||||
updated_at: string;
|
||||
}
|
||||
|
||||
export async function carregarConexao(
|
||||
admin: SupabaseClient,
|
||||
organizationId: string,
|
||||
connectionId: string,
|
||||
): Promise<LeituraConexao> {
|
||||
const { data, error } = await admin
|
||||
.from("external_db_connections")
|
||||
.select(
|
||||
"id, organization_id, label, host, port, database_name, username, password_encrypted, password_iv, password_tag, ssl_mode, enabled, updated_at",
|
||||
)
|
||||
.eq("organization_id", organizationId)
|
||||
.eq("id", connectionId)
|
||||
.maybeSingle<LinhaConexao>();
|
||||
|
||||
if (error) {
|
||||
logger.error("[external-db.credencial] leitura falhou", {
|
||||
organizationId,
|
||||
connectionId,
|
||||
error: error.message,
|
||||
});
|
||||
return { ok: false, motivo: "banco" };
|
||||
}
|
||||
if (!data) return { ok: false, motivo: "nao_encontrada" };
|
||||
if (!data.enabled) return { ok: false, motivo: "desativada" };
|
||||
|
||||
let password: string;
|
||||
try {
|
||||
password = decryptKey({
|
||||
ciphertext: byteaToBuffer(data.password_encrypted),
|
||||
iv: byteaToBuffer(data.password_iv),
|
||||
tag: byteaToBuffer(data.password_tag),
|
||||
});
|
||||
} catch {
|
||||
// Não loga o erro do decrypt com detalhe que possa conter material da chave.
|
||||
// Cifra indisponível é problema de INSTALAÇÃO (falta AI_CRED_AES_KEY), e a
|
||||
// tela precisa distinguir isso de "conexão não existe".
|
||||
return { ok: false, motivo: "cifra_indisponivel" };
|
||||
}
|
||||
|
||||
return {
|
||||
ok: true,
|
||||
conexao: {
|
||||
id: data.id,
|
||||
organizationId: data.organization_id,
|
||||
label: data.label,
|
||||
host: data.host,
|
||||
port: data.port,
|
||||
database: data.database_name,
|
||||
username: data.username,
|
||||
password,
|
||||
sslMode: data.ssl_mode,
|
||||
versao: data.updated_at,
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,97 @@
|
||||
import { describe, expect, it } from "vitest";
|
||||
|
||||
import { ipDeBancoProibido, validarHostDeBanco } from "./guardas";
|
||||
|
||||
describe("ipDeBancoProibido", () => {
|
||||
it.each([
|
||||
"169.254.169.254", // metadata de nuvem
|
||||
"127.0.0.1",
|
||||
"127.10.20.30",
|
||||
"0.0.0.0",
|
||||
"100.64.0.1", // CGNAT
|
||||
"192.0.2.10", // TEST-NET
|
||||
"198.18.0.5", // benchmark
|
||||
"203.0.113.9", // TEST-NET
|
||||
"224.0.0.1", // multicast
|
||||
"240.0.0.1", // reservada
|
||||
"::1",
|
||||
"::",
|
||||
"fe80::1",
|
||||
"fc00::1",
|
||||
"ff02::1",
|
||||
"2001:db8::1",
|
||||
"::ffff:127.0.0.1",
|
||||
"::ffff:7f00:1", // IPv4-mapeado em hex — MESMO loopback que 127.0.0.1
|
||||
"::127.0.0.1", // IPv4-compatível (legado)
|
||||
"0:0:0:0:0:0:0:1", // loopback expandido
|
||||
"0000:0000:0000:0000:0000:0000:0000:0001", // loopback expandido
|
||||
"0:0:0:0:0:0:0:0", // :: expandido
|
||||
"64:ff9b::192.168.1.1", // NAT64 apontando para LAN
|
||||
])("bloqueia %s", (ip) => {
|
||||
expect(ipDeBancoProibido(ip)).toBe(true);
|
||||
});
|
||||
|
||||
it.each([
|
||||
"8.8.8.8",
|
||||
"1.1.1.1",
|
||||
"10.1.2.3", // RFC1918 — LAN é caso legítimo do self-host
|
||||
"172.16.5.5",
|
||||
"192.168.0.10",
|
||||
"2606:4700:4700::1111",
|
||||
"::ffff:8.8.8.8", // IPv4-mapeado público não é bloqueado por tabela
|
||||
])("permite %s", (ip) => {
|
||||
expect(ipDeBancoProibido(ip)).toBe(false);
|
||||
});
|
||||
|
||||
it("recusa o que não se sabe julgar", () => {
|
||||
expect(ipDeBancoProibido("nao-e-ip")).toBe(true);
|
||||
expect(ipDeBancoProibido("999.1.1.1")).toBe(true);
|
||||
});
|
||||
});
|
||||
|
||||
describe("validarHostDeBanco", () => {
|
||||
it("recusa host vazio, com barra ou com espaço", async () => {
|
||||
await expect(validarHostDeBanco("")).resolves.toEqual({ ok: false, motivo: "host_invalido" });
|
||||
await expect(validarHostDeBanco("db.local/path")).resolves.toEqual({
|
||||
ok: false,
|
||||
motivo: "host_invalido",
|
||||
});
|
||||
await expect(validarHostDeBanco("db local")).resolves.toEqual({
|
||||
ok: false,
|
||||
motivo: "host_invalido",
|
||||
});
|
||||
});
|
||||
|
||||
it("recusa literal de IP proibido", async () => {
|
||||
await expect(validarHostDeBanco("169.254.169.254")).resolves.toEqual({
|
||||
ok: false,
|
||||
motivo: "ip_especial",
|
||||
});
|
||||
});
|
||||
|
||||
it("aceita literal de IP público e IP de LAN", async () => {
|
||||
await expect(validarHostDeBanco("8.8.8.8")).resolves.toEqual({ ok: true, enderecos: ["8.8.8.8"] });
|
||||
await expect(validarHostDeBanco("10.0.0.7")).resolves.toEqual({ ok: true, enderecos: ["10.0.0.7"] });
|
||||
});
|
||||
|
||||
it("aceita IPv6 literal entre colchetes", async () => {
|
||||
await expect(validarHostDeBanco("[2606:4700::1111]")).resolves.toEqual({
|
||||
ok: true,
|
||||
enderecos: ["2606:4700::1111"],
|
||||
});
|
||||
});
|
||||
|
||||
it("resolve hostname e recusa quando cai em IP proibido (localhost)", async () => {
|
||||
const r = await validarHostDeBanco("localhost");
|
||||
expect(r.ok).toBe(false);
|
||||
if (!r.ok) expect(r.motivo).toBe("ip_especial");
|
||||
});
|
||||
|
||||
it("falha fechado quando o DNS não resolve", async () => {
|
||||
// `.invalid` é reservado por RFC 2606 e nunca resolve — sem rede no teste.
|
||||
await expect(validarHostDeBanco("db.invalid")).resolves.toEqual({
|
||||
ok: false,
|
||||
motivo: "dns_falhou",
|
||||
});
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,218 @@
|
||||
/**
|
||||
* Guarda de destino do banco externo.
|
||||
*
|
||||
* O `pg` abre TCP cru — não passa pelo allowlist HTTP de
|
||||
* `lib/agent-engine/edge/egress.ts`, nem pelo guard de webhook
|
||||
* (`lib/automation/outbound-ip.ts`). Sem uma guarda própria, cadastrar uma
|
||||
* conexão vira "posso fazer o servidor falar com qualquer coisa da rede dele".
|
||||
*
|
||||
* ⚠️ A POLÍTICA AQUI É DIFERENTE DA DO WEBHOOK, e é deliberado.
|
||||
*
|
||||
* `assertDestinoResolvidoSeguro` bloqueia TODA faixa privada, porque um webhook
|
||||
* é cadastrado por `manager` e não deve alcançar a rede interna do compose. Aqui
|
||||
* quem cadastra é `admin` — na prática, o dono da VPS — e o caso de uso legítimo
|
||||
* inclui um Postgres na mesma rede local. Então:
|
||||
*
|
||||
* - RFC1918 (10/8, 172.16/12, 192.168/16) é PERMITIDO: LAN é caso real.
|
||||
* - link-local/metadata (169.254/16), loopback, CGNAT, multicast, reservadas
|
||||
* e TEST-NET continuam BLOQUEADOS SEMPRE: nenhum deles é um banco de dados
|
||||
* de cliente, e o 169.254.169.254 entrega credencial de instância de nuvem.
|
||||
*
|
||||
* A janela de DNS-rebinding é a mesma descrita em `outbound-ip.ts`: entre a
|
||||
* resolução desta guarda e a que o `pg` faz, o DNS pode mudar. Fechar de vez
|
||||
* exigiria fixar o IP na conexão, o que muda SNI e quebra TLS com vários
|
||||
* destinos. A dívida fica declarada, não escondida.
|
||||
*/
|
||||
import { lookup } from "node:dns/promises";
|
||||
import { isIP, isIPv4, isIPv6 } from "node:net";
|
||||
|
||||
export type MotivoHostBloqueado =
|
||||
| "host_invalido"
|
||||
| "ip_especial"
|
||||
| "dns_falhou"
|
||||
| "dns_vazio";
|
||||
|
||||
export type ResultadoDeHost = { ok: true; enderecos: string[] } | { ok: false; motivo: MotivoHostBloqueado };
|
||||
|
||||
function ipv4ParaInt(ip: string): number | null {
|
||||
const partes = ip.split(".");
|
||||
if (partes.length !== 4) return null;
|
||||
let total = 0;
|
||||
for (const parte of partes) {
|
||||
const n = Number(parte);
|
||||
if (!Number.isInteger(n) || n < 0 || n > 255) return null;
|
||||
total = total * 256 + n;
|
||||
}
|
||||
return total;
|
||||
}
|
||||
|
||||
/**
|
||||
* Faixas IPv4 que NUNCA são destino de banco. Note a ausência de 10/8, 172.16/12
|
||||
* e 192.168/16: essas são permitidas de propósito (LAN do dono).
|
||||
*/
|
||||
const FAIXAS_PROIBIDAS: ReadonlyArray<readonly [string, number]> = [
|
||||
["0.0.0.0", 8], // "este host"
|
||||
["100.64.0.0", 10], // CGNAT
|
||||
["127.0.0.0", 8], // loopback
|
||||
["169.254.0.0", 16], // link-local — inclui o metadata de nuvem (169.254.169.254)
|
||||
["192.0.0.0", 24], // IETF protocol assignments
|
||||
["192.0.2.0", 24], // TEST-NET-1
|
||||
["198.18.0.0", 15], // benchmark
|
||||
["198.51.100.0", 24], // TEST-NET-2
|
||||
["203.0.113.0", 24], // TEST-NET-3
|
||||
["224.0.0.0", 4], // multicast
|
||||
["240.0.0.0", 4], // reservada (inclui 255.255.255.255)
|
||||
];
|
||||
|
||||
/**
|
||||
* Expande um IPv6 para 16 bytes. `null` se não for parseável. Cobre a compressão
|
||||
* `::`, grupos hex e IPv4 embutido (decimal, como em `::ffff:127.0.0.1`).
|
||||
*
|
||||
* A normalização é o ponto: `isIPv6` aceita VÁRIAS grafias do MESMO endereço.
|
||||
* `::1`, `0:0:0:0:0:0:0:1` e `::ffff:7f00:1` são todas loopback — decidir por
|
||||
* prefixo de string deixaria as duas últimas passarem batido.
|
||||
*/
|
||||
function ipv6ParaBytes(ip: string): number[] | null {
|
||||
const semZona = ip.split("%")[0] ?? "";
|
||||
const lados = semZona.split("::");
|
||||
if (lados.length > 2) return null;
|
||||
|
||||
const parseLado = (lado: string): number[] | null => {
|
||||
if (lado === "") return [];
|
||||
const tokens = lado.split(":");
|
||||
const bytes: number[] = [];
|
||||
for (let i = 0; i < tokens.length; i += 1) {
|
||||
const token = tokens[i] ?? "";
|
||||
if (token.includes(".")) {
|
||||
if (i !== tokens.length - 1) return null; // IPv4 embutido só no fim
|
||||
const v4 = token.split(".");
|
||||
if (v4.length !== 4) return null;
|
||||
for (const parte of v4) {
|
||||
const n = Number(parte);
|
||||
if (!Number.isInteger(n) || n < 0 || n > 255) return null;
|
||||
bytes.push(n);
|
||||
}
|
||||
} else {
|
||||
if (!/^[0-9a-f]{1,4}$/i.test(token)) return null;
|
||||
const n = Number.parseInt(token, 16);
|
||||
bytes.push((n >> 8) & 0xff, n & 0xff);
|
||||
}
|
||||
}
|
||||
return bytes;
|
||||
};
|
||||
|
||||
const esquerda = parseLado(lados[0] ?? "");
|
||||
const direita = parseLado(lados[1] ?? "");
|
||||
if (esquerda === null || direita === null) return null;
|
||||
|
||||
if (lados.length === 2) {
|
||||
const zeros = 16 - esquerda.length - direita.length;
|
||||
if (zeros < 0) return null;
|
||||
return [...esquerda, ...new Array<number>(zeros).fill(0), ...direita];
|
||||
}
|
||||
return esquerda.length === 16 ? esquerda : null;
|
||||
}
|
||||
|
||||
/** Acesso indexado que o `noUncheckedIndexedAccess` do repo não deixa mentir. */
|
||||
function byte(bytes: ReadonlyArray<number>, i: number): number {
|
||||
return bytes[i] ?? 0;
|
||||
}
|
||||
|
||||
function ipv6Proibido(bytes: ReadonlyArray<number>): boolean {
|
||||
if (bytes.length !== 16) return true; // não é IPv6: trata como perigoso
|
||||
if (bytes.every((b) => b === 0)) return true; // ::
|
||||
if (bytes.slice(0, 15).every((b) => b === 0) && byte(bytes, 15) === 1) return true; // ::1
|
||||
if (byte(bytes, 0) === 0xfe && (byte(bytes, 1) & 0xc0) === 0x80) return true; // link-local fe80::/10
|
||||
if ((byte(bytes, 0) & 0xfe) === 0xfc) return true; // ULA fc00::/7
|
||||
if (byte(bytes, 0) === 0xff) return true; // multicast ff00::/8
|
||||
if (
|
||||
byte(bytes, 0) === 0x20 &&
|
||||
byte(bytes, 1) === 0x01 &&
|
||||
byte(bytes, 2) === 0x0d &&
|
||||
byte(bytes, 3) === 0xb8 // 2001:db8::/32 (documentação)
|
||||
) {
|
||||
return true;
|
||||
}
|
||||
if (
|
||||
byte(bytes, 0) === 0x00 &&
|
||||
byte(bytes, 1) === 0x64 &&
|
||||
byte(bytes, 2) === 0xff &&
|
||||
byte(bytes, 3) === 0x9b &&
|
||||
bytes.slice(4, 12).every((b) => b === 0) // 64:ff9b::/96 (NAT64 — alcança IPv4)
|
||||
) {
|
||||
return true;
|
||||
}
|
||||
|
||||
// IPv4 embutido: `::ffff:a.b.c.d` (mapeado) e `::a.b.c.d` (compatível, legado).
|
||||
const dezZeros = bytes.slice(0, 10).every((b) => b === 0);
|
||||
if (dezZeros && byte(bytes, 10) === 0xff && byte(bytes, 11) === 0xff) {
|
||||
return ipDeBancoProibido(
|
||||
`${byte(bytes, 12)}.${byte(bytes, 13)}.${byte(bytes, 14)}.${byte(bytes, 15)}`,
|
||||
);
|
||||
}
|
||||
if (dezZeros && byte(bytes, 10) === 0 && byte(bytes, 11) === 0) {
|
||||
const embutido = `${byte(bytes, 12)}.${byte(bytes, 13)}.${byte(bytes, 14)}.${byte(bytes, 15)}`;
|
||||
if (embutido !== "0.0.0.0") return ipDeBancoProibido(embutido);
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
/** Um literal de IP é proibido? Exportado para teste direto da política. */
|
||||
export function ipDeBancoProibido(ip: string): boolean {
|
||||
if (isIPv4(ip)) {
|
||||
const alvo = ipv4ParaInt(ip);
|
||||
if (alvo === null) return true; // não parseou: trata como perigoso
|
||||
for (const [base, prefixo] of FAIXAS_PROIBIDAS) {
|
||||
const baseInt = ipv4ParaInt(base);
|
||||
if (baseInt === null) continue;
|
||||
const mascara = prefixo === 0 ? 0 : (0xffffffff << (32 - prefixo)) >>> 0;
|
||||
if (((alvo & mascara) >>> 0) === ((baseInt & mascara) >>> 0)) return true;
|
||||
}
|
||||
return false;
|
||||
}
|
||||
|
||||
if (isIPv6(ip)) {
|
||||
const bytes = ipv6ParaBytes(ip);
|
||||
if (bytes === null) return true; // não parseou: trata como perigoso
|
||||
return ipv6Proibido(bytes);
|
||||
}
|
||||
|
||||
return true; // nem IPv4 nem IPv6
|
||||
}
|
||||
|
||||
/** Tira colchetes de IPv6 literal e recusa host com espaço, barra ou vazio. */
|
||||
function normalizarHost(host: string): string | null {
|
||||
const h = host.trim().replace(/^\[/, "").replace(/\]$/, "");
|
||||
if (h === "" || h.length > 255) return null;
|
||||
if (/[\s/\\]/.test(h)) return null;
|
||||
return h;
|
||||
}
|
||||
|
||||
/**
|
||||
* Valida o destino: literal de IP é julgado direto; hostname é resolvido e
|
||||
* recusado se QUALQUER endereço cair em faixa proibida (rebinding devolve um
|
||||
* público e um privado; o `pg` pode escolher o privado, então recusa tudo).
|
||||
*/
|
||||
export async function validarHostDeBanco(host: string): Promise<ResultadoDeHost> {
|
||||
const h = normalizarHost(host);
|
||||
if (h === null) return { ok: false, motivo: "host_invalido" };
|
||||
|
||||
if (isIP(h) !== 0) {
|
||||
if (ipDeBancoProibido(h)) return { ok: false, motivo: "ip_especial" };
|
||||
return { ok: true, enderecos: [h] };
|
||||
}
|
||||
|
||||
let enderecos: Array<{ address: string }>;
|
||||
try {
|
||||
enderecos = await lookup(h, { all: true });
|
||||
} catch {
|
||||
// Falha fechado: recusar custa um cadastro; falhar aberto custa a rede interna.
|
||||
return { ok: false, motivo: "dns_falhou" };
|
||||
}
|
||||
if (enderecos.length === 0) return { ok: false, motivo: "dns_vazio" };
|
||||
|
||||
for (const { address } of enderecos) {
|
||||
if (ipDeBancoProibido(address)) return { ok: false, motivo: "ip_especial" };
|
||||
}
|
||||
return { ok: true, enderecos: enderecos.map((e) => e.address) };
|
||||
}
|
||||
@@ -0,0 +1,89 @@
|
||||
import type pg from "pg";
|
||||
import { describe, expect, it } from "vitest";
|
||||
|
||||
import { colunasDaTabela, descreverTabela, listarTabelas } from "./introspeccao";
|
||||
|
||||
interface LinhaCatalogo {
|
||||
schema: string;
|
||||
nome: string;
|
||||
tipo: string;
|
||||
coluna: string;
|
||||
tipo_dado: string;
|
||||
nulavel: string;
|
||||
posicao: number;
|
||||
chave_primaria: string[] | null;
|
||||
estimativa: number;
|
||||
}
|
||||
|
||||
/** Pool falso que devolve as linhas do catálogo para o SELECT e ignora `begin`. */
|
||||
function poolFalso(rows: LinhaCatalogo[]) {
|
||||
const consultas: string[] = [];
|
||||
const client = {
|
||||
query: async (text: string) => {
|
||||
consultas.push(text);
|
||||
return text.trimStart().startsWith("select")
|
||||
? { rows, fields: [], rowCount: rows.length }
|
||||
: { rows: [], fields: [], rowCount: 0 };
|
||||
},
|
||||
release: () => undefined,
|
||||
};
|
||||
const pool = {
|
||||
connect: async () => client,
|
||||
} as unknown as pg.Pool;
|
||||
return { pool, consultas };
|
||||
}
|
||||
|
||||
const ROWS: LinhaCatalogo[] = [
|
||||
{ schema: "public", nome: "pedidos", tipo: "BASE TABLE", coluna: "id", tipo_dado: "uuid", nulavel: "NO", posicao: 1, chave_primaria: ["id"], estimativa: 1234.5 },
|
||||
{ schema: "public", nome: "pedidos", tipo: "BASE TABLE", coluna: "total", tipo_dado: "numeric", nulavel: "YES", posicao: 2, chave_primaria: ["id"], estimativa: 1234.5 },
|
||||
{ schema: "vendas", nome: "resumo", tipo: "VIEW", coluna: "mes", tipo_dado: "text", nulavel: "YES", posicao: 1, chave_primaria: null, estimativa: -1 },
|
||||
];
|
||||
|
||||
describe("introspecção ao vivo", () => {
|
||||
it("agrupa colunas por tabela e mapeia o tipo", async () => {
|
||||
const { pool } = poolFalso(ROWS);
|
||||
const tabelas = await listarTabelas(pool);
|
||||
expect(tabelas).toHaveLength(2);
|
||||
|
||||
const pedidos = tabelas.find((t) => t.nome === "pedidos");
|
||||
expect(pedidos?.schema).toBe("public");
|
||||
expect(pedidos?.tipo).toBe("tabela");
|
||||
expect(pedidos?.colunas.map((c) => c.nome)).toEqual(["id", "total"]);
|
||||
expect(pedidos?.colunas[0]?.nulavel).toBe(false);
|
||||
expect(pedidos?.colunas[1]?.nulavel).toBe(true);
|
||||
expect(pedidos?.chavePrimaria).toEqual(["id"]);
|
||||
expect(pedidos?.estimativaLinhas).toBe(1235);
|
||||
|
||||
const resumo = tabelas.find((t) => t.nome === "resumo");
|
||||
expect(resumo?.tipo).toBe("view");
|
||||
expect(resumo?.chavePrimaria).toEqual([]); // view não tem PK
|
||||
// `reltuples` de view é -1; nunca vira contagem negativa na tela.
|
||||
expect(resumo?.estimativaLinhas).toBe(0);
|
||||
});
|
||||
|
||||
it("descreverTabela devolve a tabela pedida e `null` quando não existe", async () => {
|
||||
const { pool } = poolFalso(ROWS);
|
||||
await expect(descreverTabela(pool, "public", "pedidos")).resolves.toMatchObject({
|
||||
schema: "public",
|
||||
nome: "pedidos",
|
||||
});
|
||||
|
||||
const vazio = poolFalso([]);
|
||||
await expect(descreverTabela(vazio.pool, "public", "nao_existe")).resolves.toBeNull();
|
||||
});
|
||||
|
||||
it("colunasDaTabela devolve o conjunto de colunas — ou `null` se a tabela some", async () => {
|
||||
const { pool } = poolFalso(ROWS);
|
||||
await expect(colunasDaTabela(pool, "public", "pedidos")).resolves.toEqual(new Set(["id", "total"]));
|
||||
|
||||
const vazio = poolFalso([]);
|
||||
await expect(colunasDaTabela(vazio.pool, "public", "nao_existe")).resolves.toBeNull();
|
||||
});
|
||||
|
||||
it("a leitura passa por transação somente-leitura", async () => {
|
||||
const { pool, consultas } = poolFalso(ROWS);
|
||||
await listarTabelas(pool);
|
||||
expect(consultas[0]).toBe("begin read only");
|
||||
expect(consultas[consultas.length - 1]).toBe("commit");
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,131 @@
|
||||
/**
|
||||
* Introspecção AO VIVO do banco externo.
|
||||
*
|
||||
* Nada de schema espelhado: as tabelas mudam com frequência, então o retrato é
|
||||
* lido do catálogo na hora da consulta. Toda a leitura passa por
|
||||
* `consultar()` (transação somente-leitura) e por `information_schema` +
|
||||
* `pg_catalog` — portável entre Postgres e Supabase, sem depender de extensão.
|
||||
*
|
||||
* `pg_catalog` é excluído da listagem: o operador quer ver os dados dele, não
|
||||
* as 60 e poucas tabelas internas. Views entram junto com tabelas — muita vez é
|
||||
* por view que o outro sistema expõe o dado.
|
||||
*/
|
||||
import type pg from "pg";
|
||||
|
||||
import { consultar } from "./conexao";
|
||||
import type { TabelaExterna } from "./types";
|
||||
|
||||
interface LinhaCatalogo {
|
||||
schema: string;
|
||||
nome: string;
|
||||
tipo: string;
|
||||
coluna: string;
|
||||
tipo_dado: string;
|
||||
nulavel: string;
|
||||
posicao: number;
|
||||
chave_primaria: string[] | null;
|
||||
estimativa: string | number;
|
||||
}
|
||||
|
||||
const SQL_CATALOGO = `
|
||||
select
|
||||
c.table_schema as schema,
|
||||
c.table_name as nome,
|
||||
t.table_type as tipo,
|
||||
c.column_name as coluna,
|
||||
c.data_type as tipo_dado,
|
||||
c.is_nullable as nulavel,
|
||||
c.ordinal_position as posicao,
|
||||
pk.colunas as chave_primaria,
|
||||
coalesce(cl.reltuples, 0) as estimativa
|
||||
from information_schema.columns c
|
||||
join information_schema.tables t
|
||||
on t.table_schema = c.table_schema and t.table_name = c.table_name
|
||||
left join pg_catalog.pg_namespace n
|
||||
on n.nspname = c.table_schema
|
||||
left join pg_catalog.pg_class cl
|
||||
on cl.relname = c.table_name and cl.relnamespace = n.oid
|
||||
left join (
|
||||
select i.indrelid, array_agg(a.attname order by k.ord) as colunas
|
||||
from pg_catalog.pg_index i
|
||||
cross join lateral unnest(i.indkey) with ordinality as k(attnum, ord)
|
||||
join pg_catalog.pg_attribute a
|
||||
on a.attrelid = i.indrelid and a.attnum = k.attnum
|
||||
where i.indisprimary
|
||||
group by i.indrelid
|
||||
) pk on pk.indrelid = cl.oid
|
||||
where c.table_schema not in ('pg_catalog', 'information_schema', 'pg_toast')
|
||||
and t.table_type in ('BASE TABLE', 'VIEW')
|
||||
`;
|
||||
|
||||
function tipoDe(t: string): TabelaExterna["tipo"] {
|
||||
if (t === "BASE TABLE") return "tabela";
|
||||
if (t === "VIEW") return "view";
|
||||
return "outro";
|
||||
}
|
||||
|
||||
function agrupar(rows: LinhaCatalogo[]): TabelaExterna[] {
|
||||
const mapa = new Map<string, TabelaExterna>();
|
||||
for (const r of rows) {
|
||||
const chave = `${r.schema}\u0000${r.nome}`;
|
||||
let tabela = mapa.get(chave);
|
||||
if (!tabela) {
|
||||
tabela = {
|
||||
schema: r.schema,
|
||||
nome: r.nome,
|
||||
tipo: tipoDe(r.tipo),
|
||||
colunas: [],
|
||||
chavePrimaria: r.chave_primaria ?? [],
|
||||
estimativaLinhas: Math.max(0, Math.round(Number(r.estimativa) || 0)),
|
||||
};
|
||||
mapa.set(chave, tabela);
|
||||
}
|
||||
tabela.colunas.push({
|
||||
nome: r.coluna,
|
||||
tipo: r.tipo_dado,
|
||||
nulavel: r.nulavel === "YES",
|
||||
posicao: r.posicao,
|
||||
});
|
||||
}
|
||||
return [...mapa.values()];
|
||||
}
|
||||
|
||||
/** Todas as tabelas e views visíveis ao usuário da conexão, com suas colunas. */
|
||||
export async function listarTabelas(pool: pg.Pool): Promise<TabelaExterna[]> {
|
||||
const { rows } = await consultar<LinhaCatalogo>(
|
||||
pool,
|
||||
`${SQL_CATALOGO} order by c.table_schema, c.table_name, c.ordinal_position`,
|
||||
);
|
||||
return agrupar(rows);
|
||||
}
|
||||
|
||||
/** Retrato de uma tabela/view específica. `null` quando ela não existe. */
|
||||
export async function descreverTabela(
|
||||
pool: pg.Pool,
|
||||
schema: string,
|
||||
tabela: string,
|
||||
): Promise<TabelaExterna | null> {
|
||||
const { rows } = await consultar<LinhaCatalogo>(
|
||||
pool,
|
||||
`${SQL_CATALOGO}
|
||||
and c.table_schema = $1 and c.table_name = $2
|
||||
order by c.ordinal_position`,
|
||||
[schema, tabela],
|
||||
);
|
||||
const [primeira] = agrupar(rows);
|
||||
return primeira ?? null;
|
||||
}
|
||||
|
||||
/**
|
||||
* Nomes de coluna da tabela, para validar o pedido de leitura. `null` quando a
|
||||
* tabela não existe — o chamador distingue "não existe" de "existe sem colunas".
|
||||
*/
|
||||
export async function colunasDaTabela(
|
||||
pool: pg.Pool,
|
||||
schema: string,
|
||||
tabela: string,
|
||||
): Promise<Set<string> | null> {
|
||||
const descricao = await descreverTabela(pool, schema, tabela);
|
||||
if (!descricao) return null;
|
||||
return new Set(descricao.colunas.map((c) => c.nome));
|
||||
}
|
||||
@@ -0,0 +1,135 @@
|
||||
import { describe, expect, it } from "vitest";
|
||||
|
||||
import { LIMITE_MAX, LeituraInvalidaError, montarConsulta, quotarIdentificador } from "./leitura";
|
||||
import type { PedidoDeLeitura } from "./types";
|
||||
|
||||
const COLS = new Set(["id", "nome", "criado_em", "valor"]);
|
||||
|
||||
function pedido(over: Partial<PedidoDeLeitura> = {}): PedidoDeLeitura {
|
||||
return {
|
||||
schema: "public",
|
||||
tabela: "pedidos",
|
||||
colunas: [],
|
||||
filtros: [],
|
||||
limite: 50,
|
||||
offset: 0,
|
||||
...over,
|
||||
};
|
||||
}
|
||||
|
||||
describe("quotarIdentificador", () => {
|
||||
it("escapa aspas duplas", () => {
|
||||
expect(quotarIdentificador('a"b')).toBe('"a""b"');
|
||||
});
|
||||
});
|
||||
|
||||
describe("montarConsulta", () => {
|
||||
it("monta select básico com colunas, filtro, ordem e paginação", () => {
|
||||
const r = montarConsulta(
|
||||
pedido({
|
||||
colunas: ["id", "nome"],
|
||||
filtros: [{ coluna: "id", operador: "eq", valor: 42 }],
|
||||
ordem: { coluna: "nome", desc: true },
|
||||
limite: 10,
|
||||
offset: 20,
|
||||
}),
|
||||
COLS,
|
||||
);
|
||||
expect(r.text).toBe(
|
||||
'select "id", "nome" from "public"."pedidos" where "id" = $1 order by "nome" desc limit 10 offset 20',
|
||||
);
|
||||
expect(r.values).toEqual([42]);
|
||||
});
|
||||
|
||||
it("sem colunas vira `*`", () => {
|
||||
const r = montarConsulta(pedido(), COLS);
|
||||
expect(r.text).toContain('select * from "public"."pedidos"');
|
||||
});
|
||||
|
||||
it("recusa coluna que não existe no catálogo (injeção por identificador)", () => {
|
||||
expect(() =>
|
||||
montarConsulta(pedido({ colunas: ['nome"; drop table x; --'] }), COLS),
|
||||
).toThrow(LeituraInvalidaError);
|
||||
expect(() =>
|
||||
montarConsulta(pedido({ filtros: [{ coluna: "senha", operador: "eq", valor: "x" }] }), COLS),
|
||||
).toThrow(LeituraInvalidaError);
|
||||
});
|
||||
|
||||
it("silencia o nome da tabela mesmo com caractere estranho (caller valida existência)", () => {
|
||||
const r = montarConsulta(pedido({ tabela: 'ped idos"; --' }), COLS);
|
||||
expect(r.text).toContain('from "public"."ped idos""; --"');
|
||||
});
|
||||
|
||||
it("parametriza os valores na ordem das colunas", () => {
|
||||
const r = montarConsulta(
|
||||
pedido({
|
||||
filtros: [
|
||||
{ coluna: "nome", operador: "contem", valor: "ana" },
|
||||
{ coluna: "valor", operador: "gte", valor: 100 },
|
||||
],
|
||||
}),
|
||||
COLS,
|
||||
);
|
||||
expect(r.text).toContain(`"nome"`);
|
||||
expect(r.text).toContain(`"valor" >= $2`);
|
||||
expect(r.values).toEqual(["%ana%", 100]);
|
||||
});
|
||||
|
||||
it("escapa curingas do LIKE no `contem`", () => {
|
||||
const r = montarConsulta(
|
||||
pedido({ filtros: [{ coluna: "nome", operador: "contem", valor: "50%_a" }] }),
|
||||
COLS,
|
||||
);
|
||||
expect(r.values).toEqual(["%50\\%\\_a%"]);
|
||||
expect(r.text).toContain("escape '\\'");
|
||||
});
|
||||
|
||||
it("`in` com lista vazia vira false; com itens vira placeholders", () => {
|
||||
expect(montarConsulta(pedido({ filtros: [{ coluna: "id", operador: "in", valor: [] }] }), COLS).text).toContain(
|
||||
"where false",
|
||||
);
|
||||
const r = montarConsulta(
|
||||
pedido({ filtros: [{ coluna: "id", operador: "in", valor: [1, 2, 3] }] }),
|
||||
COLS,
|
||||
);
|
||||
expect(r.text).toContain('"id" in ($1, $2, $3)');
|
||||
expect(r.values).toEqual([1, 2, 3]);
|
||||
});
|
||||
|
||||
it("`in` sem array é erro", () => {
|
||||
expect(() =>
|
||||
montarConsulta(pedido({ filtros: [{ coluna: "id", operador: "in", valor: "1,2" }] }), COLS),
|
||||
).toThrow(LeituraInvalidaError);
|
||||
});
|
||||
|
||||
it("nulo/nao_nulo não consomem parâmetro", () => {
|
||||
const r = montarConsulta(
|
||||
pedido({
|
||||
filtros: [
|
||||
{ coluna: "nome", operador: "nulo" },
|
||||
{ coluna: "valor", operador: "nao_nulo" },
|
||||
],
|
||||
}),
|
||||
COLS,
|
||||
);
|
||||
expect(r.text).toContain('"nome" is null');
|
||||
expect(r.text).toContain('"valor" is not null');
|
||||
expect(r.values).toEqual([]);
|
||||
});
|
||||
|
||||
it("`eq` com null vira is null", () => {
|
||||
const r = montarConsulta(pedido({ filtros: [{ coluna: "nome", operador: "eq", valor: null }] }), COLS);
|
||||
expect(r.text).toContain('"nome" is null');
|
||||
});
|
||||
|
||||
it("engessa o limite no teto e cai no padrão quando inválido", () => {
|
||||
expect(montarConsulta(pedido({ limite: 999999 }), COLS).limite).toBe(LIMITE_MAX);
|
||||
expect(montarConsulta(pedido({ limite: 0 }), COLS).limite).toBe(50);
|
||||
expect(montarConsulta(pedido({ limite: Number.NaN }), COLS).limite).toBe(50);
|
||||
expect(montarConsulta(pedido({ offset: -5 }), COLS).offset).toBe(0);
|
||||
});
|
||||
|
||||
it("recusa tabela sem nome", () => {
|
||||
expect(() => montarConsulta(pedido({ tabela: "" }), COLS)).toThrow(LeituraInvalidaError);
|
||||
});
|
||||
});
|
||||
@@ -0,0 +1,179 @@
|
||||
/**
|
||||
* Leitura paginada de uma tabela/view do banco externo.
|
||||
*
|
||||
* O `SELECT` é MONTADO NO SERVIDOR: os identificadores são quotados e validados
|
||||
* contra o catálogo (`permitidas`), e os valores viram parâmetros `$n` — nunca
|
||||
* concatenação. O filtro é um vocabulário FECHADO de operadores; não existe
|
||||
* caminho por onde texto do usuário vire SQL.
|
||||
*
|
||||
* O limite é teto rígido (`LIMITE_MAX`): sem ele, um `select *` numa tabela de
|
||||
* milhões de linhas derruba o processo do worker.
|
||||
*/
|
||||
import type pg from "pg";
|
||||
|
||||
import { consultar } from "./conexao";
|
||||
import type { OperadorDeFiltro, PedidoDeLeitura } from "./types";
|
||||
|
||||
export const LIMITE_MAX = 200;
|
||||
export const LIMITE_PADRAO = 50;
|
||||
|
||||
/** Acima disso, um valor de célula é truncado antes de virar JSON. */
|
||||
const MAX_TEXTO = 20_000;
|
||||
|
||||
export class LeituraInvalidaError extends Error {
|
||||
constructor(motivo: string) {
|
||||
super(motivo);
|
||||
this.name = "LeituraInvalidaError";
|
||||
}
|
||||
}
|
||||
|
||||
/** Aspas duplas escapadas: um identificador nunca fecha a aspa por conta própria. */
|
||||
export function quotarIdentificador(nome: string): string {
|
||||
return `"${nome.replace(/"/g, '""')}"`;
|
||||
}
|
||||
|
||||
function escaparLike(valor: string): string {
|
||||
return valor.replace(/\\/g, "\\\\").replace(/%/g, "\\%").replace(/_/g, "\\_");
|
||||
}
|
||||
|
||||
function exigirColuna(coluna: string, permitidas: ReadonlySet<string>): void {
|
||||
if (!permitidas.has(coluna)) {
|
||||
throw new LeituraInvalidaError(`coluna_inexistente:${coluna}`);
|
||||
}
|
||||
}
|
||||
|
||||
/** Traduz um filtro em cláusula + parâmetros. Reutiliza o vetor de values. */
|
||||
function clausulaDeFiltro(
|
||||
operador: OperadorDeFiltro,
|
||||
colunaQuotada: string,
|
||||
valor: unknown,
|
||||
values: unknown[],
|
||||
): string {
|
||||
const placeholder = (v: unknown): string => {
|
||||
values.push(v);
|
||||
return `$${values.length}`;
|
||||
};
|
||||
|
||||
switch (operador) {
|
||||
case "eq":
|
||||
return valor === null || valor === undefined
|
||||
? `${colunaQuotada} is null`
|
||||
: `${colunaQuotada} = ${placeholder(valor)}`;
|
||||
case "ne":
|
||||
return valor === null || valor === undefined
|
||||
? `${colunaQuotada} is not null`
|
||||
: `${colunaQuotada} <> ${placeholder(valor)}`;
|
||||
case "gt":
|
||||
return `${colunaQuotada} > ${placeholder(valor)}`;
|
||||
case "gte":
|
||||
return `${colunaQuotada} >= ${placeholder(valor)}`;
|
||||
case "lt":
|
||||
return `${colunaQuotada} < ${placeholder(valor)}`;
|
||||
case "lte":
|
||||
return `${colunaQuotada} <= ${placeholder(valor)}`;
|
||||
case "contem":
|
||||
return `cast(${colunaQuotada} as text) ilike ${placeholder(`%${escaparLike(String(valor))}%`)} escape '\\'`;
|
||||
case "comeca_com":
|
||||
return `cast(${colunaQuotada} as text) ilike ${placeholder(`${escaparLike(String(valor))}%`)} escape '\\'`;
|
||||
case "in": {
|
||||
if (!Array.isArray(valor)) throw new LeituraInvalidaError("in_exige_array");
|
||||
if (valor.length === 0) return "false";
|
||||
const placeholders = valor.map((v) => placeholder(v));
|
||||
return `${colunaQuotada} in (${placeholders.join(", ")})`;
|
||||
}
|
||||
case "nulo":
|
||||
return `${colunaQuotada} is null`;
|
||||
case "nao_nulo":
|
||||
return `${colunaQuotada} is not null`;
|
||||
default: {
|
||||
const exaustivo: never = operador;
|
||||
throw new LeituraInvalidaError(`operador_desconhecido:${String(exaustivo)}`);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
export interface ConsultaMontada {
|
||||
text: string;
|
||||
values: unknown[];
|
||||
limite: number;
|
||||
offset: number;
|
||||
}
|
||||
|
||||
/**
|
||||
* Monta o SELECT. `permitidas` é o conjunto de colunas REAIS da tabela, lido do
|
||||
* catálogo — qualquer nome fora dele é recusado.
|
||||
*/
|
||||
export function montarConsulta(
|
||||
pedido: PedidoDeLeitura,
|
||||
permitidas: ReadonlySet<string>,
|
||||
): ConsultaMontada {
|
||||
if (!pedido.schema || !pedido.tabela) {
|
||||
throw new LeituraInvalidaError("tabela_obrigatoria");
|
||||
}
|
||||
|
||||
const colunas = [...new Set(pedido.colunas)];
|
||||
for (const c of colunas) exigirColuna(c, permitidas);
|
||||
|
||||
const values: unknown[] = [];
|
||||
const clausulas: string[] = [];
|
||||
for (const filtro of pedido.filtros) {
|
||||
exigirColuna(filtro.coluna, permitidas);
|
||||
clausulas.push(clausulaDeFiltro(filtro.operador, quotarIdentificador(filtro.coluna), filtro.valor, values));
|
||||
}
|
||||
|
||||
let ordem = "";
|
||||
if (pedido.ordem) {
|
||||
exigirColuna(pedido.ordem.coluna, permitidas);
|
||||
ordem = ` order by ${quotarIdentificador(pedido.ordem.coluna)} ${pedido.ordem.desc ? "desc" : "asc"}`;
|
||||
}
|
||||
|
||||
const limite = Math.min(LIMITE_MAX, Math.max(1, Math.floor(pedido.limite) || LIMITE_PADRAO));
|
||||
const offset = Math.max(0, Math.floor(pedido.offset) || 0);
|
||||
const projecao = colunas.length > 0 ? colunas.map(quotarIdentificador).join(", ") : "*";
|
||||
const onde = clausulas.length > 0 ? ` where ${clausulas.join(" and ")}` : "";
|
||||
|
||||
const text =
|
||||
`select ${projecao} from ${quotarIdentificador(pedido.schema)}.${quotarIdentificador(pedido.tabela)}` +
|
||||
`${onde}${ordem} limit ${limite} offset ${offset}`;
|
||||
|
||||
return { text, values, limite, offset };
|
||||
}
|
||||
|
||||
function serializarValor(v: unknown): unknown {
|
||||
if (v === null || v === undefined) return null;
|
||||
if (typeof v === "bigint") return v.toString();
|
||||
if (v instanceof Date) return v.toISOString();
|
||||
if (v instanceof Uint8Array) return `\\x${Buffer.from(v).toString("hex")}`;
|
||||
if (typeof v === "string" && v.length > MAX_TEXTO) {
|
||||
return `${v.slice(0, MAX_TEXTO)}…(truncado, ${v.length} chars)`;
|
||||
}
|
||||
return v;
|
||||
}
|
||||
|
||||
function serializarLinha(linha: Record<string, unknown>): Record<string, unknown> {
|
||||
const saida: Record<string, unknown> = {};
|
||||
for (const [k, v] of Object.entries(linha)) saida[k] = serializarValor(v);
|
||||
return saida;
|
||||
}
|
||||
|
||||
export interface ResultadoDeLeitura {
|
||||
colunas: string[];
|
||||
linhas: Record<string, unknown>[];
|
||||
limite: number;
|
||||
offset: number;
|
||||
}
|
||||
|
||||
export async function lerTabela(
|
||||
pool: pg.Pool,
|
||||
pedido: PedidoDeLeitura,
|
||||
permitidas: ReadonlySet<string>,
|
||||
): Promise<ResultadoDeLeitura> {
|
||||
const { text, values, limite, offset } = montarConsulta(pedido, permitidas);
|
||||
const resultado = await consultar<Record<string, unknown>>(pool, text, values);
|
||||
return {
|
||||
colunas: resultado.fields.map((f) => f.name),
|
||||
linhas: resultado.rows.map(serializarLinha),
|
||||
limite,
|
||||
offset,
|
||||
};
|
||||
}
|
||||
@@ -0,0 +1,82 @@
|
||||
/**
|
||||
* Tipos compartilhados do conector de PostgreSQL externo.
|
||||
*
|
||||
* O schema do banco externo NÃO é espelhado em TypeScript gerado: ele muda com
|
||||
* frequência e cópia de schema envelhece. Estes tipos descrevem só o que o
|
||||
* conector precisa — a conexão decifrada e o retrato do catálogo lido ao vivo.
|
||||
*/
|
||||
|
||||
/** Modos de TLS aceitos, alinhados ao CHECK de `external_db_connections.ssl_mode`. */
|
||||
export type ModoTls = "disable" | "prefer" | "require" | "verify-ca" | "verify-full";
|
||||
|
||||
/**
|
||||
* Conexão já decifrada, pronta para abrir o pool.
|
||||
*
|
||||
* A senha vive só no escopo de quem chamou `carregarConexao()` — nunca é logada,
|
||||
* cacheada em claro ou devolvida por rota. `versao` (o `updated_at` da linha)
|
||||
* existe para invalidar o pool quando o cadastro muda: editar a conexão sem
|
||||
* derrubar o pool antigo manteria credencial velha viva em memória.
|
||||
*/
|
||||
export interface ConexaoExterna {
|
||||
id: string;
|
||||
organizationId: string;
|
||||
label: string;
|
||||
host: string;
|
||||
port: number;
|
||||
database: string;
|
||||
username: string;
|
||||
password: string;
|
||||
sslMode: ModoTls;
|
||||
versao: string;
|
||||
}
|
||||
|
||||
/** Uma coluna do catálogo externo, na ordem em que aparece na tabela. */
|
||||
export interface ColunaExterna {
|
||||
nome: string;
|
||||
tipo: string;
|
||||
nulavel: boolean;
|
||||
posicao: number;
|
||||
}
|
||||
|
||||
/** Uma tabela ou view do banco externo, com as colunas do momento da leitura. */
|
||||
export interface TabelaExterna {
|
||||
schema: string;
|
||||
nome: string;
|
||||
tipo: "tabela" | "view" | "outro";
|
||||
colunas: ColunaExterna[];
|
||||
/** Colunas da chave primária, na ordem do índice. Vazio em view/sem PK. */
|
||||
chavePrimaria: string[];
|
||||
/** Estimativa do planner (`pg_class.reltuples`), não uma contagem exata. */
|
||||
estimativaLinhas: number;
|
||||
}
|
||||
|
||||
/** Operadores aceitos no filtro da consulta — vocabulário FECHADO, de propósito. */
|
||||
export type OperadorDeFiltro =
|
||||
| "eq"
|
||||
| "ne"
|
||||
| "gt"
|
||||
| "gte"
|
||||
| "lt"
|
||||
| "lte"
|
||||
| "contem"
|
||||
| "comeca_com"
|
||||
| "in"
|
||||
| "nulo"
|
||||
| "nao_nulo";
|
||||
|
||||
export interface FiltroDeLeitura {
|
||||
coluna: string;
|
||||
operador: OperadorDeFiltro;
|
||||
valor?: unknown;
|
||||
}
|
||||
|
||||
export interface PedidoDeLeitura {
|
||||
schema: string;
|
||||
tabela: string;
|
||||
/** Vazio = todas as colunas (`*`). */
|
||||
colunas: string[];
|
||||
filtros: FiltroDeLeitura[];
|
||||
ordem?: { coluna: string; desc?: boolean };
|
||||
limite: number;
|
||||
offset: number;
|
||||
}
|
||||
Reference in New Issue
Block a user