Skip to content

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).

ts
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 batchSize righe alla volta, non carica l'intero result-set.
  • Incrementale: parametrizza la query con un watermark (WHERE col > $1) preso da ctx.params.
  • Rispetta ctx.signal (annullamento) tra un batch e l'altro.
  • Errori di estrazione → ExtractError.

Sink — postgresSink

  • Idempotente: INSERT ... ON CONFLICT (conflictTarget) DO UPDATE in transazione.
    • conflictTarget vuoto → 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 → LoadError con ROLLBACK.

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.

ts
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.
  • keyColumn dice quale colonna del result set riassocia le righe alle chiavi. Una colonna chiave assente, o non scalare, è un LookupError — non zero risultati.
  • Sola lettura imposta a build(): assertReadOnlyQuery rifiuta tutto ciò che non è una singola SELECT/WITH — anche una CTE che scrive (WITH d AS (DELETE … RETURNING …)) e un SELECT … 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 solleva ConfigError.

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 in src/connectors/shared/sql-lookup.test.ts.
  • Integrazione: tests/integration/postgres.integration.test.ts — upsert idempotente, streaming e nodo lookup in un grafo contro un Postgres reale.
bash
# richiede un Postgres raggiungibile via env PG* (vedi .env)
pnpm run test:integration

Noeva è un marchio registrato di 4D S.R.L.