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:
root
2026-09-11 17:31:44 +02:00
co-authored by opencode
parent 73a274eff3
commit 2d41130223
12 changed files with 1384 additions and 9 deletions
+9 -9
View File
@@ -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/`
+54
View File
@@ -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);
});
});
+175
View File
@@ -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);
}
}
+116
View File
@@ -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",
});
});
});
+99
View File
@@ -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,
},
};
}
+97
View File
@@ -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",
});
});
});
+218
View File
@@ -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) };
}
+89
View File
@@ -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");
});
});
+131
View File
@@ -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));
}
+135
View File
@@ -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);
});
});
+179
View File
@@ -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,
};
}
+82
View File
@@ -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;
}