diff --git a/controlplane/trino_pool_operator.go b/controlplane/trino_pool_operator.go index 242526e0..a82774ca 100644 --- a/controlplane/trino_pool_operator.go +++ b/controlplane/trino_pool_operator.go @@ -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 @@ -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 } diff --git a/controlplane/trino_pool_progress.go b/controlplane/trino_pool_progress.go index ddfc6cec..5ce60703 100644 --- a/controlplane/trino_pool_progress.go +++ b/controlplane/trino_pool_progress.go @@ -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") } @@ -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. @@ -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. @@ -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, @@ -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, @@ -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. @@ -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, @@ -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 diff --git a/controlplane/trino_pool_superseded_test.go b/controlplane/trino_pool_superseded_test.go index 89e33d35..b4b78e64 100644 --- a/controlplane/trino_pool_superseded_test.go +++ b/controlplane/trino_pool_superseded_test.go @@ -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) diff --git a/docs/trino-cells.md b/docs/trino-cells.md index 50d2678b..5cf20465 100644 --- a/docs/trino-cells.md +++ b/docs/trino-cells.md @@ -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.