Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 8 additions & 0 deletions controlplane/trino_pool_operator.go
Original file line number Diff line number Diff line change
Expand Up @@ -108,6 +108,10 @@ type trinoPoolOperator struct {
nodeGuard *trinoPoolNodeGuard
nodeProtectionErrors map[string]error
nodeProtectionCursor uint64
// Registration can commit after its caller times out, even after a member 404.
// Only a confirmed higher Gateway fence clears unresolved attempts from a term.
registrationAttemptEpoch int64
registrationAttempts map[string]bool
// projection reports what this control plane currently serves: the
// authorization bundle's revision and the fingerprints of the projected
// password and group files. It is what a member is compared against when
Expand Down Expand Up @@ -508,6 +512,10 @@ func (o *trinoPoolOperator) configureGatewayPool(ctx context.Context) error {
if attempt.payload != desired {
return fmt.Errorf("%w: previous gateway configuration settled; current settings must be applied next", errTrinoPoolBackoff)
}
if o.registrationAttemptEpoch != o.lease.Epoch {
o.registrationAttempts = nil
o.registrationAttemptEpoch = o.lease.Epoch
}
return nil
}

Expand Down
50 changes: 41 additions & 9 deletions controlplane/trino_pool_progress.go
Original file line number Diff line number Diff line change
Expand Up @@ -82,17 +82,14 @@ func (o *trinoPoolOperator) progressInstance(ctx context.Context, instance confi
case trinopool.PhaseCreating:
return o.registerWhenReady(ctx, instance)
case trinopool.PhasePreparing:
if o.pool == nil || o.pool.DesiredReleaseID == "" || o.pool.DesiredBlueprintDigest == "" {
return false, errors.New("candidate supersession requires a complete desired pool specification")
}
blueprint, err := trinopool.ParseBlueprint([]byte(instance.BlueprintSnapshot))
superseded, err := o.candidateSuperseded(instance)
if err != nil {
return false, fmt.Errorf("read candidate blueprint for supersession: %w", err)
return false, err
}
// Configuration changes cannot repair an immutable candidate snapshot.
// Only PREPARING is safe here: VALIDATING can have an unknown admission outcome.
// VALIDATING can have an unknown admission outcome and must resolve it first.
// Cleanup still requires the Gateway's guarded retirement claim before deletion.
if instance.ReleaseID != o.pool.DesiredReleaseID || blueprint.Digest() != o.pool.DesiredBlueprintDigest {
if superseded {
return true, o.failCandidate(ctx, instance, trinopool.PhasePreparing,
"candidate blueprint was superseded before admission")
}
Expand Down Expand Up @@ -121,6 +118,17 @@ func (o *trinoPoolOperator) progressInstance(ctx context.Context, instance confi
}
}

func (o *trinoPoolOperator) candidateSuperseded(instance configstore.TrinoPoolInstance) (bool, error) {
if o.pool == nil || o.pool.DesiredReleaseID == "" || o.pool.DesiredBlueprintDigest == "" {
return false, errors.New("candidate supersession requires a complete desired pool specification")
}
blueprint, err := trinopool.ParseBlueprint([]byte(instance.BlueprintSnapshot))
if err != nil {
return false, fmt.Errorf("read candidate blueprint for supersession: %w", err)
}
return instance.ReleaseID != o.pool.DesiredReleaseID || blueprint.Digest() != o.pool.DesiredBlueprintDigest, nil
}

// createResources instantiates the instance's OWN blueprint snapshot, not the
// pool's current one: a release that landed after this instance was recorded
// must not change what it runs.
Expand Down Expand Up @@ -162,6 +170,18 @@ func (o *trinoPoolOperator) registerWhenReady(ctx context.Context, instance conf
if adopted, err := o.adoptRegisteredMember(ctx, instance, observed); adopted || err != nil {
return adopted, err
}
superseded, err := o.candidateSuperseded(instance)
if err != nil {
return false, err
}
if superseded {
// A same-term registration may still commit after this read-back returned 404.
if o.registrationAttempts[instance.InstanceID] {
return false, errors.New("superseded candidate registration remains unresolved in this authority term")
}
return true, o.failCandidate(ctx, instance, trinopool.PhaseCreating,
"candidate blueprint was superseded before registration")
}
if !observed.CoordinatorReady || observed.ReadyWorkers == 0 || observed.ReadyWorkers != observed.DesiredWorkers {
// Still converging. Not an error, and not something to time out into a
// failure: the plan's budgets already bound how many instances exist.
Expand Down Expand Up @@ -190,6 +210,10 @@ func (o *trinoPoolOperator) registerWhenReady(ctx context.Context, instance conf
return true, fmt.Errorf("register gateway backend for %s: %w", instance.InstanceID, err)
}

if o.registrationAttempts == nil {
o.registrationAttempts = make(map[string]bool)
}
o.registrationAttempts[instance.InstanceID] = true
member, err := o.gateway.RegisterMember(ctx, o.config.RoutingGroup, trinogateway.RegisterMemberRequest{
Step: o.step(instance.InstanceID, "register"),
InstanceID: instance.InstanceID,
Expand All @@ -203,7 +227,7 @@ func (o *trinoPoolOperator) registerWhenReady(ctx context.Context, instance conf
if err != nil {
return true, o.dropAuthority(fmt.Errorf("register member %s: %w", instance.InstanceID, err))
}
return true, o.dropAuthority(o.store.AdvanceTrinoPoolInstance(ctx, o.lease, instance.InstanceID,
err = o.dropAuthority(o.store.AdvanceTrinoPoolInstance(ctx, o.lease, instance.InstanceID,
trinopool.PhaseCreating, trinopool.PhasePreparing, map[string]any{
"coordinator_pod_uid": observed.CoordinatorPodUID,
"coordinator_boot_id": bootID,
Expand All @@ -224,6 +248,10 @@ func (o *trinoPoolOperator) registerWhenReady(ctx context.Context, instance conf
"gateway_state": member.Phase,
"gateway_generation": member.Generation,
}))
if err == nil {
delete(o.registrationAttempts, instance.InstanceID)
}
return true, err
}

// adoptRegisteredMember resolves a registration whose response was lost.
Expand Down Expand Up @@ -256,7 +284,7 @@ func (o *trinoPoolOperator) adoptRegisteredMember(
}
slog.Info("Trino pool adopted a member whose registration response was lost.",
"pool", o.config.PublicID, "instance", instance.InstanceID, "phase", member.Phase)
return true, o.dropAuthority(o.store.AdvanceTrinoPoolInstance(ctx, o.lease, instance.InstanceID,
err = o.dropAuthority(o.store.AdvanceTrinoPoolInstance(ctx, o.lease, instance.InstanceID,
trinopool.PhaseCreating, trinopool.PhasePreparing, map[string]any{
"coordinator_pod_uid": member.PodUID,
"coordinator_boot_id": member.BootID,
Expand All @@ -271,6 +299,10 @@ func (o *trinoPoolOperator) adoptRegisteredMember(
"gateway_state": member.Phase,
"gateway_generation": member.Generation,
}))
if err == nil {
delete(o.registrationAttempts, instance.InstanceID)
}
return true, err
}

// runningCoordinatorContainer is the container instance currently running in the
Expand Down
185 changes: 185 additions & 0 deletions controlplane/trino_pool_superseded_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,13 +6,198 @@ import (
"context"
"encoding/json"
"errors"
"fmt"
"strings"
"testing"

"github.com/posthog/duckgres/controlplane/configstore"
"github.com/posthog/duckgres/controlplane/trinogateway"
"github.com/posthog/duckgres/controlplane/trinopool"
)

func TestSupersededCreatingCandidateRecoversBeforeReadiness(t *testing.T) {
for _, test := range []struct {
name string
coordinator bool
workers int
newRelease bool
}{
{"unready coordinator", false, 4, true},
{"unready workers", true, 0, true},
{"changed configuration", false, 4, false},
} {
t.Run(test.name, func(t *testing.T) {
h := newOperatorHarness(t)
h.operator.config.Spec.DesiredInstances, h.operator.config.Spec.MinServing = 1, 1
h.tick(t, 2)
instance := h.store.instances[h.store.order[0]]
original := *instance
h.kube.observed.CoordinatorReady, h.kube.observed.ReadyWorkers = test.coordinator, test.workers
h.tick(t, 3)
if instance.Phase != string(trinopool.PhaseCreating) || len(h.store.order) != 1 {
t.Fatal("unchanged unready candidate did not retain its capacity")
}
current := changedCandidateConfig(t, h.operator.config, test.newRelease)
h.operator.resolveConfig = func() (trinoPoolConfig, error) { return current, nil }
h.tick(t, 1)
if instance.Phase != string(trinopool.PhaseFailedPreparing) {
t.Fatalf("superseded unready candidate phase = %s, want FAILED_PREPARING", instance.Phase)
}
if instance.BlueprintSnapshot != original.BlueprintSnapshot || instance.ReleaseID != original.ReleaseID || instance.SpecDigest != original.SpecDigest {
t.Fatal("supersession mutated the immutable candidate specification")
}
if h.kube.deleted[instance.InstanceID] || len(h.gateway.members) != 0 {
t.Fatal("supersession deleted resources immediately or registered an unready member")
}
h.tick(t, 1)
if !h.kube.deleted[instance.InstanceID] || !trinopool.Phase(instance.Phase).OccupiesCapacity() || len(h.store.order) != 1 {
t.Fatal("cleanup freed capacity before confirmed resource absence")
}
h.kube.absent = true
h.tick(t, 1)
if instance.Phase != string(trinopool.PhaseFailureRetired) {
t.Fatal("unregistered candidate did not finish cleanup after resource absence")
}
h.kube.observed.CoordinatorReady, h.kube.observed.ReadyWorkers = true, 4
h.tick(t, 6)
replacement := h.store.instances[h.store.order[1]]
if replacement.Phase != string(trinopool.PhaseServing) || replacement.ReleaseID != current.Spec.DesiredReleaseID {
t.Fatal("corrected candidate did not restore serving capacity")
}
if len(h.operator.registrationAttempts) != 0 {
t.Fatal("completed registration retained its term-local uncertainty")
}
})
}
}

func TestSupersededCreatingCandidateResolvesLostRegistrationBeforeCleanup(t *testing.T) {
h := newOperatorHarness(t)
h.operator.config.Spec.DesiredInstances, h.operator.config.Spec.MinServing = 1, 1
h.tick(t, 2)
instance := h.store.instances[h.store.order[0]]
h.gateway.loseResponse = map[string]bool{"register": true}
h.tickTolerant(1)
member := h.gateway.members[instance.InstanceID]
if instance.Phase != string(trinopool.PhaseCreating) || member == nil {
t.Fatal("fixture did not lose a committed registration response")
}
h.kube.observed.CoordinatorReady = false
h.operator.config = changedCandidateConfig(t, h.operator.config, true)
h.tick(t, 1)
if instance.Phase != string(trinopool.PhasePreparing) || instance.GatewayIncarnation != member.Incarnation || h.kube.deleted[instance.InstanceID] {
t.Fatal("supersession bypassed committed registration read-back")
}
h.tick(t, 2)
if member.Phase != "RETIRING" || h.kube.deleted[instance.InstanceID] {
t.Fatal("registered candidate cleanup bypassed Gateway retirement")
}
}

type unreadableCreatingMemberGateway struct{ *fakePoolGateway }

type delayedCreatingRegistrationGateway struct {
*fakePoolGateway
pending *trinogateway.RegisterMemberRequest
failConfigure bool
}

func (g *delayedCreatingRegistrationGateway) RegisterMember(_ context.Context, _ string, request trinogateway.RegisterMemberRequest) (trinogateway.Member, error) {
g.pending = &request
return trinogateway.Member{}, errors.New("registration response timed out before the Gateway committed")
}

func (g *delayedCreatingRegistrationGateway) ConfigurePool(ctx context.Context, pool string, request trinogateway.ConfigurePoolRequest) (trinogateway.PoolState, error) {
if g.failConfigure {
return trinogateway.PoolState{}, errors.New("Gateway fence unavailable")
}
return g.fakePoolGateway.ConfigurePool(ctx, pool, request)
}

func (g *delayedCreatingRegistrationGateway) complete(ctx context.Context, pool string) error {
if g.pending.ControllerEpoch != g.configured.ControllerEpoch {
return trinogateway.ErrStaleEpoch
}
_, err := g.fakePoolGateway.RegisterMember(ctx, pool, *g.pending)
return err
}

func TestSupersededCreatingCandidateWaitsForSameTermRegistration(t *testing.T) {
for _, takeover := range []bool{false, true} {
t.Run(fmt.Sprint("takeover=", takeover), func(t *testing.T) {
h := newOperatorHarness(t)
h.tick(t, 2)
instance := h.store.instances[h.store.order[0]]
gateway := &delayedCreatingRegistrationGateway{fakePoolGateway: h.gateway}
h.operator.gateway = gateway
h.tickTolerant(1)
if gateway.pending == nil || len(h.gateway.members) != 0 {
t.Fatal("fixture did not leave registration executing before its commit")
}
h.kube.observed.CoordinatorReady = false
h.operator.config = changedCandidateConfig(t, h.operator.config, true)
h.tickTolerant(2)
if instance.Phase != string(trinopool.PhaseCreating) || h.kube.deleted[instance.InstanceID] {
t.Fatal("Gateway 404 was mistaken for proof that a same-term registration cannot commit")
}
if takeover {
h.operator.lease = configstore.TrinoPoolLease{}
h.operator.owner = "replacement-controller"
gateway.failConfigure = true
h.tickTolerant(1)
if instance.Phase != string(trinopool.PhaseCreating) || h.kube.deleted[instance.InstanceID] {
t.Fatal("an unconfirmed new Gateway fence authorized cleanup")
}
if !h.operator.registrationAttempts[instance.InstanceID] {
t.Fatal("failed Gateway fencing cleared unresolved registration evidence")
}
gateway.failConfigure = false
h.tick(t, 1)
if instance.Phase != string(trinopool.PhaseFailedPreparing) {
t.Fatal("a confirmed higher fence did not release the obsolete registration uncertainty")
}
if len(h.operator.registrationAttempts) != 0 {
t.Fatal("confirmed new term retained its predecessor's registration attempts")
}
if err := gateway.complete(context.Background(), h.operator.config.RoutingGroup); !errors.Is(err, trinogateway.ErrStaleEpoch) {
t.Fatalf("old registration was not fenced: %v", err)
}
} else {
if err := gateway.complete(context.Background(), h.operator.config.RoutingGroup); err != nil {
t.Fatal(err)
}
h.tick(t, 1)
if instance.Phase != string(trinopool.PhasePreparing) || instance.GatewayIncarnation == "" || h.kube.deleted[instance.InstanceID] {
t.Fatal("late committed registration was not adopted before candidate cleanup")
}
if len(h.operator.registrationAttempts) != 0 {
t.Fatal("durable adoption retained registration uncertainty")
}
h.tick(t, 2)
if h.gateway.members[instance.InstanceID].Phase != "RETIRING" || h.kube.deleted[instance.InstanceID] {
t.Fatal("late registered candidate bypassed guarded retirement")
}
}
})
}
}

func (g *unreadableCreatingMemberGateway) GetMember(context.Context, string, string) (trinogateway.Member, error) {
return trinogateway.Member{}, errors.New("Gateway read-back unavailable")
}

func TestSupersededCreatingCandidatePreservesUnknownRegistration(t *testing.T) {
h := newOperatorHarness(t)
h.tick(t, 2)
instance := h.store.instances[h.store.order[0]]
h.operator.config = changedCandidateConfig(t, h.operator.config, true)
h.operator.gateway = &unreadableCreatingMemberGateway{h.gateway}
h.tickTolerant(2)
if instance.Phase != string(trinopool.PhaseCreating) || h.kube.deleted[instance.InstanceID] {
t.Fatal("an unknown registration outcome authorized candidate cleanup")
}
}

func supersededCandidateHarness(t *testing.T) (*operatorHarness, *configstore.TrinoPoolInstance) {
t.Helper()
h := newOperatorHarness(t)
Expand Down
22 changes: 22 additions & 0 deletions docs/trino-cells.md
Original file line number Diff line number Diff line change
Expand Up @@ -205,6 +205,28 @@ readiness reset and the ownership check.
Gateway public exposure is also a separate gate: authenticate every externally
reachable API and UI before publishing it; keep unauthenticated probes internal.

### Replace an unready, superseded candidate

An instance in `CREATING` can hold the only creation slot even when its
coordinator or workers never become ready. Publish a corrected release or
blueprint through the normal deployment path. The operator compares the
candidate's immutable specification with the current authoritative desired
specification before waiting for readiness. An unchanged slow candidate has
no automatic expiration deadline.

Before retiring an obsolete candidate, the operator resolves any Gateway
registration. A recorded member follows the existing guarded retirement path;
an unregistered candidate retains its capacity until its owned resources are
confirmed absent. This does not alter admission replay or serving-member drains.

A registration request can commit after its caller times out and after a member
lookup returns `404`. If this controller attempted registration in the current
authority term, supersession remains blocked until the member can be adopted or
a genuine leadership takeover establishes a higher Gateway fence. A request
that never commits can therefore require a control-plane restart or normal
leadership handoff; repeated `404` responses alone never authorize deletion.
Do not clear lifecycle records or edit authority epochs to bypass that wait.

### Query obligations that outlive their clients

A client can stop polling before the Gateway observes the query's terminal response.
Expand Down
Loading