1 package sinks
2
3 import (
4 "os"
5 "regexp"
6 "testing"
7
8 "github.com/cybertec-postgresql/pgwatch/v6/internal/testutil"
9
10 "github.com/jackc/pgx/v5"
11 "github.com/stretchr/testify/assert"
12 "github.com/stretchr/testify/require"
13 )
14
15
16
17
18
19
20
21 func TestMigrationsCountInvariant(t *testing.T) {
22 a := assert.New(t)
23
24 registered := registeredMigrationsCount()
25
26
27
28
29
30
31 seeded := len(regexp.MustCompile(`(?m)^\s*\(\d+,\s*'`).FindAllString(sqlMetricAdminSchema, -1))
32 a.Equal(registered, seeded,
33 "registeredMigrationsCount (%d) must equal the number of rows seeded into admin.migration in admin_schema.sql (%d)",
34 registered, seeded)
35 }
36
37
38
39
40
41
42
43
44
45
46
47 func simulatePreV6Migration(t *testing.T, conn *pgx.Conn) {
48 t.Helper()
49 _, err := conn.Exec(ctx, `DELETE FROM admin.migration`)
50 require.NoError(t, err)
51
52 var count int
53 require.NoError(t, conn.QueryRow(ctx, `SELECT count(*) FROM admin.migration`).Scan(&count))
54 require.Equal(t, 0, count, "all migration rows should be wiped after rollback")
55 }
56
57
58
59
60
61
62
63
64
65 func oldSchemaMetricTable(t *testing.T, conn *pgx.Conn, metric string) {
66 t.Helper()
67 _, err := conn.Exec(ctx, `
68 CREATE TABLE public.`+metric+` (LIKE admin.metrics_template INCLUDING INDEXES) PARTITION BY LIST (dbname);
69 COMMENT ON TABLE public.`+metric+` IS 'pgwatch-generated-metric-lvl';
70
71 CREATE TABLE subpartitions.`+metric+`_db1 PARTITION OF public.`+metric+`
72 FOR VALUES IN ('db1') PARTITION BY RANGE (time);
73
74 CREATE TABLE subpartitions.`+metric+`_db1_2024w01 PARTITION OF subpartitions.`+metric+`_db1
75 FOR VALUES FROM ('2024-01-01') TO ('2024-01-08');
76 COMMENT ON TABLE subpartitions.`+metric+`_db1_2024w01 IS 'pgwatch-generated-metric-time-lvl';
77
78 INSERT INTO public.`+metric+` (time, dbname, data) VALUES
79 ('2024-01-03 10:00:00+00', 'db1', '{"x": 1}'::jsonb),
80 ('2024-01-04 11:00:00+00', 'db1', '{"x": 2}'::jsonb);
81 `)
82 require.NoError(t, err)
83 }
84
85
86 func isRangePartitioned(t *testing.T, conn *pgx.Conn, metric string) bool {
87 t.Helper()
88 var ok bool
89 err := conn.QueryRow(ctx,
90 `SELECT EXISTS (SELECT 1 FROM pg_partitioned_table WHERE partrelid = to_regclass($1) AND partstrat = 'r')`,
91 metric).Scan(&ok)
92 require.NoError(t, err)
93 return ok
94 }
95
96
97 func rowCount(t *testing.T, conn *pgx.Conn, metric string) int {
98 t.Helper()
99 var n int
100 require.NoError(t, conn.QueryRow(ctx, `SELECT count(*) FROM public.`+metric).Scan(&n))
101 return n
102 }
103
104
105
106
107
108 func TestMigration01409_TimeOnlyPartitioning(t *testing.T) {
109 if os.Getenv("PGWATCH_TEST_SKIP_MIGRATION") != "" {
110 t.Skip("migration integration test skipped via PGWATCH_TEST_SKIP_MIGRATION")
111 }
112 r := require.New(t)
113 a := assert.New(t)
114
115 pgContainer, pgTearDown, err := testutil.SetupPostgresContainer()
116 r.NoError(err)
117 defer pgTearDown()
118
119 connStr, err := pgContainer.ConnectionString(ctx, "sslmode=disable")
120 r.NoError(err)
121
122
123
124
125 mig, err := NewPostgresSinkMigrator(ctx, connStr)
126 r.NoError(err)
127
128 conn, err := pgx.Connect(ctx, connStr)
129 r.NoError(err)
130 defer conn.Close(ctx)
131
132 simulatePreV6Migration(t, conn)
133
134 const metric = "old_style_metric"
135 oldSchemaMetricTable(t, conn, metric)
136
137
138 a.False(isRangePartitioned(t, conn, metric), "table should start as LIST(dbname) partitioned")
139 a.Equal(2, rowCount(t, conn, metric), "seeded rows should be present before migration")
140
141
142 r.NoError(mig.Migrate())
143
144
145 a.True(isRangePartitioned(t, conn, metric), "table should be converted to RANGE(time) partitioning")
146 a.Equal(2, rowCount(t, conn, metric), "data must be preserved across the migration")
147
148
149 var leftover bool
150 r.NoError(conn.QueryRow(ctx,
151 `SELECT to_regclass($1) IS NOT NULL`, metric+"_before_v6_migration").Scan(&leftover))
152 a.False(leftover, "the *_before_v6_migration scratch table should be dropped")
153
154
155 needs, err := mig.NeedsMigration()
156 r.NoError(err)
157 a.False(needs, "no migrations should be pending immediately after a successful migrate")
158
159 r.NoError(mig.Migrate(), "re-running Migrate() must be a no-op and not error")
160 a.True(isRangePartitioned(t, conn, metric))
161 a.Equal(2, rowCount(t, conn, metric), "data must remain intact after a second migrate")
162 }
163
164
165
166 func TestMigration01409_EmptyTable(t *testing.T) {
167 if os.Getenv("PGWATCH_TEST_SKIP_MIGRATION") != "" {
168 t.Skip("migration integration test skipped via PGWATCH_TEST_SKIP_MIGRATION")
169 }
170 r := require.New(t)
171 a := assert.New(t)
172
173 pgContainer, pgTearDown, err := testutil.SetupPostgresContainer()
174 r.NoError(err)
175 defer pgTearDown()
176
177 connStr, err := pgContainer.ConnectionString(ctx, "sslmode=disable")
178 r.NoError(err)
179
180 mig, err := NewPostgresSinkMigrator(ctx, connStr)
181 r.NoError(err)
182
183 conn, err := pgx.Connect(ctx, connStr)
184 r.NoError(err)
185 defer conn.Close(ctx)
186
187 simulatePreV6Migration(t, conn)
188
189 const metric = "empty_old_metric"
190 _, err = conn.Exec(ctx, `
191 CREATE TABLE public.`+metric+` (LIKE admin.metrics_template INCLUDING INDEXES) PARTITION BY LIST (dbname);
192 COMMENT ON TABLE public.`+metric+` IS 'pgwatch-generated-metric-lvl';
193 CREATE TABLE subpartitions.`+metric+`_db1 PARTITION OF public.`+metric+`
194 FOR VALUES IN ('db1') PARTITION BY RANGE (time);
195 CREATE TABLE subpartitions.`+metric+`_db1_2024w01 PARTITION OF subpartitions.`+metric+`_db1
196 FOR VALUES FROM ('2024-01-01') TO ('2024-01-08');
197 COMMENT ON TABLE subpartitions.`+metric+`_db1_2024w01 IS 'pgwatch-generated-metric-time-lvl';
198 `)
199 r.NoError(err)
200
201 r.NoError(mig.Migrate())
202
203 a.True(isRangePartitioned(t, conn, metric), "empty table should still be converted to RANGE(time)")
204 a.Equal(0, rowCount(t, conn, metric))
205
206
207 var leaves int
208 r.NoError(conn.QueryRow(ctx,
209 `SELECT count(*) FROM pg_partition_tree($1) WHERE isleaf`, metric).Scan(&leaves))
210 a.GreaterOrEqual(leaves, 1, "a single empty time partition should exist")
211 }
212
213
214
215
216
217 func TestMigration01474_DropAllMetricTablesProcedure(t *testing.T) {
218 if os.Getenv("PGWATCH_TEST_SKIP_MIGRATION") != "" {
219 t.Skip("migration integration test skipped via PGWATCH_TEST_SKIP_MIGRATION")
220 }
221 r := require.New(t)
222 a := assert.New(t)
223
224 pgContainer, pgTearDown, err := testutil.SetupPostgresContainer()
225 r.NoError(err)
226 defer pgTearDown()
227
228 connStr, err := pgContainer.ConnectionString(ctx, "sslmode=disable")
229 r.NoError(err)
230
231 mig, err := NewPostgresSinkMigrator(ctx, connStr)
232 r.NoError(err)
233
234 conn, err := pgx.Connect(ctx, connStr)
235 r.NoError(err)
236 defer conn.Close(ctx)
237
238
239
240 _, err = conn.Exec(ctx, `
241 DROP PROCEDURE IF EXISTS admin.drop_all_metric_tables();
242 CREATE FUNCTION admin.drop_all_metric_tables() RETURNS int AS $$ BEGIN RETURN 0; END $$ LANGUAGE plpgsql;
243 DELETE FROM admin.migration WHERE id >= 3;
244 `)
245 r.NoError(err)
246
247
248 var prokindBefore string
249 r.NoError(conn.QueryRow(ctx, `
250 SELECT p.prokind FROM pg_proc p JOIN pg_namespace n ON n.oid = p.pronamespace
251 WHERE n.nspname = 'admin' AND p.proname = 'drop_all_metric_tables'
252 `).Scan(&prokindBefore))
253 a.Equal("f", prokindBefore, "legacy routine must start as a function")
254
255
256 _, err = conn.Exec(ctx, `
257 CREATE TABLE public.upgraded_metric (LIKE admin.metrics_template INCLUDING INDEXES) PARTITION BY RANGE (time);
258 COMMENT ON TABLE public.upgraded_metric IS 'pgwatch-generated-metric-lvl';
259 INSERT INTO admin.all_distinct_dbname_metrics (dbname, metric) VALUES ('db1', 'upgraded_metric');
260 CREATE TABLE subpartitions.upgraded_metric_2024w01 PARTITION OF public.upgraded_metric FOR VALUES FROM ('2024-01-01') TO ('2024-01-08');
261 COMMENT ON TABLE subpartitions.upgraded_metric_2024w01 IS 'pgwatch-generated-metric-time-lvl';
262 `)
263 r.NoError(err)
264
265 r.NoError(mig.Migrate())
266
267
268 var prokindAfter string
269 r.NoError(conn.QueryRow(ctx, `
270 SELECT p.prokind FROM pg_proc p JOIN pg_namespace n ON n.oid = p.pronamespace
271 WHERE n.nspname = 'admin' AND p.proname = 'drop_all_metric_tables'
272 `).Scan(&prokindAfter))
273 a.Equal("p", prokindAfter, "admin.drop_all_metric_tables must be upgraded to a procedure")
274
275
276
277
278 _, err = conn.Exec(ctx, "CALL admin.drop_all_metric_tables();")
279 r.NoError(err, "upgraded procedure must execute end-to-end")
280
281 var leftover *string
282 r.NoError(conn.QueryRow(ctx, "SELECT to_regclass('public.upgraded_metric')::text").Scan(&leftover))
283 a.Nil(leftover, "top-level metric table must be dropped by the upgraded procedure")
284 var subpartCount int
285 r.NoError(conn.QueryRow(ctx, "SELECT count(*) FROM pg_class WHERE relnamespace = 'subpartitions'::regnamespace").Scan(&subpartCount))
286 a.Equal(0, subpartCount, "partition must be dropped by the upgraded procedure")
287 var listingCount int
288 r.NoError(conn.QueryRow(ctx, "SELECT count(*) FROM admin.all_distinct_dbname_metrics").Scan(&listingCount))
289 a.Equal(0, listingCount, "listing table must be truncated by the upgraded procedure")
290
291 needs, err := mig.NeedsMigration()
292 r.NoError(err)
293 a.False(needs, "no migrations should be pending after the upgrade")
294 }
295
296
297
298
299
300
301
302
303
304 func TestMigration_AllMigrationsRunFromEmpty(t *testing.T) {
305 if os.Getenv("PGWATCH_TEST_SKIP_MIGRATION") != "" {
306 t.Skip("migration integration test skipped via PGWATCH_TEST_SKIP_MIGRATION")
307 }
308 r := require.New(t)
309 a := assert.New(t)
310
311 pgContainer, pgTearDown, err := testutil.SetupPostgresContainer()
312 r.NoError(err)
313 defer pgTearDown()
314
315 connStr, err := pgContainer.ConnectionString(ctx, "sslmode=disable")
316 r.NoError(err)
317
318 mig, err := NewPostgresSinkMigrator(ctx, connStr)
319 r.NoError(err)
320
321 conn, err := pgx.Connect(ctx, connStr)
322 r.NoError(err)
323 defer conn.Close(ctx)
324
325 simulatePreV6Migration(t, conn)
326
327
328 var before int
329 r.NoError(conn.QueryRow(ctx, "SELECT count(*) FROM admin.migration").Scan(&before))
330 a.Equal(0, before)
331
332 r.NoError(mig.Migrate(), "migrate from empty must run the entire chain without error")
333
334
335 var after int
336 r.NoError(conn.QueryRow(ctx, "SELECT count(*) FROM admin.migration").Scan(&after))
337 a.Equal(registeredMigrationsCount(), after,
338 "every registered migration (incl. 01474) must record itself in admin.migration after a wipe + migrate")
339
340
341 var prokind string
342 r.NoError(conn.QueryRow(ctx, `
343 SELECT p.prokind FROM pg_proc p JOIN pg_namespace n ON n.oid = p.pronamespace
344 WHERE n.nspname = 'admin' AND p.proname = 'drop_all_metric_tables'
345 `).Scan(&prokind))
346 a.Equal("p", prokind, "01474 must install the routine as a procedure")
347
348
349 var fnExists bool
350 r.NoError(conn.QueryRow(ctx, `
351 SELECT EXISTS (SELECT 1 FROM pg_proc p JOIN pg_namespace n ON n.oid = p.pronamespace
352 WHERE n.nspname = 'admin' AND p.proname = 'ensure_partition_metric_time')
353 `).Scan(&fnExists))
354 a.True(fnExists, "01180 must install admin.ensure_partition_metric_time")
355
356
357 needs, err := mig.NeedsMigration()
358 r.NoError(err)
359 a.False(needs, "no migrations should be pending after a full replay")
360
361
362 r.NoError(mig.Migrate(), "second migrate after full replay must be a no-op")
363 }
364