-
Notifications
You must be signed in to change notification settings - Fork 1
Expand file tree
/
Copy pathcreate_integration_test.go
More file actions
899 lines (816 loc) · 36.6 KB
/
Copy pathcreate_integration_test.go
File metadata and controls
899 lines (816 loc) · 36.6 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
package executor_test
import (
"context"
"errors"
"fmt"
"sync"
"testing"
"time"
"github.com/jackc/pgx/v5"
"github.com/jackc/pgx/v5/pgxpool"
"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
"github.com/block/pg-sprite/internal/testutil"
"github.com/block/pg-sprite/pkg/dbconn"
"github.com/block/pg-sprite/pkg/executor"
"github.com/block/pg-sprite/pkg/preflight"
"github.com/block/pg-sprite/pkg/progress"
"github.com/block/pg-sprite/pkg/statement"
)
// createFixture is one schema on a real server with an absence proof and
// a creation-privilege proof minted for the named table — the inputs
// ExecuteCreate requires.
type createFixture struct {
pool *pgxpool.Pool
schema string
at preflight.AbsentTarget
cr preflight.CreationRole
}
func newCreateFixture(t *testing.T, table string) createFixture {
t.Helper()
pool, err := dbconn.NewPool(t.Context(), dbconn.Config{URL: testutil.StartPostgres(t)})
require.NoError(t, err)
t.Cleanup(pool.Close)
schema := testutil.NewSchema(t, pool)
at, err := preflight.CheckTableAbsent(t.Context(), pool, schema, table)
require.NoError(t, err)
cr, err := preflight.CheckCreatePrivileges(t.Context(), pool, schema)
require.NoError(t, err)
return createFixture{pool: pool, schema: schema, at: at, cr: cr}
}
func desired(t *testing.T, sql string) statement.DesiredSchema {
t.Helper()
ds, err := statement.ParseDesired(sql)
require.NoError(t, err)
return ds
}
// relationKind returns the pg_class relkind of schema.name, or "" when no
// relation owns the name — the catalog oracle for what a create run left.
func relationKind(t *testing.T, pool *pgxpool.Pool, schema, name string) string {
t.Helper()
var relkind *string
require.NoError(t, pool.QueryRow(t.Context(),
`SELECT (SELECT c.relkind::text
FROM pg_class c
JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE n.nspname = $1 AND c.relname = $2)`,
schema, name).Scan(&relkind))
if relkind == nil {
return ""
}
return *relkind
}
// relationExists reports whether any relation owns schema.name, for
// assertions where only presence matters.
func relationExists(t *testing.T, pool *pgxpool.Pool, schema, name string) bool {
t.Helper()
var exists bool
require.NoError(t, pool.QueryRow(t.Context(),
`SELECT EXISTS (
SELECT FROM pg_class c
JOIN pg_namespace n ON n.oid = c.relnamespace
WHERE n.nspname = $1 AND c.relname = $2)`,
schema, name).Scan(&exists))
return exists
}
func TestExecuteCreateRunsTableAndIndexes(t *testing.T) {
f := newCreateFixture(t, "t")
ds := desired(t, `
CREATE TABLE t (id int PRIMARY KEY, name text);
CREATE INDEX t_name_idx ON t (name);
CREATE UNIQUE INDEX t_id_name_idx ON t (id, name);
`)
rep, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, ds, createBudget, executor.DefaultRetryPolicy())
require.NoError(t, err)
assert.Equal(t, "r", relationKind(t, f.pool, f.schema, "t"))
assert.Equal(t, "i", relationKind(t, f.pool, f.schema, "t_name_idx"))
assert.Equal(t, "i", relationKind(t, f.pool, f.schema, "t_id_name_idx"))
require.Len(t, rep.Steps, 3)
assert.Contains(t, rep.Steps[0].SQL, "CREATE TABLE")
for _, step := range rep.Steps {
assert.Equal(t, executor.StepBrief, step.Kind)
assert.GreaterOrEqual(t, step.Duration, time.Duration(0))
}
}
// A desired file may state its index before its table — declarative input
// carries no ordering contract — but an index cannot be built before its
// table exists, so the executor orders the CREATE TABLE first.
func TestExecuteCreateOrdersTableBeforeIndexes(t *testing.T) {
f := newCreateFixture(t, "t")
ds := desired(t, `
CREATE INDEX t_name_idx ON t (name);
CREATE TABLE t (id int, name text);
`)
rep, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, ds, createBudget, executor.DefaultRetryPolicy())
require.NoError(t, err)
require.Len(t, rep.Steps, 2)
assert.Contains(t, rep.Steps[0].SQL, "CREATE TABLE")
assert.Equal(t, "i", relationKind(t, f.pool, f.schema, "t_name_idx"))
}
func TestExecuteCreateWithProgressReportsQualifiedStepStatementsInOrder(t *testing.T) {
f := newCreateFixture(t, "t")
ds := desired(t, `
CREATE INDEX t_name_idx ON t (name);
CREATE TABLE t (id int, name text);
`)
tracker, err := progress.NewTracker(progress.WallClock{})
require.NoError(t, err)
// Hold each of this test's DDL statements long enough for the tracker
// to be observed mid-step.
testutil.InstallEventTrigger(t, f.pool, testutil.DDLCommandStart, f.schema, "delay_create_progress", fmt.Sprintf(`
IF current_query() LIKE '%%%s%%' THEN
PERFORM pg_sleep(0.25);
END IF;`, f.schema))
type result struct {
rep executor.SequenceReport
err error
}
results := make(chan result, 1)
var workers sync.WaitGroup
workers.Go(func() {
rep, executeErr := executor.ExecuteCreateWithProgress(t.Context(), f.pool, f.at, f.cr, ds,
createBudget, executor.DefaultRetryPolicy(), tracker)
results <- result{rep: rep, err: executeErr}
})
t.Cleanup(workers.Wait)
var observed []string
require.Eventually(t, func() bool {
snapshot, progressErr := tracker.Progress(t.Context())
if progressErr != nil || snapshot.Detail.Statement == "" {
return false
}
if len(observed) == 0 || observed[len(observed)-1] != snapshot.Detail.Statement {
observed = append(observed, snapshot.Detail.Statement)
}
return len(observed) == 2
}, 5*time.Second, 10*time.Millisecond, "both active create steps must publish their statements in order")
execution := <-results
workers.Wait()
require.NoError(t, execution.err)
require.Len(t, execution.rep.Steps, 2)
assert.Equal(t, []string{
fmt.Sprintf("CREATE TABLE %s.t (id int, name text)", f.schema),
fmt.Sprintf("CREATE INDEX t_name_idx ON %s.t USING btree (name)", f.schema),
}, observed)
}
// The absence proof is time-of-check: a create that takes the name after
// the check surfaces as the typed collision, and the caller re-diffs
// rather than assuming what the occupant looks like.
func TestExecuteCreateReportsCollisionAsTyped(t *testing.T) {
f := newCreateFixture(t, "t")
_, err := f.pool.Exec(t.Context(), fmt.Sprintf("CREATE TABLE %s.t (other int)", f.schema))
require.NoError(t, err)
ds := desired(t, "CREATE TABLE t (id int)")
rep, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, ds, createBudget, executor.DefaultRetryPolicy())
require.Error(t, err)
var stepErr *executor.SequenceStepError
require.ErrorAs(t, err, &stepErr)
assert.Equal(t, 1, stepErr.Step)
assert.ErrorIs(t, err, executor.ErrCreateCollision)
assert.Equal(t, executor.CodeCreateCollision, executor.OutcomeCode(err))
assert.Empty(t, rep.Steps)
}
// The first-choice name of an index-backed constraint is part of the
// desired set's contract. An unrelated catalog occupant must refuse the
// whole set rather than make PostgreSQL suffix the constraint name.
func TestExecuteCreateRefusesOccupiedImplicitIndexNameBeforeExecution(t *testing.T) {
f := newCreateFixture(t, "t")
_, err := f.pool.Exec(t.Context(), fmt.Sprintf("CREATE TABLE %s.other (v int)", f.schema))
require.NoError(t, err)
_, err = f.pool.Exec(t.Context(), fmt.Sprintf("CREATE INDEX t_pkey ON %s.other (v)", f.schema))
require.NoError(t, err)
ds := desired(t, "CREATE TABLE t (id int PRIMARY KEY, v text)")
rep, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, ds, createBudget, executor.DefaultRetryPolicy())
require.ErrorIs(t, err, executor.ErrCreateCollision)
assert.Empty(t, rep.Steps)
assert.False(t, relationExists(t, f.pool, f.schema, "t"), "the catalog preflight runs before every step")
}
// A column-owned sequence's first-choice name is part of the desired set's
// contract. An occupant must refuse the whole set rather than make
// PostgreSQL silently suffix the sequence name.
func TestExecuteCreateRefusesOccupiedImplicitSequenceNameBeforeExecution(t *testing.T) {
tests := []struct {
name string
sql string
}{
{name: "serial", sql: "CREATE TABLE t (id serial PRIMARY KEY)"},
{name: "identity", sql: "CREATE TABLE t (id bigint GENERATED BY DEFAULT AS IDENTITY)"},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
f := newCreateFixture(t, "t")
_, err := f.pool.Exec(t.Context(), fmt.Sprintf("CREATE SEQUENCE %s.t_id_seq", f.schema))
require.NoError(t, err)
rep, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, desired(t, tt.sql), createBudget, executor.DefaultRetryPolicy())
require.ErrorIs(t, err, executor.ErrCreateCollision)
assert.Empty(t, rep.Steps)
assert.False(t, relationExists(t, f.pool, f.schema, "t"), "the catalog preflight runs before every step")
})
}
}
// A claimed-name probe that cannot complete says nothing about whether the
// names are free: it is the caller's operational failure, never a
// collision, so the caller retries rather than being told a free name is
// taken.
func TestExecuteCreateProbeFailureIsNotACollision(t *testing.T) {
f := newCreateFixture(t, "t")
ds := desired(t, "CREATE TABLE t (id int PRIMARY KEY, v text)")
ctx, cancel := context.WithCancel(t.Context())
cancel()
rep, err := executor.ExecuteCreate(ctx, f.pool, f.at, f.cr, ds, createBudget, executor.DefaultRetryPolicy())
require.ErrorIs(t, err, context.Canceled)
assert.NotErrorIs(t, err, executor.ErrCreateCollision)
var stepErr *executor.SequenceStepError
assert.False(t, errors.As(err, &stepErr), "nothing started, so there is no step to blame")
assert.Empty(t, rep.Steps)
assert.False(t, relationExists(t, f.pool, f.schema, "t"))
}
// A failed step ends the run; the steps before it committed and remain,
// and the report covers exactly that prefix so the caller can disclose
// what already happened. An index on a column the table does not have
// passes admission — admission checks shape and target, not column
// existence — and fails only when the server executes it.
func TestExecuteCreateFailedStepKeepsCommittedPrefix(t *testing.T) {
f := newCreateFixture(t, "t")
ds := desired(t, `
CREATE TABLE t (id int, name text);
CREATE INDEX t_id_idx ON t (id);
CREATE INDEX t_missing_idx ON t (missing);
`)
rep, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, ds, createBudget, executor.DefaultRetryPolicy())
require.Error(t, err)
var stepErr *executor.SequenceStepError
require.ErrorAs(t, err, &stepErr)
assert.Equal(t, 3, stepErr.Step)
assert.Equal(t, 3, stepErr.Total)
// The server error is not a collision; it passes through untyped.
assert.NotErrorIs(t, err, executor.ErrCreateCollision)
assert.Equal(t, executor.CodeExecutionFailed, executor.OutcomeCode(err))
assert.True(t, relationExists(t, f.pool, f.schema, "t"))
assert.True(t, relationExists(t, f.pool, f.schema, "t_id_idx"))
require.Len(t, rep.Steps, 2)
// The committed prefix is the rerun contract: the absence check now
// refuses, which is the declarative front door's signal to re-diff.
_, err = preflight.CheckTableAbsent(t.Context(), f.pool, f.schema, "t")
assert.ErrorIs(t, err, preflight.ErrRelationExists)
}
// A name claimed twice within the desired set is decidable at admission,
// so the whole set refuses before the first step runs — never a mid-run
// failure with a committed prefix.
func TestExecuteCreateRefusesDuplicateNamesAtAdmission(t *testing.T) {
tests := []struct {
name string
sql string
}{
{
name: "two indexes under one name",
sql: `CREATE TABLE t (id int, name text);
CREATE INDEX dup_idx ON t (id);
CREATE INDEX dup_idx ON t (name)`,
},
{
name: "index named after the table",
sql: `CREATE TABLE t (id int);
CREATE INDEX t ON t (id)`,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
f := newCreateFixture(t, "t")
ds := desired(t, tt.sql)
_, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, ds, createBudget, executor.DefaultRetryPolicy())
require.ErrorIs(t, err, executor.ErrDuplicateCreateName)
assert.Equal(t, executor.CodeDuplicateCreateName, executor.OutcomeCode(err))
assert.False(t, relationExists(t, f.pool, f.schema, "t"),
"admission covers the whole set before the first step executes")
})
}
}
// A standalone type occupying the table's name raises a different SQLSTATE
// than a relation would — every table also mints a composite type — and
// still surfaces as the typed collision.
func TestExecuteCreateReportsTypeCollisionAsTyped(t *testing.T) {
f := newCreateFixture(t, "t")
_, err := f.pool.Exec(t.Context(), fmt.Sprintf("CREATE TYPE %s.t AS ENUM ('a')", f.schema))
require.NoError(t, err)
ds := desired(t, "CREATE TABLE t (id int)")
_, err = executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, ds, createBudget, executor.DefaultRetryPolicy())
require.ErrorIs(t, err, executor.ErrCreateCollision)
assert.Equal(t, executor.CodeCreateCollision, executor.OutcomeCode(err))
}
func TestExecuteCreateAdmissionRefusals(t *testing.T) {
tests := []struct {
name string
sql string
wantErr error
}{
{
name: "if not exists on the table",
sql: "CREATE TABLE IF NOT EXISTS t (id int)",
wantErr: executor.ErrIfNotExistsUnsupported,
},
{
name: "if not exists on an index",
sql: "CREATE TABLE t (id int); CREATE INDEX IF NOT EXISTS t_idx ON t (id)",
wantErr: executor.ErrIfNotExistsUnsupported,
},
// INHERITS, LIKE, and OF bind to a secondary relation or type the
// qualification never touches: the name resolves via search_path
// to an existing object the absence proof says nothing about.
{
name: "inherits from an existing parent",
sql: "CREATE TABLE t (id int) INHERITS (parent)",
wantErr: executor.ErrUnsupportedCreateStep,
},
{
name: "like an existing source table",
sql: "CREATE TABLE t (LIKE src INCLUDING ALL)",
wantErr: executor.ErrUnsupportedCreateStep,
},
{
name: "of an existing composite type",
sql: "CREATE TABLE t OF ty",
wantErr: executor.ErrUnsupportedCreateStep,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
f := newCreateFixture(t, "t")
ds := desired(t, tt.sql)
_, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, ds, createBudget, executor.DefaultRetryPolicy())
require.ErrorIs(t, err, tt.wantErr)
assert.False(t, relationExists(t, f.pool, f.schema, "t"),
"admission covers the whole set before the first step executes")
})
}
}
func TestExecuteCreateRefusesPartitionOf(t *testing.T) {
f := newCreateFixture(t, "t_part")
_, err := f.pool.Exec(t.Context(),
fmt.Sprintf("CREATE TABLE %s.parent (id int) PARTITION BY RANGE (id)", f.schema))
require.NoError(t, err)
ds := desired(t, "CREATE TABLE t_part PARTITION OF parent FOR VALUES FROM (1) TO (10)")
_, err = executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, ds, createBudget, executor.DefaultRetryPolicy())
require.ErrorIs(t, err, executor.ErrPartitionOfUnsupported)
assert.Equal(t, executor.CodePartitionOfUnsupported, executor.OutcomeCode(err))
}
// A desired schema for one table can never run against a proof minted for
// another: the mismatch is an invariant breach, not a refusal.
func TestExecuteCreateRefusesProofTargetMismatch(t *testing.T) {
f := newCreateFixture(t, "other")
ds := desired(t, "CREATE TABLE t (id int)")
_, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, ds, createBudget, executor.DefaultRetryPolicy())
require.ErrorIs(t, err, executor.ErrInvariantViolation)
assert.False(t, relationExists(t, f.pool, f.schema, "t"))
}
// Create steps run with search_path pinned to the proof's schema then
// public — the same policy the introspection read path sets — so a
// desired file's unqualified type reference resolves in the target
// schema, and resolves there even when public holds a type of the same
// name. Without the pin the steps would run under the session default and
// the target schema's type would be invisible (SQLSTATE 42704).
func TestExecuteCreateResolvesTypesInTargetSchema(t *testing.T) {
f := newCreateFixture(t, "t")
// The type lives in the target schema and, under a unique name, in
// public too — resolution must pick the target schema's copy.
typeName := f.schema + "_mood"
_, err := f.pool.Exec(t.Context(), fmt.Sprintf(
"CREATE TYPE %s.%s AS ENUM ('happy', 'sad')", f.schema, typeName))
require.NoError(t, err)
_, err = f.pool.Exec(t.Context(), fmt.Sprintf(
"CREATE TYPE public.%s AS ENUM ('decoy')", typeName))
require.NoError(t, err)
t.Cleanup(func() {
_, err := f.pool.Exec(context.WithoutCancel(t.Context()),
fmt.Sprintf("DROP TYPE IF EXISTS public.%s", typeName))
assert.NoError(t, err)
})
ds := desired(t, fmt.Sprintf("CREATE TABLE t (id int, m %s)", typeName))
_, err = executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, ds, createBudget, executor.DefaultRetryPolicy())
require.NoError(t, err)
var udtSchema string
require.NoError(t, f.pool.QueryRow(t.Context(),
`SELECT udt_schema FROM information_schema.columns
WHERE table_schema = $1 AND table_name = 't' AND column_name = 'm'`,
f.schema).Scan(&udtSchema))
assert.Equal(t, f.schema, udtSchema,
"the column's type must resolve in the proof's schema, not public")
}
// An explicit CREATE INDEX whose name is the first choice of an implicit
// constraint index is a decidable conflict: admission refuses the whole
// set before anything runs, rather than letting the server suffix its way
// around one name or fail mid-run after the table committed.
func TestExecuteCreateRefusesImplicitIndexNameCollision(t *testing.T) {
tests := []struct {
name string
sql string
}{
{
name: "explicit index named after the primary key's index",
sql: `CREATE TABLE t (id int PRIMARY KEY);
CREATE INDEX t_pkey ON t (id);`,
},
{
name: "explicit index named after a unique constraint's index",
sql: `CREATE TABLE t (a int, b int, UNIQUE (a, b));
CREATE INDEX t_a_b_key ON t (a);`,
},
{
name: "explicit index named after a named constraint",
sql: `CREATE TABLE t (id int, CONSTRAINT my_uni UNIQUE (id));
CREATE INDEX my_uni ON t (id);`,
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
f := newCreateFixture(t, "t")
ds := desired(t, tt.sql)
_, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, ds, createBudget, executor.DefaultRetryPolicy())
require.ErrorIs(t, err, executor.ErrDuplicateCreateName)
assert.Equal(t, executor.CodeDuplicateCreateName, executor.OutcomeCode(err))
assert.False(t, relationExists(t, f.pool, f.schema, "t"),
"admission covers the whole set before the first step executes")
})
}
}
// A zero CreationRole is forgeable by any package: only
// CheckCreatePrivileges mints one with a schema, so the executor refuses
// it as an invariant breach before anything runs.
func TestExecuteCreateRefusesZeroCreationRole(t *testing.T) {
f := newCreateFixture(t, "t")
ds := desired(t, "CREATE TABLE t (id int)")
_, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, preflight.CreationRole{}, ds, createBudget, executor.DefaultRetryPolicy())
require.ErrorIs(t, err, executor.ErrInvariantViolation)
assert.False(t, relationExists(t, f.pool, f.schema, "t"))
}
// A creation-privilege proof minted for one schema can never authorize a
// run whose absence proof names another: the mismatch is an invariant
// breach, not a refusal.
func TestExecuteCreateRefusesCreationRoleSchemaMismatch(t *testing.T) {
f := newCreateFixture(t, "t")
otherSchema := testutil.NewSchema(t, f.pool)
otherCR, err := preflight.CheckCreatePrivileges(t.Context(), f.pool, otherSchema)
require.NoError(t, err)
ds := desired(t, "CREATE TABLE t (id int)")
_, err = executor.ExecuteCreate(t.Context(), f.pool, f.at, otherCR, ds, createBudget, executor.DefaultRetryPolicy())
require.ErrorIs(t, err, executor.ErrInvariantViolation)
assert.False(t, relationExists(t, f.pool, f.schema, "t"))
}
// Unnamed CREATE INDEX steps claim no name — the server invents one,
// suffixing around occupants — so two of them in one desired set are not
// a duplicate-name conflict.
func TestExecuteCreateAllowsMultipleUnnamedIndexes(t *testing.T) {
f := newCreateFixture(t, "t")
ds := desired(t, `
CREATE TABLE t (a int, b int);
CREATE INDEX ON t (a);
CREATE INDEX ON t (b);
`)
rep, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, ds, createBudget, executor.DefaultRetryPolicy())
require.NoError(t, err)
require.Len(t, rep.Steps, 3)
var indexes int
require.NoError(t, f.pool.QueryRow(t.Context(),
`SELECT count(*) FROM pg_indexes WHERE schemaname = $1 AND tablename = 't'`,
f.schema).Scan(&indexes))
assert.Equal(t, 2, indexes, "both server-named indexes exist")
}
// A first-choice name taken after the probe but before the server picks
// names is the one race the probe cannot see. The CREATE TABLE has
// committed by then, so the outcome is a step-1 failure that names the
// suffixed replacement and leaves the born table in place — never a
// collision refusal, which would claim nothing ran, and never a drop of a
// table the executor cannot prove nobody has started to use. The remedy
// the code prescribes — free the first-choice name and rename the suffixed
// relation onto it — must leave the table owning exactly what was claimed.
// The proof is part of step 1: the steps after the CREATE TABLE never run
// on a table whose names are in question.
func TestExecuteCreateReportsNameTakenInsideProbeWindowAsMismatch(t *testing.T) {
tests := []struct {
name string
sql string
total int
setupSQL string
occupant string
missing []string
unclaimed []string
// unbuilt is an index step after the CREATE TABLE that must not
// have run.
unbuilt string
// repairSQL is the operator's remedy: drop the occupant and rename
// the server's suffixed relation onto the first-choice name.
repairSQL []string
wantOwned preflight.OwnedRelationNames
}{
{
name: "constraint index",
sql: `
CREATE TABLE t (id int PRIMARY KEY, v text);
CREATE INDEX t_v_idx ON t (v);`,
total: 2,
setupSQL: "CREATE TABLE %[1]s.other (v int)",
occupant: "CREATE INDEX t_pkey ON %[1]s.other (v)",
missing: []string{"t_pkey"},
unclaimed: []string{"t_pkey1"},
unbuilt: "t_v_idx",
repairSQL: []string{
"DROP INDEX %[1]s.t_pkey",
"ALTER INDEX %[1]s.t_pkey1 RENAME TO t_pkey",
},
wantOwned: preflight.OwnedRelationNames{ConstraintIndexes: []string{"t_pkey"}, Sequences: []string{}},
},
{
name: "column-owned sequence",
sql: "CREATE TABLE t (id serial PRIMARY KEY)",
total: 1,
occupant: "CREATE SEQUENCE %[1]s.t_id_seq",
missing: []string{"t_id_seq"},
unclaimed: []string{"t_id_seq1"},
repairSQL: []string{
"DROP SEQUENCE %[1]s.t_id_seq",
"ALTER SEQUENCE %[1]s.t_id_seq1 RENAME TO t_id_seq",
},
wantOwned: preflight.OwnedRelationNames{ConstraintIndexes: []string{"t_pkey"}, Sequences: []string{"t_id_seq"}},
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
f := newCreateFixture(t, "t")
if tt.setupSQL != "" {
_, err := f.pool.Exec(t.Context(), fmt.Sprintf(tt.setupSQL, f.schema))
require.NoError(t, err)
}
testutil.RunDuringDDL(t, f.pool, testutil.DDLCommandStart, "CREATE TABLE", f.schema, "t",
fmt.Sprintf(tt.occupant, f.schema))
rep, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, desired(t, tt.sql), createBudget, executor.DefaultRetryPolicy())
require.Error(t, err)
var stepErr *executor.SequenceStepError
require.ErrorAs(t, err, &stepErr)
assert.Equal(t, 1, stepErr.Step)
assert.Equal(t, tt.total, stepErr.Total)
assert.Equal(t, executor.StepBrief, stepErr.Kind)
assert.Contains(t, stepErr.SQL, "CREATE TABLE")
assert.ErrorIs(t, err, executor.ErrCreateNameMismatch)
assert.NotErrorIs(t, err, executor.ErrCreateCollision,
"a collision says nothing ran; here the CREATE TABLE committed")
assert.Equal(t, executor.CodeCreateNameMismatch, executor.OutcomeCode(err))
var mismatch *executor.CreateNameMismatchError
require.ErrorAs(t, err, &mismatch)
assert.Equal(t, f.schema, mismatch.Schema)
assert.Equal(t, "t", mismatch.Table)
assert.Equal(t, tt.missing, mismatch.Missing)
assert.Equal(t, tt.unclaimed, mismatch.Unclaimed)
assert.Empty(t, rep.Steps, "the failed step's own state is named by the code, not the committed prefix")
assert.Equal(t, "r", relationKind(t, f.pool, f.schema, "t"), "the born table is left in place")
for _, name := range tt.unclaimed {
assert.True(t, relationExists(t, f.pool, f.schema, name), "the server's suffixed replacement %s remains for the operator to rename", name)
}
if tt.unbuilt != "" {
assert.False(t, relationExists(t, f.pool, f.schema, tt.unbuilt), "the run stops at step 1; %s is never built on a table whose names are unproven", tt.unbuilt)
}
// The operator follows the code's remedy; the table then owns
// exactly the first-choice names the desired file claimed.
for _, repair := range tt.repairSQL {
_, err := f.pool.Exec(t.Context(), fmt.Sprintf(repair, f.schema))
require.NoError(t, err)
}
owned, err := preflight.LookupOwnedRelationNames(t.Context(), f.pool, f.schema, "t")
require.NoError(t, err)
assert.Equal(t, tt.wantOwned, owned)
})
}
}
// A run whose first-choice names all land is unaffected by the
// verification: the table owns exactly what the desired file claimed,
// including the shapes whose server-chosen names depend on more than the
// key columns — INCLUDE columns, repeated exclusion elements, and an
// identity column's stated sequence name.
func TestExecuteCreateVerifiesOwnedNamesWithoutFalsePositives(t *testing.T) {
tests := []struct {
name string
sql string
steps int
owned []string
}{
{
name: "serial primary key, unique column, and index",
sql: `
CREATE TABLE t (id serial PRIMARY KEY, v text UNIQUE);
CREATE INDEX t_v_idx ON t (v);`,
steps: 2,
owned: []string{"t_pkey", "t_v_key", "t_id_seq", "t_v_idx"},
},
{
name: "unique with INCLUDE",
sql: "CREATE TABLE t (id int, v text, UNIQUE (id) INCLUDE (v))",
steps: 1,
owned: []string{"t_id_v_key"},
},
{
name: "exclusion constraint with a repeated element",
sql: "CREATE TABLE t (a int, EXCLUDE USING btree (a WITH =, a WITH =))",
steps: 1,
owned: []string{"t_a_a1_excl"},
},
{
name: "exclusion constraint over two expressions",
sql: "CREATE TABLE t (a int, b int, EXCLUDE USING btree ((a + 1) WITH =, (b + 1) WITH =))",
steps: 1,
owned: []string{"t_expr_expr1_excl"},
},
{
name: "identity column with a stated sequence name",
sql: "CREATE TABLE t (id int GENERATED ALWAYS AS IDENTITY (SEQUENCE NAME t_custom_seq) PRIMARY KEY)",
steps: 1,
owned: []string{"t_pkey", "t_custom_seq"},
},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
f := newCreateFixture(t, "t")
rep, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, desired(t, tt.sql), createBudget, executor.DefaultRetryPolicy())
require.NoError(t, err)
require.Len(t, rep.Steps, tt.steps)
for _, name := range tt.owned {
assert.True(t, relationExists(t, f.pool, f.schema, name), "%s exists under its first-choice name", name)
}
})
}
}
// A table that owns more than the desired file claimed is not a mismatch:
// only a claim the table failed to honour is. Relations another actor
// attaches inside the CREATE TABLE's own transaction — a sequence owned by
// a column, a constraint added before the commit — are the table's to
// keep, and the run passes.
func TestExecuteCreateAcceptsOwnedNamesBeyondTheClaims(t *testing.T) {
f := newCreateFixture(t, "t")
testutil.RunDuringDDL(t, f.pool, testutil.DDLCommandEnd, "CREATE TABLE", f.schema, "t", fmt.Sprintf(`
CREATE SEQUENCE %[1]s.t_extra_seq OWNED BY %[1]s.t.id;
ALTER TABLE %[1]s.t ADD CONSTRAINT t_extra_key UNIQUE (v)`, f.schema))
rep, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr,
desired(t, "CREATE TABLE t (id int PRIMARY KEY, v text)"), createBudget, executor.DefaultRetryPolicy())
require.NoError(t, err)
require.Len(t, rep.Steps, 1)
owned, err := preflight.LookupOwnedRelationNames(t.Context(), f.pool, f.schema, "t")
require.NoError(t, err)
assert.Equal(t, preflight.OwnedRelationNames{
ConstraintIndexes: []string{"t_extra_key", "t_pkey"},
Sequences: []string{"t_extra_seq"},
}, owned, "the unclaimed relations the table gained stay owned by it")
}
// A CREATE TABLE that commits but whose table cannot then be read at its
// name — here renamed from under the executor inside the statement's own
// transaction — leaves the name set unproven. An unproven set is not a
// passing one: the run fails at step 1 under its own code with the read's
// cause in the chain, the index step after it never runs, and the table
// (under whatever name it now bears) is left for the operator.
func TestExecuteCreateReportsUnreadableOwnedNamesAsUnverified(t *testing.T) {
f := newCreateFixture(t, "t")
testutil.RunDuringDDL(t, f.pool, testutil.DDLCommandEnd, "CREATE TABLE", f.schema, "t",
fmt.Sprintf("ALTER TABLE %s.t RENAME TO t_moved", f.schema))
rep, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, desired(t, `
CREATE TABLE t (id int PRIMARY KEY, v text);
CREATE INDEX t_v_idx ON t (v);`), createBudget, executor.DefaultRetryPolicy())
require.Error(t, err)
var stepErr *executor.SequenceStepError
require.ErrorAs(t, err, &stepErr)
assert.Equal(t, 1, stepErr.Step)
assert.Equal(t, 2, stepErr.Total)
assert.Equal(t, executor.StepBrief, stepErr.Kind)
assert.Contains(t, stepErr.SQL, "CREATE TABLE")
assert.ErrorIs(t, err, executor.ErrCreateNamesUnverified)
assert.ErrorIs(t, err, preflight.ErrTableNotFound, "the read's own failure is the cause")
assert.NotErrorIs(t, err, executor.ErrCreateNameMismatch, "nothing was compared, so nothing mismatched")
assert.Equal(t, executor.CodeCreateNamesUnverified, executor.OutcomeCode(err),
"the code names the state the step left, not the read's fault")
assert.Empty(t, rep.Steps, "the failed step's own state is named by the code, not the committed prefix")
assert.Equal(t, "", relationKind(t, f.pool, f.schema, "t"), "nothing stands at the claimed name")
assert.Equal(t, "r", relationKind(t, f.pool, f.schema, "t_moved"), "the committed table stands under the name it was moved to")
assert.False(t, relationExists(t, f.pool, f.schema, "t_v_idx"), "the run stops at step 1; the index is never built on an unproven table")
}
// A CREATE TABLE that claims no server-chosen name — no index-backed
// constraint, no serial or identity column — has nothing to prove after
// it commits, so no read runs and nothing the read could hit can fail the
// run. Here the table is renamed inside its own transaction, which would
// leave a read with nothing at the claimed name; the run still passes.
// Without a create owner the proof names the session role, so no owner
// read runs either.
func TestExecuteCreateSkipsTheOwnedNameReadWithoutClaims(t *testing.T) {
f := newCreateFixture(t, "t")
testutil.RunDuringDDL(t, f.pool, testutil.DDLCommandEnd, "CREATE TABLE", f.schema, "t",
fmt.Sprintf("ALTER TABLE %s.t RENAME TO t_moved", f.schema))
rep, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr,
desired(t, "CREATE TABLE t (a int, v text)"), createBudget, executor.DefaultRetryPolicy())
require.NoError(t, err)
require.Len(t, rep.Steps, 1)
assert.Equal(t, "r", relationKind(t, f.pool, f.schema, "t_moved"))
}
// newCreateOwnerFixture mints the creation proof for a dedicated owner role
// that holds USAGE and CREATE on the schema. The pool stays the fixture's
// session role; the proof is what makes the create step SET LOCAL ROLE to
// the owner.
func newCreateOwnerFixture(t *testing.T, table string) (createFixture, string) {
t.Helper()
f := newCreateFixture(t, table)
owner := testutil.NewRole(t, f.pool, "NOLOGIN")
// The schema was created before the roles, so its drop would run after
// theirs; the roles will own and hold privileges on objects in it, so
// release the schema first.
t.Cleanup(func() {
_, err := f.pool.Exec(context.WithoutCancel(t.Context()),
fmt.Sprintf("DROP SCHEMA IF EXISTS %s CASCADE", pgx.Identifier{f.schema}.Sanitize()))
assert.NoError(t, err)
})
_, err := f.pool.Exec(t.Context(), fmt.Sprintf("GRANT USAGE, CREATE ON SCHEMA %s TO %s",
pgx.Identifier{f.schema}.Sanitize(), pgx.Identifier{owner}.Sanitize()))
require.NoError(t, err)
cr, err := preflight.CheckCreatePrivilegesAs(t.Context(), f.pool, f.schema, owner)
require.NoError(t, err)
require.True(t, cr.SetsRole())
f.cr = cr
return f, owner
}
// relationOwner returns the pg_roles name that owns schema.name — the
// catalog oracle for who a create run left holding the table.
func relationOwner(t *testing.T, pool *pgxpool.Pool, schema, name string) string {
t.Helper()
var owner string
require.NoError(t, pool.QueryRow(t.Context(),
`SELECT r.rolname
FROM pg_class c
JOIN pg_namespace n ON n.oid = c.relnamespace
JOIN pg_roles r ON r.oid = c.relowner
WHERE n.nspname = $1 AND c.relname = $2`, schema, name).Scan(&owner))
return owner
}
// With a create owner in the proof, the CREATE TABLE and its index run
// under SET LOCAL ROLE, so every relation the run creates — the table, its
// constraint index, its sequence, and the follow-on index — is owned by
// that role rather than the session role.
func TestExecuteCreateWithOwnerCreatesEveryRelationAsOwner(t *testing.T) {
f, owner := newCreateOwnerFixture(t, "t")
rep, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr, desired(t, `
CREATE TABLE t (
id serial PRIMARY KEY,
v text
);
CREATE INDEX t_v_idx ON t (v);
`), createBudget, executor.DefaultRetryPolicy())
require.NoError(t, err)
require.Len(t, rep.Steps, 2)
for _, relation := range []string{"t", "t_pkey", "t_id_seq", "t_v_idx"} {
assert.Equal(t, owner, relationOwner(t, f.pool, f.schema, relation), relation)
}
}
// When the proof names an owner but the committed table is not owned by
// it — here the table is transferred to a second role inside the CREATE's
// own transaction, which runs as the owner and so needs membership in
// that role — the run fails closed with the expected and actual owners,
// and the table is left as committed: the executor never repairs
// ownership with ALTER ... OWNER TO.
func TestExecuteCreateWithOwnerFailsClosedOnOwnerMismatch(t *testing.T) {
f, owner := newCreateOwnerFixture(t, "t")
other := testutil.NewRole(t, f.pool, "NOLOGIN")
// OWNER TO needs the transferring role to be a member of the new owner
// and the new owner to hold CREATE on the schema.
_, err := f.pool.Exec(t.Context(), fmt.Sprintf("GRANT %s TO %s; GRANT USAGE, CREATE ON SCHEMA %s TO %s",
pgx.Identifier{other}.Sanitize(), pgx.Identifier{owner}.Sanitize(),
pgx.Identifier{f.schema}.Sanitize(), pgx.Identifier{other}.Sanitize()))
require.NoError(t, err)
t.Cleanup(func() {
// The second role is created after the fixture, so its drop runs
// before the fixture's schema drop; release the schema first.
_, err := f.pool.Exec(context.WithoutCancel(t.Context()),
fmt.Sprintf("DROP SCHEMA IF EXISTS %s CASCADE", pgx.Identifier{f.schema}.Sanitize()))
assert.NoError(t, err)
})
testutil.RunDuringDDL(t, f.pool, testutil.DDLCommandEnd, "CREATE TABLE", f.schema, "t",
fmt.Sprintf("ALTER TABLE %s.t OWNER TO %s", pgx.Identifier{f.schema}.Sanitize(), pgx.Identifier{other}.Sanitize()))
rep, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr,
desired(t, "CREATE TABLE t (a int, v text)"), createBudget, executor.DefaultRetryPolicy())
var mismatch *executor.CreateOwnerMismatchError
require.ErrorAs(t, err, &mismatch)
assert.ErrorIs(t, err, executor.ErrCreateOwnerMismatch)
assert.Equal(t, executor.CodeCreateOwnerMismatch, executor.OutcomeCode(err),
"the step-1 failure carries the permanent ownership code, not the execution-failed fallback")
assert.Equal(t, owner, mismatch.Expected)
assert.Equal(t, other, mismatch.Actual)
assert.Empty(t, rep.Steps)
assert.Equal(t, other, relationOwner(t, f.pool, f.schema, "t"), "the committed owner is reported, not repaired")
}
// When the proof names an owner and the committed table cannot be read
// back — here it is renamed inside its own transaction — the owner is
// unproven and the run fails closed rather than passing on an unverified
// owner.
func TestExecuteCreateWithOwnerReportsUnreadableOwnerAsUnverified(t *testing.T) {
f, _ := newCreateOwnerFixture(t, "t")
testutil.RunDuringDDL(t, f.pool, testutil.DDLCommandEnd, "CREATE TABLE", f.schema, "t",
fmt.Sprintf("ALTER TABLE %s.t RENAME TO t_moved", pgx.Identifier{f.schema}.Sanitize()))
rep, err := executor.ExecuteCreate(t.Context(), f.pool, f.at, f.cr,
desired(t, "CREATE TABLE t (a int, v text)"), createBudget, executor.DefaultRetryPolicy())
require.ErrorIs(t, err, executor.ErrCreateOwnerUnverified)
assert.Equal(t, executor.CodeCreateOwnerUnverified, executor.OutcomeCode(err),
"an unreadable owner is retryable, so it must not carry the permanent mismatch code")
assert.Empty(t, rep.Steps)
assert.Equal(t, "r", relationKind(t, f.pool, f.schema, "t_moved"))
}