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
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@ load("@rules_go//go:def.bzl", "go_library", "go_test")
go_library(
name = "go_default_library",
srcs = ["fakemarker.go"],
importpath = "github.com/uber/submitqueue/submitqueue/core/fakemarker",
importpath = "github.com/uber/submitqueue/platform/fakemarker",
visibility = ["//visibility:public"],
deps = ["//platform/base/change:go_default_library"],
)
Expand Down
29 changes: 29 additions & 0 deletions runway/extension/merger/fake/BUILD.bazel
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
load("@rules_go//go:def.bzl", "go_library", "go_test")

go_library(
name = "go_default_library",
srcs = ["fake.go"],
importpath = "github.com/uber/submitqueue/runway/extension/merger/fake",
visibility = ["//visibility:public"],
deps = [
"//api/runway/messagequeue:go_default_library",
"//api/runway/messagequeue/protopb:go_default_library",
"//platform/fakemarker:go_default_library",
"//runway/extension/merger:go_default_library",
],
)

go_test(
name = "go_default_test",
srcs = ["fake_test.go"],
embed = [":go_default_library"],
deps = [
"//api/base/change/protopb:go_default_library",
"//api/base/mergestrategy/protopb:go_default_library",
"//api/runway/messagequeue:go_default_library",
"//api/runway/messagequeue/protopb:go_default_library",
"//runway/extension/merger:go_default_library",
"@com_github_stretchr_testify//assert:go_default_library",
"@com_github_stretchr_testify//require:go_default_library",
],
)
115 changes: 115 additions & 0 deletions runway/extension/merger/fake/fake.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,115 @@
// Copyright (c) 2025 Uber Technologies, Inc.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

// Package fake provides a merger.Merger whose outcome is driven by the request
// payload. With no marker it succeeds like the noop merger, behaving as a
// best-case stub for wiring and baselines. A failure can be injected end-to-end
// (e.g. from an e2e merge request) by embedding a marker token in a change URI
// of the form "sq-fake=<token>":
//
// sq-fake=merge-conflict -> merger.ErrConflict
// sq-fake=merge-invalid -> merger.ErrInvalidRequest
// sq-fake=merge-error -> a plain (non-retryable) error
//
// The first token found across the request's steps, in order, decides the
// outcome for the whole request. This lets a single running stack exercise
// Runway's terminal-failure and dead-letter paths purely by varying request
// payloads. It is intended for examples and tests only, never production.
package fake

import (
"context"
"fmt"
"sync/atomic"

runwaymq "github.com/uber/submitqueue/api/runway/messagequeue"
runwaypb "github.com/uber/submitqueue/api/runway/messagequeue/protopb"
"github.com/uber/submitqueue/platform/fakemarker"
"github.com/uber/submitqueue/runway/extension/merger"
)

// Recognized marker tokens. See the package doc for the convention.
const (
tokenConflict = "merge-conflict"
tokenInvalid = "merge-invalid"
tokenError = "merge-error"
)

var _ merger.Merger = (*Merger)(nil)

// Merger is a Merger that succeeds unless a marker token in a change URI
// requests otherwise.
type Merger struct {
seq atomic.Uint64
}

// New returns a Merger that defaults to success and honors marker tokens
// embedded in change URIs.
func New() *Merger { return &Merger{} }

// CheckMergeability reports the request as mergeable unless a recognized marker
// token asks for a failure. Outputs are empty, as for any dry run.
func (m *Merger) CheckMergeability(_ context.Context, req *runwaymq.MergeRequest) (*runwaymq.MergeResult, error) {
if err := injectedFailure(req); err != nil {
return nil, err
}

steps := make([]*runwaymq.StepResult, len(req.GetSteps()))
for i, s := range req.GetSteps() {
steps[i] = &runwaymq.StepResult{StepId: s.GetStepId()}
}
return &runwaymq.MergeResult{
Id: req.GetId(),
Outcome: runwaypb.Outcome_SUCCEEDED,
Steps: steps,
}, nil
}

