forked from libredb/libredb-studio
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathpostgres.ts
More file actions
1607 lines (1448 loc) · 59 KB
/
Copy pathpostgres.ts
File metadata and controls
1607 lines (1448 loc) · 59 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
812
813
814
815
816
817
818
819
820
821
822
823
824
825
826
827
828
829
830
831
832
833
834
835
836
837
838
839
840
841
842
843
844
845
846
847
848
849
850
851
852
853
854
855
856
857
858
859
860
861
862
863
864
865
866
867
868
869
870
871
872
873
874
875
876
877
878
879
880
881
882
883
884
885
886
887
888
889
890
891
892
893
894
895
896
897
898
899
900
901
902
903
904
905
906
907
908
909
910
911
912
913
914
915
916
917
918
919
920
921
922
923
924
925
926
927
928
929
930
931
932
933
934
935
936
937
938
939
940
941
942
943
944
945
946
947
948
949
950
951
952
953
954
955
956
957
958
959
960
961
962
963
964
965
966
967
968
969
970
971
972
973
974
975
976
977
978
979
980
981
982
983
984
985
986
987
988
989
990
991
992
993
994
995
996
997
998
999
1000
/**
* PostgreSQL Database Provider
* Full PostgreSQL support with connection pooling
*/
import { Pool, type PoolClient, type PoolConfig as PgPoolConfig, type QueryConfig } from "pg";
import { SQLBaseProvider } from "./sql-base";
import {
type DatabaseConnection,
type TableSchema,
type TableRelations,
type QueryResult,
type HealthInfo,
type MaintenanceType,
type MaintenanceResult,
type ProviderOptions,
type ProviderCapabilities,
type ProviderExecutionContext,
type ReadOnlyStatementBudget,
type SlowQuery,
type ActiveSession,
type DatabaseOverview,
type PerformanceMetrics,
type SlowQueryStats,
type ActiveSessionDetails,
type TableStats,
type IndexStats,
type StorageStats,
} from "../../types";
import {
DatabaseConfigError,
ConnectionError,
ExecutionProfileError,
QueryError,
mapDatabaseError,
} from "../../errors";
import { assertReadOnlyBudget, measureResultBytes } from "./read-only-budget";
import { formatBytes } from "../../utils/pool-manager";
// ============================================================================
// Type Definitions
// ============================================================================
interface PgStatActivityRow {
datname?: string;
pid?: number;
usename?: string;
application_name?: string;
client_addr?: string;
backend_start?: string | Date;
state?: string;
query?: string;
[key: string]: unknown;
}
// Row shapes returned by the schema introspection queries below.
interface SchemaRow {
table_schema: string;
table_name: string;
row_count: string;
total_size: string;
pk_columns: string[];
columns?: Array<{ name: string; type: string; nullable: boolean; defaultValue?: string | null }>;
indexes?: Array<{ name: string; columns: string[]; unique: boolean }>;
foreign_keys?: Array<{
columnName: string;
referencedSchema: string;
referencedTable: string;
referencedColumn: string;
}>;
}
type SchemaListRow = Omit<SchemaRow, "indexes" | "foreign_keys">;
type SchemaRelationRow = Pick<SchemaRow, "table_schema" | "table_name" | "foreign_keys" | "indexes">;
// ============================================================================
// Schema introspection SQL
// ----------------------------------------------------------------------------
// Hoisted to module scope (not inlined in the methods) on purpose. bun's
// coverage instruments the interior lines of a *multi-line template literal in
// a function body* as 0-hit in any test process that imports this file but does
// not exercise the method — and the merged lcov then reports those SQL lines as
// uncovered even though the method is tested. Evaluated once at module load,
// these consts are reported as covered everywhere, so coverage stays accurate.
//
// All CTEs are MATERIALIZED on purpose: PG12+ inlines single-reference CTEs,
// which lets the planner re-execute these information_schema-based CTEs inside
// nested-loop joins (it estimates rows=1 for them). On large schemas (100+
// tables/constraints/indexes) that explodes to minutes. MATERIALIZED forces
// each CTE to compute once — ~295s -> ~2.6s on a 122-table schema.
// ============================================================================
// Reusable CTE fragments (no trailing comma). Composed into the queries below;
// kept single-sourced so the shared CTEs aren't duplicated across queries.
const CTE_TABLES_INFO = `
tables_info AS MATERIALIZED (
SELECT
t.table_schema,
t.table_name,
COALESCE(c.reltuples::bigint, 0) as row_count,
COALESCE(pg_total_relation_size(c.oid), 0) as total_size
FROM information_schema.tables t
LEFT JOIN pg_class c ON c.oid = (quote_ident(t.table_schema) || '.' || quote_ident(t.table_name))::regclass
WHERE t.table_schema NOT IN ('pg_catalog', 'information_schema', 'pg_toast')
AND t.table_type = 'BASE TABLE'
)`;
const CTE_COLUMNS_INFO = `
columns_info AS MATERIALIZED (
SELECT
c.table_schema,
c.table_name,
json_agg(
json_build_object(
'name', c.column_name,
'type', c.data_type,
'nullable', c.is_nullable = 'YES',
'defaultValue', c.column_default
) ORDER BY c.ordinal_position
) FILTER (WHERE c.ordinal_position <= 100) as columns
FROM information_schema.columns c
WHERE c.table_schema NOT IN ('pg_catalog', 'information_schema', 'pg_toast')
GROUP BY c.table_schema, c.table_name
)`;
const CTE_PK_INFO = `
pk_info AS MATERIALIZED (
SELECT
tc.table_schema,
tc.table_name,
array_agg(kcu.column_name) as pk_columns
FROM information_schema.table_constraints tc
JOIN information_schema.key_column_usage kcu
ON tc.constraint_name = kcu.constraint_name
AND tc.table_schema = kcu.table_schema
WHERE tc.constraint_type = 'PRIMARY KEY'
GROUP BY tc.table_schema, tc.table_name
)`;
const CTE_FK_INFO = `
fk_info AS MATERIALIZED (
SELECT
tc.table_schema,
tc.table_name,
json_agg(
json_build_object(
'columnName', kcu.column_name,
'referencedSchema', ccu.table_schema,
'referencedTable', ccu.table_name,
'referencedColumn', ccu.column_name
)
) as foreign_keys
FROM information_schema.table_constraints tc
JOIN information_schema.key_column_usage kcu
ON tc.constraint_name = kcu.constraint_name
AND tc.table_schema = kcu.table_schema
JOIN information_schema.constraint_column_usage ccu
ON ccu.constraint_name = tc.constraint_name
AND ccu.constraint_schema = tc.constraint_schema
WHERE tc.constraint_type = 'FOREIGN KEY'
GROUP BY tc.table_schema, tc.table_name
)`;
const CTE_INDEX_INFO = `
index_info AS MATERIALIZED (
SELECT
n.nspname as table_schema,
t.relname as table_name,
json_agg(
json_build_object(
'name', i.relname,
'columns', (
SELECT array_agg(a.attname ORDER BY array_position(ix.indkey, a.attnum))
FROM pg_attribute a
WHERE a.attrelid = t.oid AND a.attnum = ANY(ix.indkey)
),
'unique', ix.indisunique
)
) as indexes
FROM pg_index ix
JOIN pg_class t ON t.oid = ix.indrelid
JOIN pg_class i ON i.oid = ix.indexrelid
JOIN pg_namespace n ON n.oid = t.relnamespace
WHERE n.nspname NOT IN ('pg_catalog', 'information_schema', 'pg_toast')
GROUP BY n.nspname, t.relname
)`;
// Full schema: tables + columns + PKs + foreign keys + indexes in one query.
const SCHEMA_FULL_SQL = `
WITH ${CTE_TABLES_INFO},${CTE_COLUMNS_INFO},${CTE_PK_INFO},${CTE_FK_INFO},${CTE_INDEX_INFO}
SELECT
ti.table_schema,
ti.table_name,
ti.row_count,
ti.total_size,
COALESCE(ci.columns, '[]'::json) as columns,
COALESCE(pk.pk_columns, ARRAY[]::text[]) as pk_columns,
COALESCE(fk.foreign_keys, '[]'::json) as foreign_keys,
COALESCE(ii.indexes, '[]'::json) as indexes
FROM tables_info ti
LEFT JOIN columns_info ci ON ci.table_schema = ti.table_schema AND ci.table_name = ti.table_name
LEFT JOIN pk_info pk ON pk.table_schema = ti.table_schema AND pk.table_name = ti.table_name
LEFT JOIN fk_info fk ON fk.table_schema = ti.table_schema AND fk.table_name = ti.table_name
LEFT JOIN index_info ii ON ii.table_schema = ti.table_schema AND ii.table_name = ti.table_name
ORDER BY ti.table_schema, ti.table_name ASC;
`;
// Fast structural list: tables + columns + PKs only (no FK/index joins).
const SCHEMA_LIST_SQL = `
WITH ${CTE_TABLES_INFO},${CTE_COLUMNS_INFO},${CTE_PK_INFO}
SELECT
ti.table_schema,
ti.table_name,
ti.row_count,
ti.total_size,
COALESCE(ci.columns, '[]'::json) as columns,
COALESCE(pk.pk_columns, ARRAY[]::text[]) as pk_columns
FROM tables_info ti
LEFT JOIN columns_info ci ON ci.table_schema = ti.table_schema AND ci.table_name = ti.table_name
LEFT JOIN pk_info pk ON pk.table_schema = ti.table_schema AND pk.table_name = ti.table_name
ORDER BY ti.table_schema, ti.table_name ASC;
`;
// Heavy relationship/index introspection (foreign keys + indexes).
const SCHEMA_RELATIONS_SQL = `
WITH ${CTE_FK_INFO},${CTE_INDEX_INFO}
SELECT
COALESCE(fk.table_schema, ii.table_schema) as table_schema,
COALESCE(fk.table_name, ii.table_name) as table_name,
COALESCE(fk.foreign_keys, '[]'::json) as foreign_keys,
COALESCE(ii.indexes, '[]'::json) as indexes
FROM fk_info fk
FULL OUTER JOIN index_info ii
ON ii.table_schema = fk.table_schema AND ii.table_name = fk.table_name;
`;
// ============================================================================
// Monitoring & maintenance SQL
// ----------------------------------------------------------------------------
// Hoisted to module scope for the same coverage reason as the schema SQL
// above: bun reports interior lines of method-body template literals as 0-hit
// in test processes that import this module without executing the method.
// ============================================================================
// getHealth: buffer cache hit ratio across user tables.
const HEALTH_CACHE_HIT_SQL = `
SELECT
sum(heap_blks_read) as heap_read,
sum(heap_blks_hit) as heap_hit,
COALESCE(
ROUND((sum(heap_blks_hit) * 100.0 / NULLIF(sum(heap_blks_hit) + sum(heap_blks_read), 0)), 1),
100
) as ratio
FROM pg_statio_user_tables;
`;
// getHealth: top slow queries from pg_stat_statements (optional extension).
const HEALTH_SLOW_QUERIES_SQL = `
SELECT
LEFT(query, 100) as query,
calls,
ROUND((mean_exec_time)::numeric, 2)::text || 'ms' as avgTime
FROM pg_stat_statements
WHERE calls > 0
ORDER BY total_exec_time DESC
LIMIT 5;
`;
// getHealth: recent sessions for the current database ($1 = database).
const HEALTH_SESSIONS_SQL = `
SELECT
pid,
usename as user,
datname as database,
COALESCE(state, 'unknown') as state,
LEFT(COALESCE(query, ''), 100) as query,
CASE
WHEN xact_start IS NOT NULL THEN
EXTRACT(EPOCH FROM (NOW() - xact_start))::text || 's'
ELSE 'N/A'
END as duration
FROM pg_stat_activity
WHERE datname = $1
AND pid != pg_backend_pid()
ORDER BY xact_start DESC NULLS LAST
LIMIT 10;
`;
// getOverview: server version, start time, and uptime.
const OVERVIEW_INFO_SQL = `
SELECT
version() as version,
pg_postmaster_start_time() as start_time,
EXTRACT(EPOCH FROM (now() - pg_postmaster_start_time()))::bigint as uptime_seconds
`;
// getOverview: active vs max connections ($1 = database).
const OVERVIEW_CONNECTIONS_SQL = `
SELECT
count(*) as active_connections,
(SELECT setting::int FROM pg_settings WHERE name = 'max_connections') as max_connections
FROM pg_stat_activity
WHERE datname = $1
`;
// getOverview: database size, pretty-printed and raw bytes ($1 = database).
const OVERVIEW_SIZE_SQL = `
SELECT
pg_size_pretty(pg_database_size($1)) as database_size,
pg_database_size($1) as database_size_bytes
`;
// getOverview: user table and index counts across all user schemas.
const OVERVIEW_COUNTS_SQL = `
SELECT
(SELECT count(*) FROM pg_tables WHERE schemaname NOT IN ('pg_catalog', 'information_schema', 'pg_toast')) as table_count,
(SELECT count(*) FROM pg_indexes WHERE schemaname NOT IN ('pg_catalog', 'information_schema', 'pg_toast')) as index_count
`;
// getPerformanceMetrics: buffer cache hit ratio.
const PERF_CACHE_HIT_SQL = `
SELECT
COALESCE(
ROUND(sum(heap_blks_hit) * 100.0 / NULLIF(sum(heap_blks_hit) + sum(heap_blks_read), 0), 2),
100
) as cache_hit_ratio
FROM pg_statio_user_tables
`;
// getPerformanceMetrics: transaction stats for the database ($1 = database).
const PERF_TRANSACTION_STATS_SQL = `
SELECT
xact_commit,
xact_rollback,
deadlocks,
blks_read,
blks_hit
FROM pg_stat_database
WHERE datname = $1
`;
// getPerformanceMetrics: checkpoint timings (columns absent on older PG).
const PERF_CHECKPOINT_SQL = `
SELECT
checkpoint_write_time,
checkpoint_sync_time
FROM pg_stat_bgwriter
`;
// getSlowQueries: pg_stat_statements stats ($1 = database, $2 = limit).
const SLOW_QUERIES_SQL = `
SELECT
queryid::text as query_id,
LEFT(query, 500) as query,
calls,
ROUND(total_exec_time::numeric, 2) as total_time,
ROUND(mean_exec_time::numeric, 2) as avg_time,
ROUND(min_exec_time::numeric, 2) as min_time,
ROUND(max_exec_time::numeric, 2) as max_time,
rows,
shared_blks_hit,
shared_blks_read
FROM pg_stat_statements
WHERE calls > 0
AND dbid = (SELECT oid FROM pg_database WHERE datname = $1)
ORDER BY total_exec_time DESC
LIMIT $2
`;
// getSlowQueries fallback: currently running queries from pg_stat_activity
// ($1 = database, $2 = limit).
const SLOW_QUERIES_FALLBACK_SQL = `
SELECT
pid::text as query_id,
LEFT(COALESCE(query, ''), 500) as query,
1 as calls,
COALESCE(EXTRACT(EPOCH FROM (now() - query_start)) * 1000, 0) as total_time,
COALESCE(EXTRACT(EPOCH FROM (now() - query_start)) * 1000, 0) as avg_time,
0 as rows
FROM pg_stat_activity
WHERE datname = $1
AND pid != pg_backend_pid()
AND state = 'active'
AND query IS NOT NULL
AND query != ''
AND query NOT LIKE '%pg_stat_activity%'
ORDER BY query_start ASC NULLS LAST
LIMIT $2
`;
// getActiveSessions: detailed session list ($1 = database, $2 = limit).
const ACTIVE_SESSIONS_SQL = `
SELECT
pid,
usename as user,
datname as database,
application_name,
client_addr::text,
COALESCE(state, 'unknown') as state,
LEFT(COALESCE(query, ''), 500) as query,
query_start,
wait_event_type,
wait_event,
CASE
WHEN state = 'active' THEN
EXTRACT(EPOCH FROM (now() - query_start))::text || 's'
WHEN xact_start IS NOT NULL THEN
EXTRACT(EPOCH FROM (now() - xact_start))::text || 's'
ELSE 'N/A'
END as duration,
CASE
WHEN state = 'active' THEN
EXTRACT(EPOCH FROM (now() - query_start)) * 1000
WHEN xact_start IS NOT NULL THEN
EXTRACT(EPOCH FROM (now() - xact_start)) * 1000
ELSE 0
END as duration_ms
FROM pg_stat_activity
WHERE datname = $1
AND pid != pg_backend_pid()
ORDER BY
CASE state WHEN 'active' THEN 0 ELSE 1 END,
query_start DESC NULLS LAST
LIMIT $2
`;
// getTableStats: per-table stats. A schema WHERE clause is interpolated
// between the two fragments at the call site.
const TABLE_STATS_SELECT_SQL = `
SELECT
schemaname as schema_name,
relname as table_name,
n_live_tup as live_row_count,
n_dead_tup as dead_row_count,
n_live_tup + n_dead_tup as row_count,
pg_size_pretty(pg_table_size(quote_ident(schemaname) || '.' || quote_ident(relname))) as table_size,
pg_table_size(quote_ident(schemaname) || '.' || quote_ident(relname)) as table_size_bytes,
pg_size_pretty(pg_indexes_size(quote_ident(schemaname) || '.' || quote_ident(relname))) as index_size,
pg_indexes_size(quote_ident(schemaname) || '.' || quote_ident(relname)) as index_size_bytes,
pg_size_pretty(pg_total_relation_size(quote_ident(schemaname) || '.' || quote_ident(relname))) as total_size,
pg_total_relation_size(quote_ident(schemaname) || '.' || quote_ident(relname)) as total_size_bytes,
last_vacuum,
last_autovacuum,
last_analyze,
last_autoanalyze,
CASE
WHEN n_live_tup > 0 THEN
ROUND(n_dead_tup * 100.0 / (n_live_tup + n_dead_tup), 2)
ELSE 0
END as bloat_ratio
FROM pg_stat_user_tables
`;
const TABLE_STATS_ORDER_SQL = `
ORDER BY pg_total_relation_size(quote_ident(schemaname) || '.' || quote_ident(relname)) DESC
`;
// getIndexStats: per-index stats. A schema WHERE clause is interpolated
// between the two fragments at the call site.
const INDEX_STATS_SELECT_SQL = `
SELECT
s.schemaname as schema_name,
s.relname as table_name,
s.indexrelname as index_name,
am.amname as index_type,
pg_size_pretty(pg_relation_size(s.indexrelid)) as index_size,
pg_relation_size(s.indexrelid) as index_size_bytes,
s.idx_scan as scans,
s.idx_tup_read as tuples_read,
s.idx_tup_fetch as tuples_fetched,
ix.indisunique as is_unique,
ix.indisprimary as is_primary,
array_agg(a.attname ORDER BY array_position(ix.indkey, a.attnum)) as columns,
CASE
WHEN (SELECT seq_scan + idx_scan FROM pg_stat_user_tables t WHERE t.relid = s.relid) > 0
THEN ROUND(
s.idx_scan * 100.0 /
(SELECT seq_scan + idx_scan FROM pg_stat_user_tables t WHERE t.relid = s.relid),
2
)
ELSE 0
END as usage_ratio
FROM pg_stat_user_indexes s
JOIN pg_index ix ON ix.indexrelid = s.indexrelid
JOIN pg_class i ON i.oid = s.indexrelid
JOIN pg_am am ON am.oid = i.relam
JOIN pg_attribute a ON a.attrelid = s.relid AND a.attnum = ANY(ix.indkey)
`;
const INDEX_STATS_GROUP_ORDER_SQL = `
GROUP BY s.schemaname, s.relname, s.indexrelname, am.amname,
s.indexrelid, s.idx_scan, s.idx_tup_read, s.idx_tup_fetch,
ix.indisunique, ix.indisprimary, s.relid
ORDER BY s.idx_scan DESC
`;
// getStorageStats: tablespace sizes.
const STORAGE_TABLESPACES_SQL = `
SELECT
spcname as name,
pg_tablespace_location(oid) as location,
pg_size_pretty(pg_tablespace_size(oid)) as size,
pg_tablespace_size(oid) as size_bytes,
spcname = 'pg_default' as is_default
FROM pg_tablespace
WHERE spcname NOT LIKE 'pg_global'
`;
// getStorageStats: WAL size (requires superuser; the caller ignores failures).
const STORAGE_WAL_SQL = `
SELECT
pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), '0/0')) as wal_size,
pg_wal_lsn_diff(pg_current_wal_lsn(), '0/0') as wal_size_bytes
`;
// ============================================================================
// Agent read-only execution profile (#328)
// ============================================================================
/**
* The capabilities a read-only TRANSACTION cannot contain, asked of the role
* the profile would run as.
*
* `to_regrole` keeps the query safe on a server where a predefined role is
* absent (it yields NULL, and `COALESCE` makes that a `false` rather than an
* error), so the check does not depend on the server's major version.
*
* Every catalog FUNCTION is schema-qualified because this query decides a
* security boundary. `pg_catalog` is searched implicitly first only while it is
* not named in `search_path`; once it is named explicitly, any schema ahead of
* it shadows built-ins, so `search_path = attacker_schema, pg_catalog` plus a
* shadow `pg_has_role()` would answer four falses for a superuser and defeat
* the one check meant to catch that role. `COALESCE` and `current_user` need no
* qualification (and accept none): they are SQL constructs the parser resolves,
* not functions that name resolution can redirect.
*/
const AGENT_ROLE_PRIVILEGE_SQL = `
SELECT pg_catalog.current_setting('is_superuser') = 'on' AS is_superuser,
COALESCE(
pg_catalog.pg_has_role(current_user, pg_catalog.to_regrole('pg_read_server_files'), 'USAGE'),
false
) AS reads_server_files,
COALESCE(
pg_catalog.pg_has_role(current_user, pg_catalog.to_regrole('pg_write_server_files'), 'USAGE'),
false
) AS writes_server_files,
COALESCE(
pg_catalog.pg_has_role(current_user, pg_catalog.to_regrole('pg_execute_server_program'), 'USAGE'),
false
) AS executes_programs
`;
const AGENT_ROLE_FORBIDDEN_CAPABILITIES = [
"is_superuser",
"reads_server_files",
"writes_server_files",
"executes_programs",
] as const;
/**
* Refuses a role whose privileges reach past the read-only transaction.
*
* `BEGIN READ ONLY` forbids changing the DATABASE. It does not forbid writing
* somewhere else: verified on PostgreSQL 18, a superuser session inside a
* read-only transaction still ran `COPY (…) TO '<path>'` (an arbitrary
* server-side file write), `COPY (…) TO PROGRAM '<cmd>'` (command execution as
* the server's OS user) and `pg_read_file()` (an arbitrary server-side file
* read). Only privileges refuse those — the same lesson the SQLite profile
* learned from `VACUUM INTO`: a control is a claim about one resource, so the
* question is always what else the statement can reach.
*
* Hence a least-privilege agent role is part of this profile's boundary rather
* than a recommendation, and it is VERIFIED at open instead of assumed from
* configuration: an admin can point `agentUser` at a superuser, and a
* connection's own user very often is one.
*
* Fails closed on anything it cannot read as four explicit `false`s — a server
* that answers nothing, or answers something else, leaves the boundary
* unproven.
*/
function assertAgentRoleIsUnprivileged(rows: unknown[]): void {
const row = rows[0] as Record<string, unknown> | undefined;
const held = AGENT_ROLE_FORBIDDEN_CAPABILITIES.filter((capability) => row?.[capability] !== false);
if (held.length > 0) {
throw new ExecutionProfileError(
`The agent read-only execution profile requires a least-privilege PostgreSQL role; this role is unverified or too broad (${held.join(", ")}). A read-only transaction does not stop server-side file access or program execution.`,
"PROFILE_PRIVILEGES_TOO_BROAD",
);
}
}
// ============================================================================
// PostgreSQL Provider
// ============================================================================
export class PostgresProvider extends SQLBaseProvider {
private pool: Pool | null = null;
// Transaction support: dedicated client held outside pool
private txClient: PoolClient | null = null;
private txActive = false;
private txTimeout: ReturnType<typeof setTimeout> | null = null;
private static readonly TX_TIMEOUT_MS = 5 * 60 * 1000; // 5 minutes
/** True when this instance was opened under the agent read-only profile. */
private readonly readOnlyProfile: boolean;
constructor(config: DatabaseConnection, options: ProviderOptions = {}, execution: ProviderExecutionContext = {}) {
super(config, options);
// Server-injected only (see ProviderExecutionContext): the editor path
// builds providers from caller-supplied ProviderOptions, which has no route
// to this flag in either direction.
this.readOnlyProfile = execution.readOnly === true;
this.validate();
}
// ============================================================================
// Provider Metadata
// ============================================================================
public override getCapabilities(): ProviderCapabilities {
return {
...super.getCapabilities(),
defaultPort: 5432,
supportsExplain: true,
explainFormat: "postgres-json",
supportsConnectionString: true,
supportsInlineRowEdit: true,
maintenanceOperations: ["vacuum", "analyze", "reindex", "kill"],
};
}
// ============================================================================
// Validation
// ============================================================================
public validate(): void {
super.validate();
if (!this.config.connectionString) {
if (!this.config.host) {
throw new DatabaseConfigError("Host is required for PostgreSQL", "postgres");
}
if (!this.config.database) {
throw new DatabaseConfigError("Database name is required for PostgreSQL", "postgres");
}
}
}
// ============================================================================
// Connection Management
// ============================================================================
public async connect(): Promise<void> {
if (this.pool) {
return;
}
try {
const poolConfig = this.buildPoolConfig();
this.pool = new Pool(poolConfig);
this.attachPoolErrorListener(this.pool);
const client = await this.pool.connect();
try {
// Under the profile, the role itself is part of the boundary — verify it
// on the same client this connect already borrowed.
if (this.readOnlyProfile) {
assertAgentRoleIsUnprivileged((await client.query(AGENT_ROLE_PRIVILEGE_SQL)).rows);
}
} finally {
client.release();
}
this.setConnected(true);
} catch (error) {
this.setError(error instanceof Error ? error : new Error(String(error)));
// The pool is built before anything that can fail here — the borrow, and
// under the profile the role probe. Whatever went wrong, the caller ends
// up without a usable provider, and `acquireExecutionProfileProvider`
// drops it WITHOUT calling disconnect(), so a pool left open here leaks
// its idle socket and timers with no reference left to close them. End it
// on every failed connect, not just the typed refusal.
const failedPool = this.pool;
this.pool = null;
await failedPool?.end().catch(() => {});
// A typed profile refusal keeps its identity: wrapping it would strip the
// deny reason code callers branch on.
if (error instanceof ExecutionProfileError) {
throw error;
}
throw new ConnectionError(
`Failed to connect to PostgreSQL: ${error instanceof Error ? error.message : error}`,
"postgres",
this.config.host,
this.config.port,
);
}
}
public async disconnect(): Promise<void> {
if (this.pool) {
await this.pool.end();
this.pool = null;
this.setConnected(false);
}
}
/**
* A client that fails while CHECKED OUT rejects its own query, which is why ordinary
* query failures behave correctly. A client that fails while IDLE — the server dropped
* it, the network went away — has no query to reject, so `pg` removes and destroys it
* and then emits `error` on the pool. An `error` event with no listener is an uncaught
* exception, so without this handler a dropped idle connection takes the whole server
* process down (#298).
*
* The client is already gone by the time this fires, so the handler exists to keep the
* event non-fatal and visible — not to reconnect. The pool opens a fresh client on the
* next acquire by itself.
*/
private attachPoolErrorListener(pool: Pool): void {
pool.on("error", (error) => {
console.error("[Postgres] Idle pool client error:", error);
});
}
private buildPoolConfig(): PgPoolConfig {
const sslConfig = this.buildSSLConfig();
const baseConfig: PgPoolConfig = {
min: this.poolConfig.min,
max: this.poolConfig.max,
idleTimeoutMillis: this.poolConfig.idleTimeout,
connectionTimeoutMillis: this.poolConfig.acquireTimeout,
statement_timeout: this.queryTimeout,
ssl: sslConfig,
};
if (this.config.connectionString) {
return {
...baseConfig,
connectionString: this.config.connectionString,
};
}
return {
...baseConfig,
host: this.config.host,
port: this.config.port ?? 5432,
user: this.config.user,
password: this.config.password,
database: this.config.database,
};
}
private buildSSLConfig(): PgPoolConfig["ssl"] {
const connSSL = this.config.ssl;
// Explicit SSL config from connection takes priority
if (connSSL) {
if (connSSL.mode === "disable") return false;
const ssl: Record<string, unknown> = {
rejectUnauthorized: connSSL.mode === "verify-ca" || connSSL.mode === "verify-full",
};
if (connSSL.caCert) ssl.ca = connSSL.caCert;
if (connSSL.clientCert) ssl.cert = connSSL.clientCert;
if (connSSL.clientKey) ssl.key = connSSL.clientKey;
return ssl as PgPoolConfig["ssl"];
}
// Auto-detect for cloud providers
if (this.shouldEnableSSL()) {
return { rejectUnauthorized: false };
}
// Provider options fallback
if (this.options.ssl === false) return false;
return undefined;
}
// ============================================================================
// Query Execution
// ============================================================================
// Track running query PIDs for cancellation
private runningQueryPids = new Map<string, number>();
public async query(sql: string, params?: unknown[], queryId?: string): Promise<QueryResult> {
this.ensureConnected();
return this.trackQuery(async () => {
const { result, executionTime } = await this.measureExecution(async () => {
try {
const client = await this.pool!.connect();
try {
// Track PID for cancellation support
if (queryId) {
const pidRes = await client.query("SELECT pg_backend_pid() as pid");
this.runningQueryPids.set(queryId, pidRes.rows[0].pid);
}
const res = await client.query(sql, params);
return res;
} finally {
if (queryId) this.runningQueryPids.delete(queryId);
client.release();
}
} catch (error) {
if (queryId) this.runningQueryPids.delete(queryId);
throw mapDatabaseError(error, "postgres", sql);
}
});
return {
rows: result.rows,
fields: result.fields?.map((f) => f.name) ?? [],
rowCount: result.rowCount ?? 0,
executionTime,
};
});
}
public async cancelQuery(queryId: string): Promise<boolean> {
const pid = this.runningQueryPids.get(queryId);
if (!pid) return false;
try {
const client = await this.pool!.connect();
try {
const res = await client.query("SELECT pg_cancel_backend($1) as cancelled", [pid]);
return res.rows[0]?.cancelled === true;
} finally {
client.release();
}
} catch (error) {
console.error("[Postgres] Failed to cancel query:", error);
return false;
}
}
// ============================================================================
// Agent Read-Only Execution Profile (#328)
// ============================================================================
/**
* Runs exactly one statement inside `BEGIN READ ONLY` with a
* transaction-local timeout, then rolls back and releases the client. The
* DATABASE is the boundary, twice over:
*
* - The read-only transaction makes the server itself reject any write that
* reaches it (SQLSTATE 25006) — no SQL classification happens here.
* - The statement travels on the extended query protocol (`queryMode:
* "extended"`, pg >= 8.11), whose Parse message the server refuses for
* multi-command strings (SQLSTATE 42601) BEFORE executing anything. That
* is what stops `SELECT 1; COMMIT; INSERT ...` from committing its way out
* of the read-only transaction — on the simple protocol the server would
* execute each command in turn, honoring the smuggled COMMIT.
*
* A single hostile statement cannot escape either: `SET TRANSACTION READ
* WRITE` would be the transaction's only statement before ROLLBACK, a lone
* COMMIT merely ends an empty read-only transaction, and a session-level
* `SET` reverts with the rollback (GUC changes are transactional).
*
* The row/byte caps are enforced result-side after the statement returns;
* the timeout is `SET LOCAL`, so it dies with the transaction.
*/
public async queryReadOnly(sql: string, budget: ReadOnlyStatementBudget): Promise<QueryResult> {
this.ensureConnected();
assertReadOnlyBudget(budget, "postgres");
if (!this.readOnlyProfile) {
// A provider opened outside the profile has had no role verification, so
// its session may be able to write server files or run programs from
// inside a read-only transaction. Refuse rather than serve agent
// semantics without the boundary that makes them true.
throw new QueryError(
"Read-only execution requires a provider opened under the agent read-only profile",
"postgres",
sql,
);
}
return this.trackQuery(async () => {
const { result, executionTime } = await this.measureExecution(async () => {
const client = await this.pool!.connect();
try {
await client.query("BEGIN READ ONLY");
// SET cannot take bind parameters; the value is proven a positive
// integer by assertReadOnlyBudget above, so no text can pass through.
await client.query(`SET LOCAL statement_timeout = ${budget.statementTimeoutMs}`);
// @types/pg does not model queryMode yet; the runtime supports it
// since pg 8.11 (node_modules/pg/lib/query.js requiresPreparation).
const extendedQuery = { text: sql, queryMode: "extended" } as QueryConfig & { queryMode: "extended" };
return await client.query(extendedQuery);
} catch (error) {
throw mapDatabaseError(error, "postgres", sql);
} finally {
// The profile never commits. A client that cannot be reset is
// destroyed (release(error)), never returned to the pool mid-transaction.
try {
await client.query("ROLLBACK");
// Session state a rollback does NOT undo: an advisory lock taken
// inside the transaction survives it (verified on PostgreSQL 18) and
// no statement the agent path admits could release it, so a pooled
// client would carry it into every later execution. DISCARD ALL
// cannot run inside a transaction block, hence after the ROLLBACK.
await client.query("DISCARD ALL");
client.release();
} catch (cleanupError) {
client.release(cleanupError instanceof Error ? cleanupError : new Error(String(cleanupError)));
}
}
});
if (result.rows.length > budget.maxResultRows) {
throw new QueryError(
`Read-only execution exceeded the row budget: ${result.rows.length} rows > ${budget.maxResultRows} allowed`,
"postgres",
sql,
);
}
const resultBytes = measureResultBytes(result.rows);
if (resultBytes > budget.maxResultBytes) {
throw new QueryError(
`Read-only execution exceeded the byte budget: ${resultBytes} bytes > ${budget.maxResultBytes} allowed`,
"postgres",
sql,
);
}
return {
rows: result.rows,
fields: result.fields?.map((f) => f.name) ?? [],
rowCount: result.rowCount ?? 0,
executionTime,
};
});
}
// ============================================================================
// Transaction Support
// ============================================================================
private clearTxTimeout(): void {
if (this.txTimeout) {
clearTimeout(this.txTimeout);
this.txTimeout = null;
}
}
/**
* Force-expire an active transaction (auto-rollback).
* Called by the timeout timer, but also available for testing.
*/
public async expireTransaction(): Promise<void> {
if (this.txActive && this.txClient) {
console.warn("[Postgres] Transaction timed out, auto-rolling back");
try {
await this.txClient.query("ROLLBACK");
} catch {
/* ignore */
} finally {
this.txClient.release();
this.txClient = null;
this.txActive = false;
this.clearTxTimeout();
}
}
}
public async beginTransaction(): Promise<void> {
this.ensureConnected();
if (this.txActive) throw new QueryError("Transaction already active", "postgres");
this.txClient = await this.pool!.connect();
await this.txClient.query("BEGIN");
this.txActive = true;
// Auto-rollback after timeout to prevent leaked locks. Single-line callback
// on purpose: bun lcov attributes a multi-line arrow's opening line as 0-hit.
this.txTimeout = setTimeout(() => void this.expireTransaction(), PostgresProvider.TX_TIMEOUT_MS);
}
public async commitTransaction(): Promise<void> {
if (!this.txClient || !this.txActive) throw new QueryError("No active transaction", "postgres");
this.clearTxTimeout();
try {
await this.txClient.query("COMMIT");
} finally {
this.txClient.release();
this.txClient = null;
this.txActive = false;
}
}
public async rollbackTransaction(): Promise<void> {
if (!this.txClient || !this.txActive) throw new QueryError("No active transaction", "postgres");
this.clearTxTimeout();
try {