From 935e87107b41d4f3cb380c8c625ea58dedd6fb03 Mon Sep 17 00:00:00 2001 From: sergeyb Date: Fri, 4 Sep 2026 22:05:31 +0000 Subject: [PATCH] fix(runway): reconcile stale merge retries MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Summary: Intent: - Let at-least-once merge deliveries converge when a prior attempt already updated Git but failed to publish its result. - Preserve staleness protection for superseded requests that would still modify the target. Changes: - Apply transforming merge strategies locally before staleness validation and return success when the target already satisfies every step. - Check PROMOTE containment before freshness and account for head refs moved by a prior contention attempt. - Add real-Git fail-once end-to-end coverage, strategy-level regression tests, and adjacent reproduction documentation. --- Generated by the 🪄 pr-create skill in devexp-agent-marketplace --- runway/extension/merger/git/git_merger.go | 116 ++++++++------- .../extension/merger/git/git_merger_test.go | 82 ++++++++++- runway/extension/merger/git/objects.go | 6 +- test/e2e/submitqueue/ISS-001.md | 58 ++++++++ test/e2e/submitqueue/git_suite_test.go | 132 ++++++++++++++++-- 5 files changed, 329 insertions(+), 65 deletions(-) create mode 100644 test/e2e/submitqueue/ISS-001.md diff --git a/runway/extension/merger/git/git_merger.go b/runway/extension/merger/git/git_merger.go index f0d6575f6..f9b08487d 100644 --- a/runway/extension/merger/git/git_merger.go +++ b/runway/extension/merger/git/git_merger.go @@ -420,18 +420,13 @@ func (o *providerCheck) check(ref changeRef, stepID string) error { // strategies (REBASE, SQUASH_REBASE, MERGE), retrying on remote contention when // committing. For a dry run it applies the steps locally then discards them. func (m *gitMerger) applyTransforming(ctx context.Context, req *runwaymq.MergeRequest, steps []resolvedStep, commit bool) (*runwaymq.MergeResult, error) { - // Fetch and vet every change the request names before the first attempt, so - // an unusable request fails without having touched the checkout. This sits - // outside the retry loop deliberately: the refs are fixed for the whole - // request, and re-checking them per attempt would re-query the remote to - // learn what it already told us. + // Fetch every change the request names before the first attempt. Freshness + // is checked after local application, when the merger knows whether the + // target already satisfies the request. refs := stepChangeRefs(steps) if err := m.ensureObjects(ctx, refs); err != nil { return nil, err } - if err := m.checkStale(ctx, refs); err != nil { - return nil, err - } var lastErr error // Head branches an attempt moved before failing to push the target stay @@ -439,8 +434,27 @@ func (m *gitMerger) applyTransforming(ctx context.Context, req *runwaymq.MergeRe // the branch is moved on rather than stranded on a commit that never landed. tracked := make(headBranchTracker) for attempt := 1; attempt <= m.maxPushAttempts; attempt++ { - baseSHA, stepResults, err := m.tryApply(ctx, steps, commit, tracked) - if err == nil { + baseSHA, stepResults, heads, err := m.tryApply(ctx, steps) + if err != nil { + if staleErr := m.checkStale(ctx, refs, tracked); staleErr != nil { + return nil, staleErr + } + + // A conflict is terminal — no retry. The failing apply function + // already discarded its in-progress Git operation. + if errors.Is(err, merger.ErrConflict) { + return nil, err + } + if !commit || baseSHA == "" { + return nil, err + } + } else if !stepResultsHaveOutputs(stepResults) { + m.logger.Debugw("merge complete", "id", req.GetId(), "target", m.target, "commit", commit) + return successResult(req, stepResults), nil + } else { + if staleErr := m.checkStale(ctx, refs, tracked); staleErr != nil { + return nil, staleErr + } if !commit { // Discard the local commits the dry run created so the checkout // is clean for the next operation, and report empty Outputs. @@ -448,25 +462,28 @@ func (m *gitMerger) applyTransforming(ctx context.Context, req *runwaymq.MergeRe return nil, fmt.Errorf("discard after dry-run: %w", derr) } stripOutputs(stepResults) + m.logger.Debugw("merge complete", "id", req.GetId(), "target", m.target, "commit", commit) + return successResult(req, stepResults), nil } - m.logger.Debugw("merge complete", "id", req.GetId(), "target", m.target, "commit", commit) - return successResult(req, stepResults), nil - } - // A conflict is terminal — no retry. Discard any partial dry-run state. - if errors.Is(err, merger.ErrConflict) { - if !commit { - _ = m.resetToRemote(ctx) + if m.updateHeadBranch { + err = m.updateHeadBranches(ctx, heads, tracked) + } + if err == nil { + if pushErr := m.push(ctx); pushErr != nil { + coremetrics.NamedCounter(m.metricsScope, "merge", "git_push_errors", 1) + err = pushErr + } + } + if err == nil { + m.logger.Debugw("merge complete", "id", req.GetId(), "target", m.target, "commit", commit) + return successResult(req, stepResults), nil } - return nil, err } // Only a push failure caused by the remote tip moving under us (between // reset and push) is worth retrying; everything else is fatal. baseSHA // is empty when the failure happened before reset captured a base. - if !commit || baseSHA == "" { - return nil, err - } currentSHA, refetchErr := m.refetchTipSHA(ctx) if refetchErr != nil { return nil, fmt.Errorf("refetch after push failure failed: %v (original push error: %w)", refetchErr, err) @@ -490,43 +507,26 @@ func (m *gitMerger) applyTransforming(ctx context.Context, req *runwaymq.MergeRe return nil, fmt.Errorf("exceeded %d merge attempts due to remote contention: %w", m.maxPushAttempts, lastErr) } -// tryApply runs one full reset+apply(+push) cycle. The returned baseSHA is the -// SHA the cycle was based on (set as soon as resetToRemote completes) so the -// caller can distinguish concurrent-push contention from other failures. The -// tracker carries head-branch state across attempts; see headBranchTracker. -func (m *gitMerger) tryApply(ctx context.Context, steps []resolvedStep, commit bool, tracked headBranchTracker) (string, []*runwaymq.StepResult, error) { +// tryApply resets to the current target and applies every step locally. Remote +// writes remain with the caller so freshness can be checked after application +// but before any ref is updated. +func (m *gitMerger) tryApply(ctx context.Context, steps []resolvedStep) (string, []*runwaymq.StepResult, []headUpdate, error) { if err := m.resetToRemote(ctx); err != nil { coremetrics.NamedCounter(m.metricsScope, "merge", "reset_errors", 1) - return "", nil, err + return "", nil, nil, err } baseSHA, err := m.headSHA(ctx) if err != nil { - return "", nil, err + return "", nil, nil, err } stepResults, heads, err := m.applySteps(ctx, steps) if err != nil { // The failing apply function aborts its own in-progress git operation; // the next attempt starts with resetToRemote regardless. - return baseSHA, nil, err - } - - if commit { - // The head branches move first, as their own push. A provider decides - // merged-versus-closed while processing the push to the target, against - // the head it has recorded at that moment, so a head that moves later — - // or in the same atomic push — is recorded too late. See headbranch.go. - if m.updateHeadBranch { - if err := m.updateHeadBranches(ctx, heads, tracked); err != nil { - return baseSHA, nil, err - } - } - if err := m.push(ctx); err != nil { - coremetrics.NamedCounter(m.metricsScope, "merge", "git_push_errors", 1) - return baseSHA, nil, err - } + return baseSHA, nil, nil, err } - return baseSHA, stepResults, nil + return baseSHA, stepResults, heads, nil } // applied is what one step produced: the commits created on the target, and the @@ -714,17 +714,12 @@ func (m *gitMerger) promote(ctx context.Context, req *runwaymq.MergeRequest, rs sha := ref.SHA // PROMOTE does not go through tryApply, so it performs the same availability - // and freshness checks itself. Without them a commit the remote cannot - // supply turns every containment query into a plain error, which the - // consumer retries forever instead of reporting. They sit outside the retry - // loop because the commit under promotion is fixed for the whole request — - // only the target tip moves between attempts. + // check itself. Without it a commit the remote cannot supply turns every + // containment query into a plain error, which the consumer retries forever + // instead of reporting. if err := m.ensureObjects(ctx, []changeRef{ref}); err != nil { return nil, err } - if err := m.checkStale(ctx, []changeRef{ref}); err != nil { - return nil, err - } var lastErr error for attempt := 1; attempt <= m.maxPushAttempts; attempt++ { @@ -748,6 +743,10 @@ func (m *gitMerger) promote(ctx context.Context, req *runwaymq.MergeRequest, rs return promoteResult(req, rs, sha, commit), nil } + if err := m.checkStale(ctx, []changeRef{ref}, nil); err != nil { + return nil, err + } + // Only a true fast-forward is allowed; divergence is a terminal conflict. fastForward, err := m.isAncestor(ctx, tip, sha) if err != nil { @@ -1167,6 +1166,15 @@ func stripOutputs(steps []*runwaymq.StepResult) { } } +func stepResultsHaveOutputs(steps []*runwaymq.StepResult) bool { + for _, step := range steps { + if len(step.GetOutputs()) > 0 { + return true + } + } + return false +} + // successResult builds a SUCCEEDED MergeResult echoing the request id. func successResult(req *runwaymq.MergeRequest, steps []*runwaymq.StepResult) *runwaymq.MergeResult { return &runwaymq.MergeResult{ diff --git a/runway/extension/merger/git/git_merger_test.go b/runway/extension/merger/git/git_merger_test.go index b1e6e2224..81ffc3c4f 100644 --- a/runway/extension/merger/git/git_merger_test.go +++ b/runway/extension/merger/git/git_merger_test.go @@ -19,6 +19,7 @@ import ( "context" "errors" "fmt" + "net/url" "os" "os/exec" "path/filepath" @@ -374,8 +375,13 @@ func TestMerge_Rebase_RetriesWhenRemoteMovesUnderUs(t *testing.T) { featureSHA := f.pushPRCommit(t, "feature/a", "hello.txt", "hello\nearth\n", "tweak hello") f.installRaceHook(t, []string{raceSHA}) - m := f.newMerger(t, mergestrategypb.Strategy_REBASE) - res, err := m.Merge(context.Background(), req("b", stepOf(mergestrategypb.Strategy_REBASE, "s1", uri(featureSHA)))) + m := f.newMergerWith(t, func(p *Params) { + p.CheckStaleness = true + p.UpdateHeadBranch = true + }) + res, err := m.Merge(context.Background(), req("b", + stepOf(mergestrategypb.Strategy_REBASE, "s1", gitURI("feature/a", featureSHA)), + )) require.NoError(t, err) require.Len(t, res.GetSteps(), 1) require.Len(t, res.GetSteps()[0].GetOutputs(), 1) @@ -388,6 +394,7 @@ func TestMerge_Rebase_RetriesWhenRemoteMovesUnderUs(t *testing.T) { assert.Equal(t, raceSHA, commits[0], "race commit landed first via the hook") assert.Equal(t, res.GetSteps()[0].GetOutputs()[0].GetId(), commits[1], "our cherry-pick landed on top after the retry") + assert.Equal(t, res.GetSteps()[0].GetOutputs()[0].GetId(), f.remoteSHA(t, "feature/a")) assert.Equal(t, "hello\nearth\n", f.remoteFile(t, "hello.txt")) } @@ -1189,6 +1196,73 @@ func TestMerge_StaleChangeRejected(t *testing.T) { require.NoError(t, err) } +func TestMerge_StaleAlreadySatisfiedSucceeds(t *testing.T) { + tests := []struct { + name string + strategy mergestrategypb.Strategy + }{ + {name: "rebase", strategy: mergestrategypb.Strategy_REBASE}, + {name: "squash rebase", strategy: mergestrategypb.Strategy_SQUASH_REBASE}, + {name: "merge", strategy: mergestrategypb.Strategy_MERGE}, + {name: "promote", strategy: mergestrategypb.Strategy_PROMOTE}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + f := setupGitFixture(t) + stale := f.pushPRCommit(t, "feature/satisfied", "satisfied.txt", "satisfied\n", "satisfied") + switch tt.strategy { + case mergestrategypb.Strategy_REBASE, mergestrategypb.Strategy_SQUASH_REBASE: + f.landOnMain(t, stale) + case mergestrategypb.Strategy_MERGE, mergestrategypb.Strategy_PROMOTE: + f.advanceMain(t, stale) + } + mainBefore := f.remoteHEAD(t) + current := f.pushPRCommitFrom(t, stale, "feature/satisfied", "current.txt", "current\n", "current") + require.NotEqual(t, stale, current) + + m := f.newMergerWith(t, func(p *Params) { p.CheckStaleness = true }) + res, err := m.Merge(context.Background(), req("b", + stepOf(tt.strategy, "s1", gitURI("feature/satisfied", stale)), + )) + require.NoError(t, err) + assert.Equal(t, runwaypb.Outcome_SUCCEEDED, res.GetOutcome()) + assert.Equal(t, mainBefore, f.remoteHEAD(t)) + }) + } +} + +func TestMerge_StalePendingChangeRejected(t *testing.T) { + tests := []struct { + name string + strategy mergestrategypb.Strategy + }{ + {name: "rebase", strategy: mergestrategypb.Strategy_REBASE}, + {name: "squash rebase", strategy: mergestrategypb.Strategy_SQUASH_REBASE}, + {name: "merge", strategy: mergestrategypb.Strategy_MERGE}, + {name: "promote", strategy: mergestrategypb.Strategy_PROMOTE}, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + f := setupGitFixture(t) + base := f.remoteHEAD(t) + stale := f.pushPRCommitFrom(t, base, "feature/pending", "stale.txt", "stale\n", "stale") + mustGit(t, f.authorDir, "push", "origin", stale+":refs/heads/archive/stale") + current := f.pushPRCommitFrom(t, base, "feature/pending", "current.txt", "current\n", "current") + require.NotEqual(t, stale, current) + + m := f.newMergerWith(t, func(p *Params) { p.CheckStaleness = true }) + _, err := m.Merge(context.Background(), req("b", + stepOf(tt.strategy, "s1", gitURI("feature/pending", stale)), + )) + require.Error(t, err) + assert.True(t, errors.Is(err, merger.ErrInvalidRequest)) + assert.Equal(t, base, f.remoteHEAD(t)) + }) + } +} + func TestMerge_StalenessCheckOffByDefault(t *testing.T) { f := setupGitFixture(t) stale := f.pushPRCommit(t, "feature/s", "s.txt", "v1\n", "v1") @@ -1603,6 +1677,10 @@ func uri(sha string) string { return fmt.Sprintf("github://github.example.com/uber/submitqueue/pull/1/%s", sha) } +func gitURI(branch, sha string) string { + return fmt.Sprintf("git://git.example.com/uber/submitqueue/%s/%s", url.PathEscape("refs/heads/"+branch), sha) +} + // installRaceHook writes a pre-receive hook on the bare remote that simulates // concurrent pushes. On its Nth invocation it reads the Nth line of race-shas, // points refs/heads/main at that SHA via update-ref, and exits 1 (rejecting the diff --git a/runway/extension/merger/git/objects.go b/runway/extension/merger/git/objects.go index 3a919e5c5..3c806811a 100644 --- a/runway/extension/merger/git/objects.go +++ b/runway/extension/merger/git/objects.go @@ -93,7 +93,7 @@ func (m *gitMerger) hasCommit(ctx context.Context, sha string) bool { // that no longer exists yields no verdict rather than a failure: the change may // legitimately have been closed, and ensureObjects has already established that // the commit itself is present. -func (m *gitMerger) checkStale(ctx context.Context, refs []changeRef) error { +func (m *gitMerger) checkStale(ctx context.Context, refs []changeRef, tracked headBranchTracker) error { if !m.checkStaleness { return nil } @@ -110,6 +110,10 @@ func (m *gitMerger) checkStale(ctx context.Context, refs []changeRef) error { continue } if current := fields[0]; current != ref.SHA { + branch, movedTo, moved := tracked.lookup(ref.SHA) + if moved && branch == ref.Ref && movedTo == current { + continue + } coremetrics.NamedCounter(m.metricsScope, "merge", "stale_changes", 1) return fmt.Errorf("%w: change is stale: %s names commit %s but %s now points at %s", merger.ErrInvalidRequest, ref.Label, ref.SHA, ref.Ref, current) diff --git a/test/e2e/submitqueue/ISS-001.md b/test/e2e/submitqueue/ISS-001.md new file mode 100644 index 000000000..db9a32db3 --- /dev/null +++ b/test/e2e/submitqueue/ISS-001.md @@ -0,0 +1,58 @@ +# ISS-001: stale retry after a successful Git push + +## Reproduction + +The regression test uses the real Git-backed SubmitQueue E2E stack with `checkStaleness: true` and `updateHeadBranch: true`. + +Initial Git refs: + +```text +refs/heads/main = M0 +refs/heads/feature/retry-result = A +``` + +SubmitQueue eventually publishes a `runway-merge` message whose `id` is the batch ID and whose change URI pins the feature branch to `A`: + +```json +{ + "id": "e2e-git-queue/batch/1", + "queue_name": "e2e-git-queue", + "steps": [{ + "change": { + "uris": ["git://git.example.com/sandbox/refs%2Fheads%2Ffeature%2Fretry-result/A"] + }, + "strategy": "REBASE" + }] +} +``` + +The test closes the `runway-merge` consumer gate until it can read the exact batch ID from the application database's `request_batch` table. It then seeds a batch-scoped MyISAM guard row and installs a `BEFORE INSERT` trigger on the queue database's `queue_messages` table. + +The trigger counts every attempted `SUCCEEDED` publication for the batch and rejects only the first one. It decrements the MyISAM guard before raising MySQL error 1213, so the counter changes survive the failed InnoDB insert and the next successful result is allowed. + +On the first delivery, Runway rebases `A` into `A'` and updates both refs before result publication: + +```text +refs/heads/main = A' +refs/heads/feature/retry-result = A' +``` + +The injected publication failure causes the original `runway-merge` delivery to be retried. Its unchanged URI still pins `A`, while the feature ref now points to `A'`. + +Before the fix, Runway rejects this retry as stale and publishes `FAILED`, leaving the contradictory state: + +```text +Git main: A' (the code landed) +SubmitQueue request: error +``` + +After the fix, Runway applies the request locally against current `main` before enforcing staleness. The application is a no-op because `A` is already represented by `A'`, so Runway publishes `SUCCEEDED` without another remote write. The final state converges: + +```text +Git main: A' +feature branch: A' +SubmitQueue request: landed +successful-result attempts: 2 +``` + +The two result-publication attempts are the durable witness that `runway-merge` was delivered twice. The queue's acknowledged message and delivery-state rows may be garbage-collected before the terminal request status is observable. diff --git a/test/e2e/submitqueue/git_suite_test.go b/test/e2e/submitqueue/git_suite_test.go index 321051ec0..710549a39 100644 --- a/test/e2e/submitqueue/git_suite_test.go +++ b/test/e2e/submitqueue/git_suite_test.go @@ -38,12 +38,16 @@ import ( "path/filepath" "strings" "testing" + "time" "github.com/stretchr/testify/require" "github.com/stretchr/testify/suite" changepb "github.com/uber/submitqueue/api/base/change/protopb" mergestrategypb "github.com/uber/submitqueue/api/base/mergestrategy/protopb" + runwaymq "github.com/uber/submitqueue/api/runway/messagequeue" gatewaypb "github.com/uber/submitqueue/api/submitqueue/gateway/protopb" + "github.com/uber/submitqueue/platform/extension/consumergate" + consumergatefile "github.com/uber/submitqueue/platform/extension/consumergate/file" gitexectest "github.com/uber/submitqueue/platform/git/exectest" "github.com/uber/submitqueue/submitqueue/entity" "github.com/uber/submitqueue/test/testutil" @@ -67,6 +71,7 @@ type GitMergeSuite struct { gatewayClient gatewaypb.SubmitQueueGatewayClient db *sql.DB queueDB *sql.DB + gate *consumergatefile.Store // git is the pinned git the test drives the bare repository with — the same // build the merger uses inside the container. @@ -89,7 +94,9 @@ func (s *GitMergeSuite) SetupSuite() { s.git = gitexectest.Git(t) containerUser := dockerContainerUser(t) t.Setenv("SQ_CONTAINER_USER", containerUser) - t.Setenv("SQ_CONSUMER_GATE_DIR", t.TempDir()) + gateDir := t.TempDir() + t.Setenv("SQ_CONSUMER_GATE_DIR", gateDir) + s.gate = consumergatefile.New(gateDir) // The bare repository lives in a host directory bind-mounted into Runway, so // the test can seed it and read back exactly what the merger pushed. @@ -218,6 +225,37 @@ func (s *GitMergeSuite) TestLand_MovesEachChangeHeadBranchToItsLandedCommit() { s.True(s.isAncestorOfMain(s.branchSHA("feature/head-2"))) } +func (s *GitMergeSuite) TestLand_RetryAfterResultPublishFailure_ReconcilesLandedState() { + const consumerGroup = "runway-merge" + gateKey := consumergate.Key{ConsumerGroup: consumerGroup, PartitionKey: gitQueue} + require.NoError(s.T(), s.gate.Close(s.ctx, gateKey, consumergate.Metadata{ + Reason: "install a batch-scoped fail-once result publisher before merge", + CreatedBy: "GitMergeSuite", + CreatedAtMs: time.Now().UnixMilli(), + })) + defer func() { + require.NoError(s.T(), s.gate.Open(s.ctx, gateKey)) + }() + + before := s.mainSHA() + head := s.pushChange("feature/retry-result", map[string]string{"retry-result.txt": "retry result\n"}, "add retry result") + sqid := s.land(gitQueue, s.uri("feature/retry-result", head)) + batchID := s.awaitBatchID(sqid) + s.awaitParkedMerge(batchID) + s.installMergeSignalFailOnce(batchID) + + require.NoError(s.T(), s.gate.Open(s.ctx, gateKey)) + s.requireStatus(sqid, entity.RequestStatusLanded) + + landed := s.mainSHA() + s.NotEqual(before, landed) + s.NotEqual(head, landed) + s.Equal("retry result\n", s.fileOnMain("retry-result.txt")) + s.Equal(landed, s.branchSHA("feature/retry-result")) + s.Equal(2, s.mergeSignalPublishAttemptCount(batchID), + "two result publications prove the runway-merge delivery was retried") +} + func (s *GitMergeSuite) TestLand_Conflict_FailsAndLeavesTheTargetUntouched() { // Two changes editing the same line from the same base: the first lands, // the second cannot be replayed onto it. @@ -236,11 +274,7 @@ func (s *GitMergeSuite) TestLand_Conflict_FailsAndLeavesTheTargetUntouched() { s.Equal(loser, s.branchSHA("feature/conflict-b")) } -func (s *GitMergeSuite) TestLand_ResubmittedAfterLanding_IsRejectedAsStale() { - // Landing a change moves its head branch to the commit it became, so the - // URI that was submitted no longer describes where that branch points. The - // staleness check catches exactly that, which is what stops a change from - // being replayed onto the target a second time. +func (s *GitMergeSuite) TestLand_ResubmittedAfterLanding_IsSuccessfulNoOp() { head := s.pushChange("feature/already", map[string]string{"already.txt": "already\n"}, "add already") s.requireStatus(s.land(gitQueue, s.uri("feature/already", head)), entity.RequestStatusLanded) @@ -249,8 +283,8 @@ func (s *GitMergeSuite) TestLand_ResubmittedAfterLanding_IsRejectedAsStale() { landedAs := s.branchSHA("feature/already") s.NotEqual(head, landedAs, "the head branch moved to the landed commit") - s.requireStatus(s.land(gitQueue, s.uri("feature/already", head)), entity.RequestStatusError) - s.Equal(settled, s.mainSHA(), "a stale resubmission must not move the target") + s.requireStatus(s.land(gitQueue, s.uri("feature/already", head)), entity.RequestStatusLanded) + s.Equal(settled, s.mainSHA(), "an already-satisfied resubmission must not move the target") s.Equal(updates, s.mainRefUpdateCount(), "and must not push at all") s.Equal(landedAs, s.branchSHA("feature/already"), "nor disturb the change's branch") } @@ -286,6 +320,87 @@ func (s *GitMergeSuite) requireStatus(sqid string, want entity.RequestStatus) { s.Require().Equal(want, got, "request %s reached the wrong terminal status", sqid) } +func (s *GitMergeSuite) awaitBatchID(sqid string) string { + var batchID string + pollUntil(persistPollInterval, func() bool { + err := s.db.QueryRowContext(s.ctx, + "SELECT batch_id FROM request_batch WHERE queue = ? AND request_id = ? ORDER BY batch_id LIMIT 1", + gitQueue, sqid, + ).Scan(&batchID) + return err == nil + }) + return batchID +} + +func (s *GitMergeSuite) awaitParkedMerge(batchID string) { + pollUntil(persistPollInterval, func() bool { + parked, err := s.gate.ListParked(s.ctx, "runway-merge") + if err != nil { + return false + } + for _, delivery := range parked { + if delivery.Topic == runwaymq.TopicKeyMerge.String() && delivery.MessageID == batchID { + return true + } + } + return false + }) +} + +func (s *GitMergeSuite) installMergeSignalFailOnce(messageID string) { + t := s.T() + const ( + triggerName = "e2e_fail_merge_signal_once" + tableName = "e2e_merge_signal_fail_once" + ) + _, err := s.queueDB.ExecContext(s.ctx, "DROP TRIGGER IF EXISTS "+triggerName) + require.NoError(t, err) + _, err = s.queueDB.ExecContext(s.ctx, "DROP TABLE IF EXISTS "+tableName) + require.NoError(t, err) + _, err = s.queueDB.ExecContext(s.ctx, "CREATE TABLE "+tableName+" (message_id VARCHAR(191) CHARACTER SET ascii PRIMARY KEY, remaining INT NOT NULL, attempts INT NOT NULL) ENGINE=MyISAM") + require.NoError(t, err) + _, err = s.queueDB.ExecContext(s.ctx, "INSERT INTO "+tableName+" (message_id, remaining, attempts) VALUES (?, 1, 0)", messageID) + require.NoError(t, err) + _, err = s.queueDB.ExecContext(s.ctx, ` +CREATE TRIGGER `+triggerName+` +BEFORE INSERT ON queue_messages +FOR EACH ROW +BEGIN + IF NEW.topic = 'merge-signal' + AND JSON_UNQUOTE(JSON_EXTRACT(CONVERT(NEW.payload USING utf8mb4), '$.outcome')) = 'SUCCEEDED' + THEN + UPDATE `+tableName+` + SET attempts = attempts + 1 + WHERE message_id = NEW.id; + UPDATE `+tableName+` + SET remaining = remaining - 1 + WHERE message_id = NEW.id AND remaining > 0; + IF ROW_COUNT() > 0 THEN + SIGNAL SQLSTATE '40001' + SET MYSQL_ERRNO = 1213, + MESSAGE_TEXT = 'injected fail-once merge-signal publication failure'; + END IF; + END IF; +END`) + require.NoError(t, err) + t.Cleanup(func() { + _, dropTriggerErr := s.queueDB.ExecContext(s.ctx, "DROP TRIGGER IF EXISTS "+triggerName) + require.NoError(t, dropTriggerErr) + _, dropTableErr := s.queueDB.ExecContext(s.ctx, "DROP TABLE IF EXISTS "+tableName) + require.NoError(t, dropTableErr) + }) +} + +func (s *GitMergeSuite) mergeSignalPublishAttemptCount(batchID string) int { + var attempts int + err := s.queueDB.QueryRowContext(s.ctx, + "SELECT attempts FROM e2e_merge_signal_fail_once WHERE message_id = ?", + batchID, + ).Scan(&attempts) + s.Require().NoError(err) + return attempts +} + // uri builds the git:// change URI for a branch pinned at a commit. The ref is // percent-encoded so a branch name containing slashes stays one path segment. func (s *GitMergeSuite) uri(branch, sha string) string { @@ -424,6 +539,7 @@ func (s *GitMergeSuite) runGit(dir string, args ...string) string { s.T().Helper() cmd := exec.Command(s.git, args...) cmd.Dir = dir + cmd.Env = append(os.Environ(), "GIT_EXEC_PATH="+filepath.Dir(s.git)) var stderr strings.Builder cmd.Stderr = &stderr out, err := cmd.Output()