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
9 changes: 8 additions & 1 deletion e2e/k8s/k8s_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -864,7 +864,14 @@ func hasRowCopyProgress(rowsTotal, rowsCopied int64, percentComplete int32) bool

func storedK8sApplyAndTaskStates(t *testing.T, dsn, applyID string) (string, string) {
t.Helper()
db := testutil.OpenMySQL(t, dsn)
return storedApplyAndTaskStates(t, testutil.OpenMySQL(t, dsn), applyID)
}

// storedApplyAndTaskStates is the handle-taking form of
// storedK8sApplyAndTaskStates for callers that already hold the storage pool,
// such as poll loops that would otherwise open one per tick.
func storedApplyAndTaskStates(t *testing.T, db *sql.DB, applyID string) (string, string) {
t.Helper()

var applyState, taskState string
require.NoError(t, db.QueryRowContext(t.Context(), `
Expand Down
56 changes: 43 additions & 13 deletions e2e/k8s/retryable_pause_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -21,14 +21,22 @@ import (
// database. The data plane parks the apply between its own recovery attempts,
// and that pause must be survivable end to end:
//
// - The pause crosses the wire as STATE_FAILED_RETRYABLE, so the control
// plane can tell "paused, will self-retry" from a settled failure without
// inspecting per-table statuses.
// - The pause never settles as a failure: the wire either renders
// STATE_FAILED_RETRYABLE or has already moved on to the recovery attempt,
// and STATE_FAILED anywhere in between fails the test.
// - The control plane's stored apply stays non-terminal for the whole pause.
// A terminal verdict here would end the drive and orphan a live remote
// apply that goes on to change the schema with nobody watching.
// - Once the failure injection stops, the data plane's own recovery claims
// another attempt and both planes land completed with the DDL applied.
//
// The pause is transient by design: the data plane's recovery reclaims a
// failed_retryable apply as soon as a driver's claim poll sees it, so the wire
// may render STATE_FAILED_RETRYABLE for less than one observation interval.
// The test therefore accepts either witness of the pause — the wire state, or
// the data plane's stored attempt counter advancing past its pre-kill value,
// which only a claim out of failed_retryable does — so a pause that recovery
// has already consumed still counts as observed.
func TestK8s_DataPlaneRetryablePauseHoldsControlPlaneOpenUntilRecovery(t *testing.T) {
cleanupState(t)

Expand All @@ -43,6 +51,7 @@ func TestK8s_DataPlaneRetryablePauseHoldsControlPlaneOpenUntilRecovery(t *testin
podClient := dialDataPlanePod(t, pods[0])
killer := testutil.OpenMySQL(t, testutil.TernStagingDSN(t))
controlPlaneDB := testutil.OpenMySQL(t, testutil.SchemabotDSN(t))
dataPlaneDB := testutil.OpenMySQL(t, storageDSNs(t)[0])

// Return from the fixture at dispatch rather than waiting for the control
// plane to report running: the control plane's view lags the engine by a
Expand All @@ -55,20 +64,24 @@ func TestK8s_DataPlaneRetryablePauseHoldsControlPlaneOpenUntilRecovery(t *testin
// completes with no pause to observe.
fixture := startIndexAddApplyWithOptions(t, "k8s_retry_pause", false, nil, 2000000)
waitForPodApplyState(t, podClient, fixture.DataPlaneApplyID, ternv1.State_STATE_RUNNING, testutil.PollDeadline)
attemptBeforeKills := storedDataPlaneApplyAttempt(t, dataPlaneDB, fixture.DataPlaneApplyID)

// Kill the data plane's target connections on every poll tick until the
// pause is visible on the wire. Spirit may absorb a single kill mid-chunk,
// so the injection repeats until the pause is actually observed. Each tick
// observes the wire before it injects: once the pause is visible, the data
// pause is observed. Spirit may absorb a single kill mid-chunk, so the
// injection repeats until the pause is actually observed. Each tick
// observes before it injects: once the pause has happened, the data
// plane's own recovery attempt may already be re-driving the apply, and a
// kill landing on that attempt would sabotage the very recovery the rest
// of the test waits for. The wire trails the engine — the pause is stamped
// to storage before it renders on the wire — so a tick can still inject
// just after the engine run has already failed; observing first narrows
// that window rather than closing it. Throughout, the control plane's
// stored apply must stay non-terminal: the data plane will retry, so any
// terminal state here is the split-brain this stack prevents.
// that window rather than closing it. The pause counts as observed when
// the wire shows it or when the stored attempt counter proves recovery
// has already claimed out of it. Throughout, the control plane's stored
// apply must stay non-terminal: the data plane will retry, so any terminal
// state here is the split-brain this stack prevents.
var lastWireState ternv1.State
var lastAttempt int
kills := 0
testutil.Poll(t, 3*time.Minute, 250*time.Millisecond,
func() bool {
Expand All @@ -82,18 +95,23 @@ func TestK8s_DataPlaneRetryablePauseHoldsControlPlaneOpenUntilRecovery(t *testin
if lastWireState == ternv1.State_STATE_FAILED_RETRYABLE {
return true
}
lastAttempt = storedDataPlaneApplyAttempt(t, dataPlaneDB, fixture.DataPlaneApplyID)
if lastAttempt > attemptBeforeKills {
return true
}

killDataPlaneTargetConnections(t, killer)
kills++
return false
},
func() string {
return fmt.Sprintf("timeout waiting for the data-plane pause to cross the wire as STATE_FAILED_RETRYABLE, last wire state: %s", lastWireState)
return fmt.Sprintf("timeout waiting for the data-plane pause: wire never showed STATE_FAILED_RETRYABLE (last %s) and the stored attempt never advanced past %d (last %d)",
lastWireState, attemptBeforeKills, lastAttempt)
})
require.Positive(t, kills,
"the pause must come from injected connection kills, not a failure the apply produced on its own")

// The pause is on the wire and the control plane is still holding. Stop
// The pause has been observed and the control plane is still holding. Stop
// injecting failures: the data plane's next recovery attempt must finish
// the schema change and reconcile both planes to completed.
testutil.WaitForState(t, fixture.Endpoint, fixture.ApplyID, state.Apply.Completed, 3*time.Minute)
Expand All @@ -105,7 +123,7 @@ func TestK8s_DataPlaneRetryablePauseHoldsControlPlaneOpenUntilRecovery(t *testin
var dataPlaneApplyState, dataPlaneTaskState string
testutil.Poll(t, testutil.PollDeadline, testutil.PollInterval,
func() bool {
dataPlaneApplyState, dataPlaneTaskState = storedK8sApplyAndTaskStates(t, storageDSNs(t)[0], fixture.DataPlaneApplyID)
dataPlaneApplyState, dataPlaneTaskState = storedApplyAndTaskStates(t, dataPlaneDB, fixture.DataPlaneApplyID)
return state.IsState(dataPlaneApplyState, state.Apply.Completed) &&
state.IsState(dataPlaneTaskState, state.Task.Completed)
},
Expand All @@ -121,10 +139,22 @@ func storedControlPlaneApplyState(t *testing.T, db *sql.DB, applyID string) stri
t.Helper()
var applyState string
require.NoError(t, db.QueryRowContext(t.Context(),
"SELECT state FROM applies WHERE apply_identifier = ?", applyID).Scan(&applyState))
"SELECT `state` FROM `applies` WHERE `apply_identifier` = ?", applyID).Scan(&applyState))
return applyState
}

// storedDataPlaneApplyAttempt reads the data plane's stored attempt counter for
// its apply. The counter only advances out of failed_retryable, so an increase
// is durable proof that a retryable pause happened and recovery has already
// picked it up, even after the pause itself has left the wire.
func storedDataPlaneApplyAttempt(t *testing.T, db *sql.DB, applyID string) int {
t.Helper()
var attempt int
require.NoError(t, db.QueryRowContext(t.Context(),
"SELECT `attempt` FROM `applies` WHERE `apply_identifier` = ?", applyID).Scan(&attempt))
return attempt
}

// killDataPlaneTargetConnections kills every connection the data plane holds
// to the target database, failing whatever engine work those connections were
// doing. The killer's own connection is excluded.
Expand Down
Loading