← Files Aivana Database EngineerARCHIVED FILE
runtime/db/postgresAdapter.js
29.5 KB · Oct 2, 2026 · 00:34 UTC
const { BaseDbAdapter } = require("./baseAdapter");
const { validateSqlSafety } = require("../sqlSafety");
function mapRowRows(rows = []) {
return rows.map((row) => row.value || row);
}
const discoveryKinds = ["table", "view", "materialized_view", "column", "index", "primary_key", "foreign_key", "unique", "check", "default", "trigger", "procedure", "function", "sequence", "partition"];
const userNamespace = "n.nspname NOT IN ('pg_catalog','information_schema') AND n.nspname !~ '^pg_toast' AND n.nspname !~ '^pg_temp_'";
const pgDiscoveryQueries = {
table: `SELECT n.nspname AS schema, c.relname AS name, c.relkind AS relation_kind, c.relpersistence AS persistence,
c.relrowsecurity AS row_security FROM pg_catalog.pg_class c JOIN pg_catalog.pg_namespace n ON n.oid=c.relnamespace
WHERE ${userNamespace} AND c.relkind IN ('r','p','f')`,
column: `SELECT n.nspname AS schema, format('%I.%I',c.relname,a.attname) AS name, c.relname AS table_name,
a.attname AS column_name, a.attnum AS ordinal, pg_catalog.format_type(a.atttypid,a.atttypmod) AS data_type,
NOT a.attnotnull AS nullable, a.attidentity AS identity_kind
FROM pg_catalog.pg_attribute a JOIN pg_catalog.pg_class c ON c.oid=a.attrelid JOIN pg_catalog.pg_namespace n ON n.oid=c.relnamespace
WHERE ${userNamespace} AND c.relkind IN ('r','p','f','v','m') AND a.attnum>0 AND NOT a.attisdropped`,
index: `SELECT n.nspname AS schema, format('%I.%I',t.relname,c.relname) AS name, t.relname AS table_name,
c.relname AS index_name, pg_catalog.pg_get_indexdef(c.oid) AS definition, i.indisunique AS is_unique, i.indisvalid AS is_valid
FROM pg_catalog.pg_index i JOIN pg_catalog.pg_class c ON c.oid=i.indexrelid JOIN pg_catalog.pg_class t ON t.oid=i.indrelid
JOIN pg_catalog.pg_namespace n ON n.oid=t.relnamespace WHERE ${userNamespace}`,
default: `SELECT n.nspname AS schema, format('%I.%I',c.relname,a.attname) AS name, c.relname AS table_name,
a.attname AS column_name, pg_catalog.pg_get_expr(d.adbin,d.adrelid) AS definition
FROM pg_catalog.pg_attrdef d JOIN pg_catalog.pg_attribute a ON a.attrelid=d.adrelid AND a.attnum=d.adnum
JOIN pg_catalog.pg_class c ON c.oid=d.adrelid JOIN pg_catalog.pg_namespace n ON n.oid=c.relnamespace WHERE ${userNamespace}`,
trigger: `SELECT n.nspname AS schema, format('%I.%I',c.relname,t.tgname) AS name, c.relname AS table_name,
t.tgname AS trigger_name, t.tgenabled AS enabled_mode, pg_catalog.pg_get_triggerdef(t.oid) AS definition
FROM pg_catalog.pg_trigger t JOIN pg_catalog.pg_class c ON c.oid=t.tgrelid JOIN pg_catalog.pg_namespace n ON n.oid=c.relnamespace
WHERE ${userNamespace} AND NOT t.tgisinternal`,
sequence: `SELECT n.nspname AS schema, c.relname AS name, s.seqstart::text AS start_value, s.seqincrement::text AS increment,
s.seqmin::text AS minimum, s.seqmax::text AS maximum, s.seqcycle AS cycle
FROM pg_catalog.pg_sequence s JOIN pg_catalog.pg_class c ON c.oid=s.seqrelid JOIN pg_catalog.pg_namespace n ON n.oid=c.relnamespace WHERE ${userNamespace}`,
partition: `SELECT n.nspname AS schema, c.relname AS name, pn.nspname AS parent_schema, p.relname AS table_name,
pg_catalog.pg_get_expr(c.relpartbound,c.oid) AS definition
FROM pg_catalog.pg_inherits i JOIN pg_catalog.pg_class c ON c.oid=i.inhrelid JOIN pg_catalog.pg_class p ON p.oid=i.inhparent
JOIN pg_catalog.pg_namespace n ON n.oid=c.relnamespace JOIN pg_catalog.pg_namespace pn ON pn.oid=p.relnamespace
WHERE ${userNamespace} AND c.relispartition`,
};
for (const [kind, relkind] of [["view", "v"], ["materialized_view", "m"]]) {
pgDiscoveryQueries[kind] = `SELECT n.nspname AS schema, c.relname AS name, pg_catalog.pg_get_viewdef(c.oid,true) AS definition,
c.relispopulated AS populated FROM pg_catalog.pg_class c JOIN pg_catalog.pg_namespace n ON n.oid=c.relnamespace
WHERE ${userNamespace} AND c.relkind='${relkind}'`;
}
for (const [kind, type] of [["primary_key", "p"], ["foreign_key", "f"], ["unique", "u"], ["check", "c"]]) {
pgDiscoveryQueries[kind] = `SELECT n.nspname AS schema, format('%I.%I',c.relname,k.conname) AS name, c.relname AS table_name,
k.conname AS constraint_name, pg_catalog.pg_get_constraintdef(k.oid,true) AS definition, k.convalidated AS validated,
k.conkey AS column_numbers, rn.nspname AS referenced_schema, r.relname AS referenced_table, k.confkey AS referenced_column_numbers
FROM pg_catalog.pg_constraint k JOIN pg_catalog.pg_class c ON c.oid=k.conrelid JOIN pg_catalog.pg_namespace n ON n.oid=c.relnamespace
LEFT JOIN pg_catalog.pg_class r ON r.oid=k.confrelid LEFT JOIN pg_catalog.pg_namespace rn ON rn.oid=r.relnamespace
WHERE ${userNamespace} AND k.contype='${type}'`;
}
for (const [kind, type] of [["procedure", "p"], ["function", "f"]]) {
pgDiscoveryQueries[kind] = `SELECT n.nspname AS schema, format('%I(%s)',p.proname,pg_catalog.pg_get_function_identity_arguments(p.oid)) AS name,
p.proname AS routine_name, pg_catalog.pg_get_functiondef(p.oid) AS definition, p.prosecdef AS security_definer,
pg_catalog.pg_get_userbyid(p.proowner) AS owner FROM pg_catalog.pg_proc p JOIN pg_catalog.pg_namespace n ON n.oid=p.pronamespace
WHERE ${userNamespace} AND p.prokind='${type}'`;
}
class PostgresAdapter extends BaseDbAdapter {
constructor(context, connectionConfig) {
super(context);
this.connectionConfig = connectionConfig;
this.pool = null;
this.driver = null;
}
async initialize() {
try {
this.driver = require("pg");
} catch (_error) {
this.driver = null;
return { initialized: true, adapter: "mock_postgres" };
}
if (!this.connectionConfig || (!this.connectionConfig.connectionString && !this.connectionConfig.server)) {
return { initialized: true, adapter: "mock_postgres" };
}
const poolConfig = this.connectionConfig.connectionString
? { connectionString: this.connectionConfig.connectionString }
: {
host: this.connectionConfig.server,
user: this.connectionConfig.user,
password: this.connectionConfig.password,
database: this.connectionConfig.database,
port: this.connectionConfig.port ? Number(this.connectionConfig.port) : undefined,
};
this.pool = new this.driver.Pool({
...poolConfig,
connectionTimeoutMillis: Number(this.context.policy?.auth?.liveConnectionTimeoutMs || 5000),
max: Number(this.connectionConfig.poolMax || 5),
});
await this.pool.query("SELECT 1");
return { initialized: true, adapter: "postgres_adapter" };
}
async close() {
if (this.pool && this.pool.end) {
await this.pool.end();
}
return { closed: true };
}
getCapabilities() {
const live = Boolean(this.driver && this.pool);
return { engine: "postgres", live, discover: live, readOnly: live, explain: live, performance: live, security: live,
capabilitiesAreNotPermissionChecks: true, limitations: ["Performance requires pg_stat_statements and privileges; explain is not execution approval", "Discovery and security have no offline fallback"] };
}
async _readCatalog(sql) {
if (!this.driver || !this.pool) throw new Error("Live PostgreSQL connection required");
const client = await this.pool.connect();
let cleanupError;
try {
await client.query({ text: "BEGIN READ ONLY", query_timeout: 5000 });
await client.query({ text: "SET LOCAL statement_timeout = '5s'", query_timeout: 5000 });
await client.query({ text: "SET LOCAL search_path = pg_catalog", query_timeout: 5000 });
const result = await client.query({ text: sql, query_timeout: 5000 });
if (!Array.isArray(result.rows)) throw new Error("Invalid catalog response");
if (result.rows.length > 10000 || Buffer.byteLength(JSON.stringify(result.rows)) > 8 * 1024 * 1024) {
throw Object.assign(new Error("Catalog limit exceeded; discovery is incomplete"), { code: "CATALOG_LIMIT" });
}
return result.rows;
} finally {
try { await client.query({ text: "ROLLBACK", query_timeout: 5000 }); } catch (error) { cleanupError = error; }
client.release(cleanupError);
if (cleanupError) throw new Error("Catalog session cleanup failed");
}
}
async discoverDatabase() {
const [identity] = await this._readCatalog("SELECT current_database() AS database, version() AS version, current_user AS principal");
if (!identity?.database || !identity?.version) throw new Error("Missing live database identity");
const objects = [], coverage = {}, limitations = ["Visible user catalogs only; permission completeness is not independently established",
"PostgreSQL 11+ catalogs; aggregates, window routines, domain constraints and event triggers are outside V1", "Catalog queries are separate read-only transactions, not one atomic snapshot"];
const identities = new Set();
for (const kind of discoveryKinds) {
let rows;
try { rows = await this._readCatalog(`${pgDiscoveryQueries[kind]} ORDER BY 1,2 LIMIT 10001 /* discovery:${kind} */`); }
catch (error) {
if (!["42501", "42703", "42P01", "42883"].includes(error.code)) throw error;
coverage[kind] = "unavailable"; limitations.push(`${kind}: catalog permission or version unavailable`); continue;
}
coverage[kind] = "collected";
for (const row of rows) {
const { schema, name, definition, ...attributes } = row;
const id = JSON.stringify([kind, schema, name]);
if (typeof schema !== "string" || !schema || typeof name !== "string" || !name || identities.has(id)) throw new Error("Invalid or duplicate catalog identity");
identities.add(id);
objects.push({ kind, schema, name, ...(typeof definition === "string" ? { definition } : {}), attributes });
}
}
return { engine: "postgres", database: identity.database, version: identity.version, objects, coverage, limitations, source: "live", complete: false };
}
async getSecurityFindings() {
if (!this.driver || !this.pool) throw new Error("Live PostgreSQL connection required");
const queries = {
admin_roles: "SELECT rolname AS principal, rolsuper AS superuser, rolcreaterole AS create_role, rolcreatedb AS create_database, rolbypassrls AS bypass_rls FROM pg_catalog.pg_roles WHERE rolsuper OR rolcreaterole OR rolbypassrls",
admin_memberships: `SELECT m.rolname AS principal, r.rolname AS role, a.admin_option FROM pg_catalog.pg_auth_members a
JOIN pg_catalog.pg_roles r ON r.oid=a.roleid JOIN pg_catalog.pg_roles m ON m.oid=a.member
WHERE r.rolsuper OR r.rolcreaterole OR r.rolbypassrls OR a.admin_option OR r.rolname IN ('pg_read_all_data','pg_write_all_data','pg_execute_server_program','pg_read_server_files','pg_write_server_files')`,
security_definer: `SELECT n.nspname AS schema, p.proname || '(' || pg_catalog.pg_get_function_identity_arguments(p.oid) || ')' AS routine,
pg_catalog.pg_get_userbyid(p.proowner) AS owner, p.proconfig AS settings FROM pg_catalog.pg_proc p
JOIN pg_catalog.pg_namespace n ON n.oid=p.pronamespace WHERE ${userNamespace} AND p.prosecdef`,
public_grants: `SELECT n.nspname AS schema, c.relname AS object, a.privilege_type AS privilege FROM pg_catalog.pg_class c
JOIN pg_catalog.pg_namespace n ON n.oid=c.relnamespace CROSS JOIN LATERAL pg_catalog.aclexplode(c.relacl) a
WHERE ${userNamespace} AND a.grantee=0 AND a.privilege_type IN ('INSERT','UPDATE','DELETE','TRUNCATE','TRIGGER')
UNION ALL SELECT n.nspname, n.nspname, a.privilege_type FROM pg_catalog.pg_namespace n
CROSS JOIN LATERAL pg_catalog.aclexplode(n.nspacl) a WHERE ${userNamespace} AND a.grantee=0 AND a.privilege_type='CREATE'`,
};
const findings = [], coverage = {}, limitations = ["Visible direct catalog grants and memberships only; not an effective-permission or exploitability proof",
"No MFA, unused-account, nested membership, external identity or complete privilege assessment"];
for (const [category, sql] of Object.entries(queries)) {
let rows;
try { rows = await this._readCatalog(`${sql} ORDER BY 1 LIMIT 10001 /* security:${category} */`); }
catch (error) {
if (!["42501", "42703", "42P01", "42883"].includes(error.code)) throw error;
coverage[category] = "unavailable"; limitations.push(`${category}: catalog unavailable`); continue;
}
coverage[category] = "collected";
for (const row of rows) findings.push({ id: `${category}:${require("node:crypto").createHash("sha256").update(JSON.stringify(row)).digest("hex")}`,
severity: category === "security_definer" ? "medium" : "high", evidence: row,
reason: category === "security_definer" ? "Routine executes with owner privileges; review owner, search_path and execute grants" : "Elevated catalog privilege or broad grant requires least-privilege review" });
}
return require("../auditLogger").sanitizeObject({ source: "live", findings, coverage, limitations });
}
async listDatabases() {
if (!this.driver || !this.pool) {
return (this.sampleCatalog.postgres?.databases || []).map((item) => item);
}
const client = await this.pool.connect();
try {
const result = await client.query("SELECT datname FROM pg_database WHERE datistemplate = false ORDER BY datname;");
return result.rows.map((r) => ({ database: r.datname, engine: "postgres", state: "online", owner: "postgres" }));
} finally {
client.release();
}
}
async listTables({ database, schema = "public" }) {
const dbData = this.sampleCatalog.postgres?.schemas?.[database]?.[schema]?.tables || {};
if (!this.driver || !this.pool) {
return Object.entries(dbData).map(([table, meta]) => ({
database,
schema,
table,
type: "table",
rowCountHint: meta.rowCountHint,
indexCount: meta.indexes.length,
containsPIIFlag: Boolean(meta.containsPII),
}));
}
const client = await this.pool.connect();
try {
const result = await client.query(
"SELECT t.table_name, COALESCE(c.reltuples::bigint, 0) AS row_count_hint, " +
"COALESCE(i.index_count, 0) AS index_count " +
"FROM information_schema.tables t " +
"LEFT JOIN pg_class c ON c.relname = t.table_name " +
"LEFT JOIN pg_namespace n ON n.oid = c.relnamespace AND n.nspname = t.table_schema " +
"LEFT JOIN (SELECT schemaname, tablename, COUNT(*)::int AS index_count FROM pg_indexes GROUP BY schemaname, tablename) i " +
"ON i.schemaname = t.table_schema AND i.tablename = t.table_name " +
"WHERE t.table_schema=$1 AND t.table_type='BASE TABLE' ORDER BY t.table_name",
[schema]
);
return mapRowRows(result.rows).map((row) => ({
database,
schema,
table: row.table_name || row,
type: "table",
rowCountHint: Number(row.row_count_hint || 0),
indexCount: Number(row.index_count || 0),
containsPIIFlag: /email|ssn|phone|iban/i.test(String(row.table_name || "")),
}));
} finally {
client.release();
}
}
async describeTable({ database, schema = "public", table }) {
const catalog = this.sampleCatalog.postgres?.schemas?.[database]?.[schema]?.tables?.[table];
if (!this.driver || !this.pool) {
if (!catalog) {
return { error: "table_not_found", table, schema, database };
}
return {
database,
schema,
table,
columns: catalog.columns,
constraints: catalog.columns.filter((c) => c.isPrimaryKey || c.isForeignKey),
indexes: catalog.indexes,
relationships: catalog.relationships,
riskNotes: catalog.containsPII ? ["Contains PII-like columns"] : [],
};
}
const client = await this.pool.connect();
try {
const result = await client.query(
"SELECT column_name, is_nullable, data_type FROM information_schema.columns WHERE table_schema=$1 AND table_name=$2 ORDER BY ordinal_position",
[schema, table]
);
return {
database,
schema,
table,
columns: mapRowRows(result.rows).map((r) => ({ name: r.column_name, type: r.data_type, nullable: r.is_nullable === "YES" })),
};
} finally {
client.release();
}
}
async describeRelationships({ database, schema = "public" }) {
const schemaCatalog = this.sampleCatalog.postgres?.schemas?.[database]?.[schema]?.tables || {};
if (!this.driver || !this.pool) {
return {
edges: Object.entries(schemaCatalog).flatMap(([source, meta]) =>
(meta.relationships || []).map((rel) => ({
source,
target: rel.target,
cardinality: rel.cardinality,
joinHint: rel.via,
criticality: meta.containsPII ? "high" : "medium",
}))
),
};
}
const client = await this.pool.connect();
try {
const rows = await client.query(
"SELECT tc.table_name AS source_table, ccu.table_name AS target_table " +
"FROM information_schema.table_constraints tc " +
"JOIN information_schema.referential_constraints rc ON tc.constraint_name = rc.constraint_name " +
"JOIN information_schema.constraint_column_usage ccu ON ccu.constraint_name = rc.unique_constraint_name " +
"WHERE tc.constraint_type = 'FOREIGN KEY' AND tc.table_schema=$1",
[schema]
);
return {
edges: mapRowRows(rows.rows).map((row) => ({
source: row.source_table,
target: row.target_table,
cardinality: "many-to-one",
joinHint: "foreign key",
criticality: "medium",
})),
};
} finally {
client.release();
}
}
async replicationStatus() {
if (!this.driver || !this.pool) {
return {
topology: "primary-replica-2",
lagSeconds: 1.2,
trend: "stable",
consistencyRisk: "low",
actions: ["ensure_replication_slots", "monitor_lag"],
};
}
const client = await this.pool.connect();
try {
const rows = await client.query(
"SELECT CASE WHEN pg_is_in_recovery() THEN 'replica' ELSE 'primary' END AS mode, pg_is_in_recovery() AS is_replica"
);
const row = mapRowRows(rows.rows)[0] || { mode: "primary", is_replica: false };
return {
topology: `${row.mode}`,
lagSeconds: 0,
trend: "stable",
consistencyRisk: row.is_replica ? "low" : "info",
actions: ["monitor_lag"],
};
} finally {
client.release();
}
}
async queryStats() {
if (!this.driver || !this.pool) {
return [
{ queryId: "pg-001", avgMs: 120, p95Ms: 540, ioWait: "medium", cpuMs: 50, regressionScore: 0.14 },
{ queryId: "pg-002", avgMs: 15, p95Ms: 22, ioWait: "low", cpuMs: 4, regressionScore: 0.04 },
];
}
const client = await this.pool.connect();
try {
const rows = await client.query(
"SELECT queryid, calls, total_exec_time, mean_exec_time, rows, shared_blks_read, shared_blks_hit " +
"FROM pg_stat_statements WHERE dbid = (SELECT oid FROM pg_database WHERE datname = current_database()) " +
"ORDER BY total_exec_time DESC LIMIT 10"
);
return mapRowRows(rows.rows).map((r) => ({
queryId: `${r.queryId || r.queryid || "unknown"}`,
avgMs: Number(r.mean_exec_time),
totalElapsedMs: Number(r.total_exec_time),
p95Ms: null,
ioWait: "unknown",
cpuMs: null,
regressionScore: null,
sharedBlocksRead: r.shared_blks_read == null ? null : Number(r.shared_blks_read),
sharedBlocksHit: r.shared_blks_hit == null ? null : Number(r.shared_blks_hit),
executionCount: Number(r.calls || 0),
rows: Number(r.rows || 0),
source: "live",
metricEvidence: {
collector: "pg_stat_statements",
aggregation: "cumulative_since_statistics_reset",
unavailable: { p95Ms: "requires_execution_samples", cpuMs: "not_collected_by_pg_stat_statements",
ioWait: "block_counts_are_not_wait_durations", regressionScore: "requires_comparable_baseline" },
},
}));
} catch (error) {
const message = String(error?.message || error);
if (/pg_stat_statements|does not exist|undefined_table/i.test(message)) {
return {
error: "pg_stat_statements_unavailable",
queryStats: [],
remediation: "enable the pg_stat_statements extension and preload it in shared_preload_libraries",
source: "live_error",
};
}
return {
error: "query_stats_unavailable",
queryStats: [],
message,
source: "live_error",
};
} finally {
client.release();
}
}
async lockAnalysis() {
if (!this.driver || !this.pool) {
return {
deadlockRisk: "unknown",
topWaiters: [],
blockingChains: [],
remediationPlan: ["configure live connection for lock analysis"],
source: "mock",
};
}
const client = await this.pool.connect();
try {
const result = await client.query(
"SELECT a.pid, a.wait_event_type, a.wait_event, a.state, " +
"EXTRACT(MILLISECONDS FROM (now() - COALESCE(a.query_start, now())))::int AS duration_ms, " +
"c.relname AS relation_name " +
"FROM pg_stat_activity a " +
"LEFT JOIN pg_locks l ON l.pid = a.pid AND NOT l.granted " +
"LEFT JOIN pg_class c ON c.oid = l.relation " +
"WHERE a.wait_event_type IS NOT NULL OR NOT l.granted " +
"ORDER BY duration_ms DESC LIMIT 20"
);
const rows = mapRowRows(result.rows);
const topWaiters = rows.map((row) => ({
session: String(row.pid),
waitType: row.wait_event || row.wait_event_type || "unknown",
durationMs: Number(row.duration_ms || 0),
relation: row.relation_name || null,
state: row.state || "unknown",
}));
return {
deadlockRisk: topWaiters.some((row) => row.durationMs > 10000) ? "high" : topWaiters.length ? "medium" : "low",
topWaiters,
blockingChains: [],
remediationPlan: topWaiters.length
? ["inspect blocking pids with pg_blocking_pids", "shorten long transactions", "add missing indexes on lock-heavy joins"]
: ["no lock waits observed"],
source: "live",
};
} finally {
client.release();
}
}
async indexUsage({ schema = "public", table } = {}) {
if (!this.driver || !this.pool) {
return { indexes: [], recommendation: "configure live connection for index usage", source: "mock" };
}
const client = await this.pool.connect();
try {
const params = table ? [schema, table] : [schema];
const tableFilter = table ? "AND s.relname = $2 " : "";
const result = await client.query(
"SELECT s.relname, i.relname AS indexrelname, ui.idx_scan, ui.idx_tup_read, ui.idx_tup_fetch, " +
"pg_relation_size(i.oid) AS index_size_bytes " +
"FROM pg_stat_user_indexes ui " +
"JOIN pg_class s ON s.oid = ui.relid " +
"JOIN pg_class i ON i.oid = ui.indexrelid " +
"JOIN pg_namespace n ON n.oid = s.relnamespace " +
"WHERE n.nspname = $1 " +
tableFilter +
"ORDER BY ui.idx_scan DESC, pg_relation_size(i.oid) DESC LIMIT 20",
params
);
const indexes = mapRowRows(result.rows).map((row) => ({
table: row.relname,
index: row.indexrelname,
usageScore: Number(row.idx_scan || 0),
tuplesRead: Number(row.idx_tup_read || 0),
tuplesFetched: Number(row.idx_tup_fetch || 0),
storageBytes: Number(row.index_size_bytes || 0),
recommendation: Number(row.idx_scan || 0) === 0 ? "review unused index before dropping" : "keep monitoring workload benefit",
}));
return { table, indexes, source: "live" };
} finally {
client.release();
}
}
async explainQuery({ query, sql, analyze }) {
const statement = String(sql || query || "");
const safety = validateSqlSafety(statement, { readOnly: true });
if (!safety.safe) {
return {
error: "unsafe_sql",
executed: false,
violations: safety.violations,
source: this.driver && this.pool ? "live" : "mock",
};
}
if (!statement || !this.driver || !this.pool) {
return {
plan: "Adapter-ExecutionPlan(offline)",
bottlenecks: [],
rewriteHints: [],
estimatedCost: null,
confidence: "low",
source: "mock",
};
}
const client = await this.pool.connect();
try {
const explainOptions = "BUFFERS, FORMAT JSON";
const explainPrefix = analyze === true ? `EXPLAIN (ANALYZE, ${explainOptions})` : `EXPLAIN (${explainOptions})`;
const explain = await client.query(`${explainPrefix} ${statement}`);
const plan = explain.rows?.[0]?.["QUERY PLAN"];
const root = plan?.[0]?.Plan;
if (!root) throw new Error("EXPLAIN returned no plan");
const bottlenecks = [];
const rewriteHints = [];
const pending = [root];
while (pending.length) {
const node = pending.pop();
if (node["Node Type"] === "Seq Scan" && node.Filter) {
bottlenecks.push("filtered sequential scan; compare selectivity and table size before tuning");
rewriteHints.push("evaluate predicate index candidates against measured workload");
}
if (Number(node["Temp Written Blocks"]) > 0 || node["Sort Space Type"] === "Disk") {
bottlenecks.push("temporary disk usage observed in execution plan");
rewriteHints.push("inspect spilling operators and memory requirements");
}
pending.push(...(node.Plans || []));
}
return {
plan,
bottlenecks: [...new Set(bottlenecks)],
rewriteHints: [...new Set(rewriteHints)],
estimatedCost: Number.isFinite(root["Total Cost"]) ? root["Total Cost"] : null,
costUnit: "postgres_planner_cost",
confidence: "medium",
sourceQuery: statement,
source: "live",
};
} catch (_error) {
return {
error: "EXPLAIN_FAILED",
message: "Explain not available for statement",
plan: "adapter_explain_unavailable",
source: "live",
};
} finally {
client.release();
}
}
async detectPii({ database, schema = "public", table }) {
const catalog = this.sampleCatalog.postgres?.schemas?.[database]?.[schema]?.tables?.[table];
if (!this.driver || !this.pool) {
if (!catalog) {
return { piiColumns: [], sensitivityLevel: "unknown", policyViolations: [], remediation: "no schema context", source: "mock" };
}
const piiColumns = catalog.columns.filter((c) => c.pii).map((c) => c.name);
return {
piiColumns,
sensitivityLevel: piiColumns.length ? "high" : "low",
policyViolations: piiColumns.length ? ["PII exposure candidate detected"] : [],
remediation: piiColumns.length ? "mask output and apply least privilege" : "no direct PII in sampled metadata",
source: "mock",
};
}
if (!table) {
return { piiColumns: [], sensitivityLevel: "unknown", policyViolations: ["missing_target_table"], remediation: "provide table" };
}
const client = await this.pool.connect();
try {
const result = await client.query(
"SELECT column_name FROM information_schema.columns WHERE table_schema=$1 AND table_name=$2 AND " +
"(column_name ILIKE '%email%' OR column_name ILIKE '%ssn%' OR column_name ILIKE '%phone%' OR column_name ILIKE '%iban%')",
[schema, table]
);
const piiColumns = mapRowRows(result.rows).map((r) => r.column_name || r);
return {
piiColumns,
sensitivityLevel: piiColumns.length ? "high" : "low",
policyViolations: piiColumns.length ? ["PII exposure candidate detected"] : [],
remediation: piiColumns.length ? "mask output and narrow projection" : "no direct PII columns detected",
source: "live",
};
} finally {
client.release();
}
}
async executeSql(sql, options = {}) {
const statement = String(sql || "").trim();
if (!statement) {
return { error: "missing_sql", executed: false, source: "mock" };
}
const safety = validateSqlSafety(statement, {
readOnly: !options.isMigration,
allowWrite: Boolean(options.isMigration || options.allowWrite),
isMigration: Boolean(options.isMigration),
});
if (!safety.safe) {
return {
error: "unsafe_sql",
executed: false,
violations: safety.violations,
statement,
database: options.database,
schema: options.schema,
tool: options.tool,
source: this.driver && this.pool ? "live" : "mock",
};
}
const fallback = {
statement,
executed: false,
affectedRows: 0,
durationMs: 0,
database: options.database,
schema: options.schema,
tool: options.tool,
source: this.driver && this.pool ? "live" : "mock",
};
if (!this.driver || !this.pool) {
return {
...fallback,
mockReason: "no_connection",
};
}
const start = Date.now();
const client = await this.pool.connect();
try {
const result = await client.query(statement);
return {
...fallback,
executed: true,
affectedRows: result?.rowCount || 0,
durationMs: Date.now() - start,
fields: result?.fields ? result.fields.length : 0,
source: "live",
};
} catch (error) {
return {
...fallback,
error: "execution_failed",
message: String(error?.message || error),
durationMs: Date.now() - start,
source: "live_error",
};
} finally {
client.release();
}
}
}
module.exports = {
PostgresAdapter,
};
SHA-256: 1c194359c562470a637cc7c06dd3be6bc03eab3f889f63594902a8fd3e085b13