Skip to content

Commit d3a4146

Browse files
authored
Add durable Trino cell catalog lifecycle primitives (#1186)
* Add durable Trino cell catalog lifecycle primitives * Require confirmed terminal catalog DDL outcomes in shared mode * Retain catalog intent after remote query failure * fix(trino): acknowledge finished result pages
1 parent ac00e76 commit d3a4146

9 files changed

Lines changed: 1138 additions & 6 deletions
Lines changed: 21 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,21 @@
1+
-- +goose Up
2+
-- Ownership has no timeout. Unknown remote mutation outcomes require recovery.
3+
CREATE TABLE duckgres_trino_cell_lifecycle (
4+
cell_id TEXT PRIMARY KEY,
5+
reconcile_owner TEXT NOT NULL DEFAULT '',
6+
reconcile_epoch BIGINT NOT NULL DEFAULT 0 CHECK (reconcile_epoch >= 0),
7+
intent_sequence BIGINT NOT NULL DEFAULT 0 CHECK (intent_sequence >= 0),
8+
intent JSONB NOT NULL DEFAULT '{}' CHECK (jsonb_typeof(intent) = 'object'),
9+
admission_epoch BIGINT NOT NULL DEFAULT 0 CHECK (admission_epoch >= 0),
10+
freeze_operation_id TEXT NOT NULL DEFAULT '',
11+
freeze_plan_hash TEXT NOT NULL DEFAULT '',
12+
freeze_target TEXT NOT NULL DEFAULT '',
13+
freeze_stable BOOLEAN NOT NULL DEFAULT FALSE,
14+
certificate JSONB NOT NULL DEFAULT '{}' CHECK (jsonb_typeof(certificate) = 'object'),
15+
released_operation_id TEXT NOT NULL DEFAULT '',
16+
released_admission_epoch BIGINT NOT NULL DEFAULT 0,
17+
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
18+
);
19+
20+
-- +goose Down
21+
DROP TABLE duckgres_trino_cell_lifecycle;
Lines changed: 334 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,334 @@
1+
package configstore
2+
3+
import (
4+
"context"
5+
"encoding/json"
6+
"errors"
7+
"regexp"
8+
"strings"
9+
"time"
10+
"unicode"
11+
12+
"gorm.io/gorm"
13+
"gorm.io/gorm/clause"
14+
)
15+
16+
var ErrTrinoCellConflict = errors.New("trino cell lifecycle conflict")
17+
var trinoCellHashPattern = regexp.MustCompile(`^[a-f0-9]{64}$`)
18+
19+
type TrinoCellLease struct {
20+
CellID string
21+
Owner string
22+
ReconcileEpoch int64
23+
AdmissionEpoch int64
24+
IntentSequence int64
25+
}
26+
27+
type TrinoCatalogIntent struct {
28+
Sequence int64 `json:"sequence"`
29+
ID string `json:"id"`
30+
Backend string `json:"backend"`
31+
Action string `json:"action"`
32+
Catalog string `json:"catalog"`
33+
}
34+
35+
type TrinoCellCertificate struct {
36+
TargetBackend string `json:"targetBackend"`
37+
NodeID string `json:"nodeId"`
38+
CoordinatorID string `json:"coordinatorId"`
39+
RosterHash string `json:"rosterHash"`
40+
AdmittedCount int `json:"admittedCount"`
41+
}
42+
43+
type TrinoCellFreeze struct {
44+
OperationID string `json:"operationId"`
45+
PlanHash string `json:"planHash"`
46+
TargetBackend string `json:"targetBackend"`
47+
AdmissionEpoch int64 `json:"admissionEpoch"`
48+
Certificate *TrinoCellCertificate `json:"certificate,omitempty"`
49+
Stable bool `json:"stable"`
50+
}
51+
52+
type trinoCellLifecycle struct {
53+
CellID string `gorm:"primaryKey"`
54+
ReconcileOwner string
55+
ReconcileEpoch int64
56+
IntentSequence int64
57+
Intent string `gorm:"type:jsonb"`
58+
AdmissionEpoch int64
59+
FreezeOperationID string
60+
FreezePlanHash string
61+
FreezeTarget string
62+
FreezeStable bool
63+
Certificate string `gorm:"type:jsonb"`
64+
ReleasedOperationID string
65+
ReleasedAdmissionEpoch int64
66+
UpdatedAt time.Time
67+
}
68+
69+
func (trinoCellLifecycle) TableName() string { return "duckgres_trino_cell_lifecycle" }
70+
71+
func validTrinoCellValue(s string) bool {
72+
return s != "" && len(s) <= 255 && strings.IndexFunc(s, unicode.IsControl) == -1
73+
}
74+
75+
func (cs *ConfigStore) withTrinoCell(ctx context.Context, cell string, fn func(*gorm.DB, *trinoCellLifecycle) error) error {
76+
if !strings.HasPrefix(cell, "registered:") || !validTrinoCellValue(cell) {
77+
return ErrTrinoCellConflict
78+
}
79+
return cs.db.WithContext(ctx).Transaction(func(tx *gorm.DB) error {
80+
seed := trinoCellLifecycle{CellID: cell, Intent: "{}", Certificate: "{}", UpdatedAt: time.Now().UTC()}
81+
if err := tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&seed).Error; err != nil {
82+
return err
83+
}
84+
var row trinoCellLifecycle
85+
if err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&row, "cell_id = ?", cell).Error; err != nil {
86+
return err
87+
}
88+
return fn(tx, &row)
89+
})
90+
}
91+
92+
func saveTrinoCell(tx *gorm.DB, row *trinoCellLifecycle) error {
93+
row.UpdatedAt = time.Now().UTC()
94+
return tx.Save(row).Error
95+
}
96+
97+
func ownsTrinoCell(row *trinoCellLifecycle, lease TrinoCellLease) bool {
98+
return row.CellID == lease.CellID && row.ReconcileOwner != "" && row.ReconcileOwner == lease.Owner && row.ReconcileEpoch == lease.ReconcileEpoch
99+
}
100+
101+
// BeginTrinoCellReconcile must precede the tenant and Gateway snapshots.
102+
// A held owner never expires, including after a controller process disappears.
103+
func (cs *ConfigStore) BeginTrinoCellReconcile(ctx context.Context, cell, owner string) (*TrinoCellLease, bool, error) {
104+
if !validTrinoCellValue(owner) {
105+
return nil, false, ErrTrinoCellConflict
106+
}
107+
var lease *TrinoCellLease
108+
err := cs.withTrinoCell(ctx, cell, func(tx *gorm.DB, row *trinoCellLifecycle) error {
109+
if row.ReconcileOwner != "" {
110+
return nil
111+
}
112+
row.ReconcileOwner = owner
113+
row.ReconcileEpoch++
114+
if err := saveTrinoCell(tx, row); err != nil {
115+
return err
116+
}
117+
lease = &TrinoCellLease{CellID: cell, Owner: owner, ReconcileEpoch: row.ReconcileEpoch, AdmissionEpoch: row.AdmissionEpoch, IntentSequence: row.IntentSequence}
118+
return nil
119+
})
120+
return lease, lease != nil && err == nil, err
121+
}
122+
123+
// SetTrinoCellIntent authorizes one submission only after its durable commit.
124+
// A duplicate call is a conflict even when its payload is identical.
125+
func (cs *ConfigStore) SetTrinoCellIntent(ctx context.Context, lease TrinoCellLease, intent TrinoCatalogIntent) error {
126+
if !validTrinoCellValue(intent.ID) || !validTrinoCellValue(intent.Backend) || !validTrinoCellValue(intent.Catalog) || (intent.Action != "create" && intent.Action != "drop") {
127+
return ErrTrinoCellConflict
128+
}
129+
encoded, err := json.Marshal(intent)
130+
if err != nil {
131+
return err
132+
}
133+
return cs.withTrinoCell(ctx, lease.CellID, func(tx *gorm.DB, row *trinoCellLifecycle) error {
134+
if !ownsTrinoCell(row, lease) || row.Intent != "{}" || row.FreezeOperationID != "" || row.AdmissionEpoch != lease.AdmissionEpoch || intent.Sequence <= 0 || intent.Sequence != row.IntentSequence+1 {
135+
return ErrTrinoCellConflict
136+
}
137+
row.Intent = string(encoded)
138+
row.IntentSequence = intent.Sequence
139+
return saveTrinoCell(tx, row)
140+
})
141+
}
142+
143+
// ClearTrinoCellIntent requires a confirmed terminal response for this intent.
144+
func (cs *ConfigStore) ClearTrinoCellIntent(ctx context.Context, lease TrinoCellLease, intentID string) error {
145+
return cs.withTrinoCell(ctx, lease.CellID, func(tx *gorm.DB, row *trinoCellLifecycle) error {
146+
var intent TrinoCatalogIntent
147+
if !ownsTrinoCell(row, lease) || json.Unmarshal([]byte(row.Intent), &intent) != nil || intent.ID == "" || intent.ID != intentID {
148+
return ErrTrinoCellConflict
149+
}
150+
row.Intent = "{}"
151+
return saveTrinoCell(tx, row)
152+
})
153+
}
154+
155+
func (cs *ConfigStore) FinishTrinoCellReconcile(ctx context.Context, lease TrinoCellLease) error {
156+
return cs.withTrinoCell(ctx, lease.CellID, func(tx *gorm.DB, row *trinoCellLifecycle) error {
157+
if !ownsTrinoCell(row, lease) || row.Intent != "{}" {
158+
return ErrTrinoCellConflict
159+
}
160+
row.ReconcileOwner = ""
161+
if row.FreezeOperationID != "" {
162+
row.FreezeStable = true
163+
}
164+
return saveTrinoCell(tx, row)
165+
})
166+
}
167+
168+
func freezeFromRow(row *trinoCellLifecycle) (*TrinoCellFreeze, error) {
169+
if row.FreezeOperationID == "" {
170+
return nil, nil
171+
}
172+
freeze := &TrinoCellFreeze{OperationID: row.FreezeOperationID, PlanHash: row.FreezePlanHash, TargetBackend: row.FreezeTarget, AdmissionEpoch: row.AdmissionEpoch, Stable: row.FreezeStable}
173+
if row.Certificate != "{}" {
174+
if err := json.Unmarshal([]byte(row.Certificate), &freeze.Certificate); err != nil {
175+
return nil, err
176+
}
177+
}
178+
return freeze, nil
179+
}
180+
181+
func (cs *ConfigStore) FreezeTrinoCellAdmissions(ctx context.Context, cell, operation, planHash, target string, expectedEpoch int64) (*TrinoCellFreeze, error) {
182+
if !validTrinoCellValue(operation) || !validTrinoCellValue(target) || !trinoCellHashPattern.MatchString(planHash) {
183+
return nil, ErrTrinoCellConflict
184+
}
185+
var result *TrinoCellFreeze
186+
err := cs.withTrinoCell(ctx, cell, func(tx *gorm.DB, row *trinoCellLifecycle) error {
187+
if row.ReleasedOperationID == operation {
188+
return ErrTrinoCellConflict
189+
}
190+
if row.FreezeOperationID != "" {
191+
if row.FreezeOperationID != operation || row.FreezePlanHash != planHash || row.FreezeTarget != target {
192+
return ErrTrinoCellConflict
193+
}
194+
} else {
195+
if row.AdmissionEpoch != expectedEpoch {
196+
return ErrTrinoCellConflict
197+
}
198+
row.AdmissionEpoch++
199+
row.FreezeOperationID, row.FreezePlanHash, row.FreezeTarget = operation, planHash, target
200+
row.FreezeStable = row.ReconcileOwner == ""
201+
row.Certificate = "{}"
202+
if err := saveTrinoCell(tx, row); err != nil {
203+
return err
204+
}
205+
}
206+
var err error
207+
result, err = freezeFromRow(row)
208+
return err
209+
})
210+
return result, err
211+
}
212+
213+
func (cs *ConfigStore) GetTrinoCellFreeze(ctx context.Context, cell string) (*TrinoCellFreeze, error) {
214+
var row trinoCellLifecycle
215+
err := cs.db.WithContext(ctx).First(&row, "cell_id = ?", cell).Error
216+
if errors.Is(err, gorm.ErrRecordNotFound) {
217+
return nil, nil
218+
}
219+
if err != nil {
220+
return nil, err
221+
}
222+
return freezeFromRow(&row)
223+
}
224+
225+
type TrinoCellLifecycleStatus struct {
226+
CellID string
227+
AdmissionEpoch int64
228+
ReconcileOwner string
229+
ReconcileEpoch int64
230+
IntentSequence int64
231+
Intent *TrinoCatalogIntent
232+
Freeze *TrinoCellFreeze
233+
ReleasedOperationID string
234+
ReleasedAdmissionEpoch int64
235+
}
236+
237+
// GetTrinoCellLifecycle returns the current epoch even when no freeze is active.
238+
// Reading an uninitialized managed cell does not create a lifecycle row.
239+
func (cs *ConfigStore) GetTrinoCellLifecycle(ctx context.Context, cell string) (*TrinoCellLifecycleStatus, error) {
240+
if !strings.HasPrefix(cell, "registered:") || !validTrinoCellValue(cell) {
241+
return nil, ErrTrinoCellConflict
242+
}
243+
var row trinoCellLifecycle
244+
err := cs.db.WithContext(ctx).First(&row, "cell_id = ?", cell).Error
245+
if errors.Is(err, gorm.ErrRecordNotFound) {
246+
return &TrinoCellLifecycleStatus{CellID: cell}, nil
247+
}
248+
if err != nil {
249+
return nil, err
250+
}
251+
freeze, err := freezeFromRow(&row)
252+
if err != nil {
253+
return nil, err
254+
}
255+
result := &TrinoCellLifecycleStatus{CellID: cell, AdmissionEpoch: row.AdmissionEpoch, ReconcileOwner: row.ReconcileOwner, ReconcileEpoch: row.ReconcileEpoch, IntentSequence: row.IntentSequence, Freeze: freeze, ReleasedOperationID: row.ReleasedOperationID, ReleasedAdmissionEpoch: row.ReleasedAdmissionEpoch}
256+
if row.Intent != "{}" {
257+
if err := json.Unmarshal([]byte(row.Intent), &result.Intent); err != nil {
258+
return nil, err
259+
}
260+
}
261+
return result, nil
262+
}
263+
264+
func (cs *ConfigStore) CertifyTrinoCellTarget(ctx context.Context, lease TrinoCellLease, operation string, certificate TrinoCellCertificate) error {
265+
if !validTrinoCellValue(certificate.TargetBackend) || !validTrinoCellValue(certificate.NodeID) || !validTrinoCellValue(certificate.CoordinatorID) || !trinoCellHashPattern.MatchString(certificate.RosterHash) || certificate.AdmittedCount < 0 || certificate.AdmittedCount > 100000 {
266+
return ErrTrinoCellConflict
267+
}
268+
encoded, err := json.Marshal(certificate)
269+
if err != nil {
270+
return err
271+
}
272+
return cs.withTrinoCell(ctx, lease.CellID, func(tx *gorm.DB, row *trinoCellLifecycle) error {
273+
if !ownsTrinoCell(row, lease) || row.Intent != "{}" || !row.FreezeStable || row.FreezeOperationID != operation || row.FreezeTarget != certificate.TargetBackend || row.AdmissionEpoch != lease.AdmissionEpoch {
274+
return ErrTrinoCellConflict
275+
}
276+
if row.Certificate != "{}" {
277+
var existing TrinoCellCertificate
278+
if json.Unmarshal([]byte(row.Certificate), &existing) != nil || existing != certificate {
279+
return ErrTrinoCellConflict
280+
}
281+
return nil
282+
}
283+
row.Certificate = string(encoded)
284+
return saveTrinoCell(tx, row)
285+
})
286+
}
287+
288+
// ReleaseTrinoCellAdmissions follows verification of the exact Gateway target route.
289+
// The caller must validate the active operation and certified coordinator process.
290+
func (cs *ConfigStore) ReleaseTrinoCellAdmissions(ctx context.Context, cell, operation string, epoch int64) error {
291+
return cs.withTrinoCell(ctx, cell, func(tx *gorm.DB, row *trinoCellLifecycle) error {
292+
if row.FreezeOperationID == "" && row.ReleasedOperationID == operation && row.ReleasedAdmissionEpoch == epoch {
293+
return nil
294+
}
295+
if operation == "" || row.FreezeOperationID != operation || row.AdmissionEpoch != epoch || row.Certificate == "{}" {
296+
return ErrTrinoCellConflict
297+
}
298+
row.ReleasedOperationID, row.ReleasedAdmissionEpoch = operation, epoch
299+
row.FreezeOperationID, row.FreezePlanHash, row.FreezeTarget = "", "", ""
300+
row.FreezeStable = false
301+
row.Certificate = "{}"
302+
row.AdmissionEpoch++
303+
return saveTrinoCell(tx, row)
304+
})
305+
}
306+
307+
// UpdateManagedTrinoState fences new admission in the same database transaction.
308+
// A backend health observation cannot implicitly grant new admission after cutover.
309+
func (cs *ConfigStore) UpdateManagedTrinoState(ctx context.Context, lease TrinoCellLease, org string, update TrinoStateUpdate) (bool, error) {
310+
updated := false
311+
err := cs.withTrinoCell(ctx, lease.CellID, func(tx *gorm.DB, row *trinoCellLifecycle) error {
312+
if !ownsTrinoCell(row, lease) {
313+
return ErrTrinoCellConflict
314+
}
315+
var tenant ManagedWarehouseTrino
316+
err := tx.Clauses(clause.Locking{Strength: "UPDATE"}).First(&tenant, "org_id = ? AND enabled = ? AND trino_cell_id = ?", org, true, lease.CellID).Error
317+
if errors.Is(err, gorm.ErrRecordNotFound) {
318+
return nil
319+
}
320+
if err != nil {
321+
return err
322+
}
323+
if update.State == ManagedWarehouseStateReady && tenant.State != ManagedWarehouseStateReady && (row.FreezeOperationID != "" || row.AdmissionEpoch != lease.AdmissionEpoch) {
324+
return nil
325+
}
326+
temporary := &ConfigStore{db: tx}
327+
if err := temporary.UpdateTrinoState(org, update); err != nil {
328+
return err
329+
}
330+
updated = true
331+
return nil
332+
})
333+
return updated, err
334+
}

‎controlplane/provisioner/opa/policy.rego‎

Lines changed: 13 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -736,6 +736,19 @@ allow if {
736736
observer_nodes_table(input.action.resource.table)
737737
}
738738

739+
# The provisioner reads catalog startup status without querying tenant data.
740+
allow if {
741+
is_admin
742+
input.action.operation == "SelectFromColumns"
743+
table := input.action.resource.table
744+
table.catalogName == "system"
745+
table.schemaName == "metadata"
746+
table.tableName == "catalogs"
747+
every column in table.columns {
748+
column in {"catalog_name", "state"}
749+
}
750+
}
751+
739752
# ---------------------------------------------------------------------------
740753
# Hard denies for customer principals.
741754
#

0 commit comments

Comments
 (0)