Connettore: PostgreSQL
Source, Sink e resolver del nodo lookup per PostgreSQL, basati su pg + pg-cursor (vedi ADR-0003). Codice: src/connectors/postgres/.
Configurazione
loadPgPoolConfig(prefix?, env?) costruisce una PoolConfig dalle variabili <PREFIX>HOST|PORT|USER|PASSWORD|DATABASE (default prefisso PG) e opzionale <PREFIX>SSL=true. Usa prefissi distinti per DB sorgente e destinazione. Le credenziali stanno solo in .env (vedi .env.example).
import { Pool } from "pg";
import {
loadPgPoolConfig,
postgresSource,
postgresSink,
} from "../connectors/postgres/index.js";
const pool = new Pool(loadPgPoolConfig("PG"));
try {
// estrazione (incrementale via watermark in ctx.params)
const source = postgresSource<Order>({
pool,
query:
"SELECT id, amount, updated_at FROM orders WHERE updated_at > $1 ORDER BY updated_at",
values: [ctx.params.since],
batchSize: 1000,
});
// caricamento idempotente (upsert su chiave naturale)
const sink = postgresSink<FactOrder>({
pool,
table: "dwh.fact_orders",
columns: ["orderId", "amount", "date"],
conflictTarget: ["orderId"],
});
} finally {
await pool.end(); // il chiamante possiede il ciclo di vita del Pool
}Source — postgresSource
- Streaming via cursor server-side: legge
batchSizerighe alla volta, non carica l'intero result-set. - Incrementale: parametrizza la query con un watermark (
WHERE col > $1) preso dactx.params. - Rispetta
ctx.signal(annullamento) tra un batch e l'altro. - Errori di estrazione →
ExtractError.
Sink — postgresSink
- Idempotente:
INSERT ... ON CONFLICT (conflictTarget) DO UPDATEin transazione.conflictTargetvuoto → INSERT non idempotente (sconsigliato).- tutte le colonne nella chiave →
DO NOTHING.
- Scrive a batch (
batchSize, default 500); l'intera load è in un'unica transazione (atomica). - Errori di caricamento →
LoadErrorconROLLBACK.
Nota: la transazione singola garantisce atomicità ma tiene aperto un lock sul carico molto grande; per volumi enormi valutare commit a finestre (futuro ADR).
Lookup — postgresLookup
Alimenta il nodo lookup del grafo (ADR-0041): risolve un batch di chiavi con un solo round-trip, in sola lettura. È il gemello di mssqlLookup: stessa query, stesso controllo di sola lettura, stessa normalizzazione — la parte comune vive in src/connectors/shared/sql-lookup.ts, qui cambia solo il driver.
postgresLookup({
pool,
instanceId: "dwh",
query:
"SELECT UPPER(codice) AS chiave, ordine FROM v WHERE UPPER(codice) IN (@chiavi)",
keyColumn: "chiave",
});- Il segnaposto
@chiavi(KEYS_PLACEHOLDER) è lo stesso di SQL Server e viene espanso in un parametro posizionale per chiave ($1, $2, …): i valori restano bindati. Con zero chiavi diventa un insieme vuoto valido, non SQL invalido. keyColumndice quale colonna del result set riassocia le righe alle chiavi. Una colonna chiave assente, o non scalare, è unLookupError— non zero risultati.- Sola lettura imposta a
build():assertReadOnlyQueryrifiuta tutto ciò che non è una singolaSELECT/WITH— anche una CTE che scrive (WITH d AS (DELETE … RETURNING …)) e unSELECT … FOR UPDATE.
⚠️ Maiuscole e spazi. Le chiavi arrivano al DB normalizzate (trim + maiuscolo) e senza duplicati; le righe vengono riassociate con la stessa normalizzazione e restituite sulla chiave così come l'ha passata il nodo. Su SQL Server la collation tipica è case-insensitive e ignora gli spazi in coda, quindi WHERE codice IN (@chiavi) basta; su Postgres no: se la colonna può contenere minuscole o spazi, confronta su UPPER(TRIM(codice)) (o su un indice espressione equivalente), altrimenti quelle righe non vengono trovate — senza errore, come un non-trovato.
⚠️ È un controllo sintattico, non un modello di sicurezza: non vede, per esempio, una funzione con effetti collaterali dentro la SELECT (nextval(), set_config()). La garanzia vera che la pipeline non scriva resta l'utenza di sola lettura sul database.
Sicurezza SQL
- I valori passano sempre come parametri
$n(mai concatenati). - Identificatori (tabella/colonne) provengono dal codice e sono validati/quotati (
quoteIdentifier): una stringa non conforme sollevaConfigError.
Test
- Unit (CI):
src/connectors/postgres/{query,config,lookup}.test.ts— SQL building, parsing config, resolver del lookup con pool finto. La parte comune del lookup è testata insrc/connectors/shared/sql-lookup.test.ts. - Integrazione:
tests/integration/postgres.integration.test.ts— upsert idempotente, streaming e nodolookupin un grafo contro un Postgres reale.
# richiede un Postgres raggiungibile via env PG* (vedi .env)
pnpm run test:integration