// Merge reports the request as merged unless a recognized marker token asks for
// a failure, producing one synthetic revision id per step.
func (m *Merger) Merge(_ context.Context, req *runwaymq.MergeRequest) (*runwaymq.MergeResult, error) {
if err := injectedFailure(req); err != nil {
return nil, err
}

steps := make([]*runwaymq.StepResult, len(req.GetSteps()))
for i, s := range req.GetSteps() {
n := m.seq.Add(1)
steps[i] = &runwaymq.StepResult{
StepId: s.GetStepId(),
Outputs: []*runwaymq.StepOutput{{Id: fmt.Sprintf("%040x", n)}},
}
}
return &runwaymq.MergeResult{
Id: req.GetId(),
Outcome: runwaypb.Outcome_SUCCEEDED,
Steps: steps,
}, nil
}

// injectedFailure returns the error the request's marker token asks for, or nil
// when no step carries a recognized token.
func injectedFailure(req *runwaymq.MergeRequest) error {
for _, s := range req.GetSteps() {
switch fakemarker.Token(s.GetChange().GetUris()) {
case tokenConflict:
return fmt.Errorf("fake: marked conflicting on step %s: %w", s.GetStepId(), merger.ErrConflict)
case tokenInvalid:
return fmt.Errorf("fake: marked invalid on step %s: %w", s.GetStepId(), merger.ErrInvalidRequest)
case tokenError:
return fmt.Errorf("fake: marked failing on step %s", s.GetStepId())
}
}
return nil
}
151 changes: 151 additions & 0 deletions runway/extension/merger/fake/fake_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,151 @@
// Copyright (c) 2025 Uber Technologies, Inc.
//
// Licensed under the Apache License, Version 2.0 (the "License");
// you may not use this file except in compliance with the License.
// You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing, software
// distributed under the License is distributed on an "AS IS" BASIS,
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
// See the License for the specific language governing permissions and
// limitations under the License.

package fake

import (
"context"
"testing"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
changepb "github.com/uber/submitqueue/api/base/change/protopb"
strategypb "github.com/uber/submitqueue/api/base/mergestrategy/protopb"
runwaymq "github.com/uber/submitqueue/api/runway/messagequeue"
runwaypb "github.com/uber/submitqueue/api/runway/messagequeue/protopb"
"github.com/uber/submitqueue/runway/extension/merger"
)

const baseURI = "github://github.example.com/uber/repo/pull/1/abcdef0123456789abcdef0123456789abcdef01"

// requestWith builds a two-step request whose second step carries the given
// URIs, so tests can prove the token is found on any step, not just the first.
func requestWith(uris ...string) *runwaymq.MergeRequest {
return &runwaymq.MergeRequest{
Id: "queue-a/42",
QueueName: "queue-a",
Steps: []*runwaymq.MergeStep{
{
StepId: "queue-a/1",
Change: &changepb.Change{Uris: []string{baseURI}},
Strategy: strategypb.Strategy_REBASE,
},
{
StepId: "queue-a/2",
Change: &changepb.Change{Uris: uris},
Strategy: strategypb.Strategy_REBASE,
},
},
}
}

func TestUnmarkedRequestSucceeds(t *testing.T) {
req := requestWith(baseURI)

t.Run("check mergeability reports no outputs", func(t *testing.T) {
res, err := New().CheckMergeability(context.Background(), req)
require.NoError(t, err)

assert.Equal(t, req.GetId(), res.GetId())
assert.Equal(t, runwaypb.Outcome_SUCCEEDED, res.GetOutcome())
require.Len(t, res.GetSteps(), 2)
assert.Equal(t, "queue-a/1", res.GetSteps()[0].GetStepId())
assert.Empty(t, res.GetSteps()[0].GetOutputs())
assert.Equal(t, "queue-a/2", res.GetSteps()[1].GetStepId())
assert.Empty(t, res.GetSteps()[1].GetOutputs())
})

t.Run("merge reports one output per step", func(t *testing.T) {
res, err := New().Merge(context.Background(), req)
require.NoError(t, err)

assert.Equal(t, req.GetId(), res.GetId())
assert.Equal(t, runwaypb.Outcome_SUCCEEDED, res.GetOutcome())
require.Len(t, res.GetSteps(), 2)
require.Len(t, res.GetSteps()[0].GetOutputs(), 1)
require.Len(t, res.GetSteps()[1].GetOutputs(), 1)
assert.NotEqual(t,
res.GetSteps()[0].GetOutputs()[0].GetId(),
res.GetSteps()[1].GetOutputs()[0].GetId(),
"each step should produce a distinct revision id")
})
}

func TestMarkedRequestFails(t *testing.T) {
tests := []struct {
name string
token string
want error // nil means "an error that is neither terminal sentinel"
}{
{name: "conflict", token: tokenConflict, want: merger.ErrConflict},
{name: "invalid", token: tokenInvalid, want: merger.ErrInvalidRequest},
{name: "plain error", token: tokenError, want: nil},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
req := requestWith(baseURI + "?sq-fake=" + tt.token)

for _, call := range []struct {
name string
fn func(*runwaymq.MergeRequest) (*runwaymq.MergeResult, error)
}{
{"CheckMergeability", func(r *runwaymq.MergeRequest) (*runwaymq.MergeResult, error) {
return New().CheckMergeability(context.Background(), r)
}},
{"Merge", func(r *runwaymq.MergeRequest) (*runwaymq.MergeResult, error) {
return New().Merge(context.Background(), r)
}},
} {
t.Run(call.name, func(t *testing.T) {
res, err := call.fn(req)
require.Error(t, err)
assert.Nil(t, res)

if tt.want != nil {
assert.ErrorIs(t, err, tt.want)
assert.True(t, merger.IsTerminal(err), "sentinel failures must be terminal")
return
}
assert.False(t, merger.IsTerminal(err),
"an unmarked failure must not look terminal, so it dead-letters")
})
}
})
}
}

func TestUnrecognizedTokenSucceeds(t *testing.T) {
req := requestWith(baseURI + "?sq-fake=some-other-fakes-token")

res, err := New().Merge(context.Background(), req)
require.NoError(t, err)
assert.Equal(t, runwaypb.Outcome_SUCCEEDED, res.GetOutcome())
}

func TestFirstRecognizedTokenWins(t *testing.T) {
req := &runwaymq.MergeRequest{
Id: "queue-a/42",
QueueName: "queue-a",
Steps: []*runwaymq.MergeStep{
{StepId: "queue-a/1", Change: &changepb.Change{Uris: []string{baseURI + "?sq-fake=" + tokenConflict}}},
{StepId: "queue-a/2", Change: &changepb.Change{Uris: []string{baseURI + "?sq-fake=" + tokenInvalid}}},
},
}

_, err := New().Merge(context.Background(), req)
require.Error(t, err)
assert.ErrorIs(t, err, merger.ErrConflict)
assert.NotErrorIs(t, err, merger.ErrInvalidRequest)
}
1 change: 1 addition & 0 deletions service/runway/server/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@ go_library(
"//runway/controller/merge:go_default_library",
"//runway/controller/mergeconflictcheck:go_default_library",
"//runway/extension/merger:go_default_library",
"//runway/extension/merger/fake:go_default_library",
"//runway/extension/merger/git:go_default_library",
"//runway/extension/merger/noop:go_default_library",
"@com_github_go_sql_driver_mysql//:go_default_library",
Expand Down
4 changes: 4 additions & 0 deletions service/runway/server/docker-compose.yml
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,10 @@ services:
- "8080" # Random ephemeral port to avoid conflicts
environment:
- PORT=:8080
# Merger implementation. Empty (the default) resolves from the merge
# environment: git when MERGE_CHECKOUT_PATH is set, noop otherwise. The
# e2e suite sets SQ_RUNWAY_MERGER=fake to drive outcomes from the payload.
- MERGER=${SQ_RUNWAY_MERGER:-}
# Queue infrastructure connection
- QUEUE_MYSQL_DSN=root:root@tcp(mysql-queue:3306)/submitqueue?parseTime=true
- HOSTNAME=runway-dev
Expand Down
35 changes: 31 additions & 4 deletions service/runway/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ import (
"github.com/uber/submitqueue/runway/controller/merge"
"github.com/uber/submitqueue/runway/controller/mergeconflictcheck"
"github.com/uber/submitqueue/runway/extension/merger"
"github.com/uber/submitqueue/runway/extension/merger/fake"
gitmerger "github.com/uber/submitqueue/runway/extension/merger/git"
"github.com/uber/submitqueue/runway/extension/merger/noop"
"go.uber.org/zap"
Expand Down Expand Up @@ -302,11 +303,27 @@ func run() error {
return err
}

// newMergerFactory returns a merger.Factory for the server. When
// MERGE_CHECKOUT_PATH is set it wires the git-backed merger built from the
// MERGE_* / GIT_* environment; otherwise it falls back to the noop merger so
// local development and compose runs need no git checkout.
// newMergerFactory returns a merger.Factory for the server. MERGER selects the
// implementation explicitly; when it is unset the choice falls back to the merge
// environment — MERGE_CHECKOUT_PATH wires the git-backed merger built from the
// MERGE_* / GIT_* environment, and its absence wires the noop merger so local
// development and compose runs need no git checkout.
func newMergerFactory(logger *zap.Logger, scope tally.Scope) (merger.Factory, error) {
switch impl := strings.ToLower(strings.TrimSpace(os.Getenv("MERGER"))); impl {
case "fake":
// Marker-driven outcomes, for e2e tests that need Runway to fail on
// demand without a git checkout. Never production.
logger.Info("MERGER=fake; using marker-driven fake merger")
return &fakeMergerFactory{merger: fake.New()}, nil
case "noop":
logger.Info("MERGER=noop; using noop merger")
return &noopMergerFactory{}, nil
case "", "git":
// Fall through to the merge-environment default below.
default:
return nil, fmt.Errorf("invalid MERGER %q", impl)
}

checkoutPath := os.Getenv("MERGE_CHECKOUT_PATH")
if checkoutPath == "" {
logger.Info("MERGE_CHECKOUT_PATH not set; using noop merger")
Expand Down Expand Up @@ -369,6 +386,16 @@ func (f *noopMergerFactory) For(_ merger.Config) (merger.Merger, error) {
return noop.New(), nil
}

// fakeMergerFactory shares one fake merger across queues so the synthetic
// revision ids it mints stay unique for the lifetime of the process.
type fakeMergerFactory struct {
merger merger.Merger
}

func (f *fakeMergerFactory) For(_ merger.Config) (merger.Merger, error) {
return f.merger, nil
}

// parseStrategy maps the MERGE_DEFAULT_STRATEGY env value to a concrete merge
// strategy, defaulting to REBASE when unset. DEFAULT is rejected because it
// cannot itself be the default a step resolves to.
Expand Down
2 changes: 1 addition & 1 deletion submitqueue/extension/buildrunner/fake/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -6,8 +6,8 @@ go_library(
importpath = "github.com/uber/submitqueue/submitqueue/extension/buildrunner/fake",
visibility = ["//visibility:public"],
deps = [
"//platform/fakemarker:go_default_library",
"//submitqueue/core/changeset:go_default_library",
"//submitqueue/core/fakemarker:go_default_library",
"//submitqueue/entity:go_default_library",
"//submitqueue/extension/buildrunner:go_default_library",
],
Expand Down
Loading
Loading