diff --git a/platform/errs/README.md b/platform/errs/README.md index b86c3b104..d2063d457 100644 --- a/platform/errs/README.md +++ b/platform/errs/README.md @@ -89,7 +89,7 @@ One operational consequence worth knowing before relying on any of this: **retry ## Adding a Backend-Specific Classifier -Backend classifiers live alongside the extension they classify, under `platform/errs//`. The canonical examples are `platform/errs/mysql` (MySQL driver errors), `platform/errs/http` (rejected status codes and transport failures from clients built on `platform/http`), `platform/errs/yarpc` (YARPC status codes), and `platform/errs/generic` (transport-agnostic concerns such as `context.Canceled`). +Backend classifiers live alongside the extension they classify, under `platform/errs//`. The canonical examples are `platform/errs/mysql` (MySQL driver errors), `platform/errs/http` (rejected status codes and transport failures from clients built on `platform/http`), `platform/errs/git` (structured Git process failures), `platform/errs/yarpc` (YARPC status codes), and `platform/errs/generic` (transport-agnostic concerns such as `context.Canceled`). A classifier: @@ -122,6 +122,7 @@ Servers wire each classifier into the consumer's `ErrorProcessor`. Order matters import ( "github.com/uber/submitqueue/platform/errs" genericerrs "github.com/uber/submitqueue/platform/errs/generic" + giterrs "github.com/uber/submitqueue/platform/errs/git" httperrs "github.com/uber/submitqueue/platform/errs/http" mysqlerrs "github.com/uber/submitqueue/platform/errs/mysql" yarpcerrs "github.com/uber/submitqueue/platform/errs/yarpc" @@ -130,6 +131,7 @@ import ( c := consumer.New(logger, scope, registry, errs.NewClassifierProcessor( genericerrs.Classifier, + giterrs.Classifier, httperrs.Classifier, yarpcerrs.Classifier, mysqlerrs.Classifier, @@ -143,7 +145,9 @@ Classifiers are not installed globally. A host that wants YARPC statuses classif The YARPC classifier reads the typed status code rather than matching its rendered message. Cancellation is retryable caller-side infrastructure; transient or ambiguous server codes (`Unknown`, `DeadlineExceeded`, `ResourceExhausted`, `Aborted`, `Internal`, and `Unavailable`) are retryable dependency failures; request verdicts and permanent server failures are non-retryable dependency failures. A deadline may expire after a mutating RPC succeeded, so this classification relies on the repository-wide requirement that queue-driven operations are idempotent. -Tests follow the same shape: assert per-node behaviour against `Classifier.Classify(node)` directly, and assert end-to-end behaviour by running `errs.NewClassifierProcessor(Classifier).Process(err)` and checking the helpers (`IsRetryable`, `IsUserError`, …) on the result. See `platform/errs/mysql/mysql_test.go`, `platform/errs/yarpc/yarpc_test.go`, and `platform/errs/generic/generic_test.go`. +The Git classifier reads `gitexec.CommandError`, which preserves the Git subcommand and the underlying `os/exec` error through contextual wrapping. A started `fetch`, `push`, or `ls-remote` process is a retryable dependency failure unless its diagnostic identifies a permanent authentication, repository, invocation, or configuration problem. Started local repository operations are retryable infrastructure failures under the same exception; commands that never started and unknown operations remain non-retryable by default. + +Tests follow the same shape: assert per-node behaviour against `Classifier.Classify(node)` directly, and assert end-to-end behaviour by running `errs.NewClassifierProcessor(Classifier).Process(err)` and checking the helpers (`IsRetryable`, `IsUserError`, …) on the result. See `platform/errs/mysql/mysql_test.go`, `platform/errs/git/git_test.go`, `platform/errs/yarpc/yarpc_test.go`, and `platform/errs/generic/generic_test.go`. ## Overriding Classification from a Controller diff --git a/platform/errs/git/BUILD.bazel b/platform/errs/git/BUILD.bazel new file mode 100644 index 000000000..566cbb452 --- /dev/null +++ b/platform/errs/git/BUILD.bazel @@ -0,0 +1,24 @@ +load("@rules_go//go:def.bzl", "go_library", "go_test") + +go_library( + name = "go_default_library", + srcs = ["git.go"], + importpath = "github.com/uber/submitqueue/platform/errs/git", + visibility = ["//visibility:public"], + deps = [ + "//platform/errs:go_default_library", + "//platform/git/exec:go_default_library", + ], +) + +go_test( + name = "go_default_test", + srcs = ["git_test.go"], + embed = [":go_default_library"], + deps = [ + "//platform/errs:go_default_library", + "//platform/git/exec:go_default_library", + "@com_github_stretchr_testify//assert:go_default_library", + "@com_github_stretchr_testify//require:go_default_library", + ], +) diff --git a/platform/errs/git/git.go b/platform/errs/git/git.go new file mode 100644 index 000000000..b9d5370e1 --- /dev/null +++ b/platform/errs/git/git.go @@ -0,0 +1,80 @@ +// Copyright (c) 2026 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 git provides an errs.Classifier for failures from Git processes. +package git + +import ( + "strings" + + "github.com/uber/submitqueue/platform/errs" + gitexec "github.com/uber/submitqueue/platform/git/exec" +) + +// Classifier recognises structured Git command failures. Remote exchanges are +// retryable dependency failures; local repository operations are retryable +// infrastructure failures. Commands that never started and unknown operations +// remain unclassified and therefore fail fast. +var Classifier errs.Classifier = classifier{} + +type classifier struct{} + +var permanentDiagnosticFragments = []string{ + "authentication failed", + "bad config line", + "does not appear to be a git repository", + "invalid refspec", + "not a git repository", + "permission denied (publickey)", + "repository not found", + "unknown option", + "unknown switch", +} + +func (classifier) Classify(err error) errs.Verdict { + commandErr, ok := err.(*gitexec.CommandError) + if !ok || !commandErr.ProcessExited() { + return errs.Unknown + } + + if isPermanentGitDiagnostic(commandErr.Diagnostic()) { + switch commandErr.Operation() { + case "fetch", "ls-remote", "push": + return errs.InfraDependency + default: + return errs.Infra + } + } + + switch commandErr.Operation() { + case "fetch", "ls-remote", "push": + return errs.InfraDependencyRetryable + case "cat-file", "cherry-pick", "clean", "commit", "ls-files", "merge", "merge-base", "reset", "rev-list", "rev-parse", "show": + return errs.InfraRetryable + case "config": + return errs.Infra + default: + return errs.Unknown + } +} + +func isPermanentGitDiagnostic(diagnostic string) bool { + diagnostic = strings.ToLower(diagnostic) + for _, fragment := range permanentDiagnosticFragments { + if strings.Contains(diagnostic, fragment) { + return true + } + } + return false +} diff --git a/platform/errs/git/git_test.go b/platform/errs/git/git_test.go new file mode 100644 index 000000000..08a1ac117 --- /dev/null +++ b/platform/errs/git/git_test.go @@ -0,0 +1,164 @@ +// Copyright (c) 2026 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 git + +import ( + "errors" + "fmt" + "os" + "os/exec" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "github.com/uber/submitqueue/platform/errs" + gitexec "github.com/uber/submitqueue/platform/git/exec" +) + +type classifierFixtures struct { + exitError error + startError error +} + +func setupClassifierFixtures(t *testing.T) classifierFixtures { + t.Helper() + + err := exec.Command(os.Args[0], "-test.run=[").Run() + require.Error(t, err) + var exitErr *exec.ExitError + require.ErrorAs(t, err, &exitErr) + + return classifierFixtures{ + exitError: exitErr, + startError: &exec.Error{Name: "git", Err: exec.ErrNotFound}, + } +} + +func TestClassifier(t *testing.T) { + fixtures := setupClassifierFixtures(t) + tests := []struct { + name string + err error + want errs.Verdict + }{ + { + name: "remote fetch exit is retryable dependency failure", + err: gitexec.NewCommandError("fetch", "temporary remote failure", fixtures.exitError), + want: errs.InfraDependencyRetryable, + }, + { + name: "remote push exit is retryable dependency failure", + err: gitexec.NewCommandError("push", "temporary remote failure", fixtures.exitError), + want: errs.InfraDependencyRetryable, + }, + { + name: "remote authentication failure is permanent dependency failure", + err: gitexec.NewCommandError("fetch", "fatal: Authentication failed", fixtures.exitError), + want: errs.InfraDependency, + }, + { + name: "local reset exit is retryable infrastructure failure", + err: gitexec.NewCommandError("reset", "checkout unavailable", fixtures.exitError), + want: errs.InfraRetryable, + }, + { + name: "invalid local refspec is permanent infrastructure failure", + err: gitexec.NewCommandError("reset", "fatal: invalid refspec", fixtures.exitError), + want: errs.Infra, + }, + { + name: "local configuration exit is permanent infrastructure failure", + err: gitexec.NewCommandError("config", "invalid configuration", fixtures.exitError), + want: errs.Infra, + }, + { + name: "process start failure remains non-retryable", + err: gitexec.NewCommandError("fetch", "git executable missing", fixtures.startError), + want: errs.Unknown, + }, + { + name: "unknown operation remains non-retryable", + err: gitexec.NewCommandError("unknown", "unsupported command", fixtures.exitError), + want: errs.Unknown, + }, + { + name: "plain error remains unknown", + err: errors.New("anything"), + want: errs.Unknown, + }, + { + name: "nil remains unknown", + err: nil, + want: errs.Unknown, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, tt.want, Classifier.Classify(tt.err)) + }) + } +} + +func TestClassifier_AppliedViaProcessor(t *testing.T) { + fixtures := setupClassifierFixtures(t) + tests := []struct { + name string + err error + wantRetryable bool + wantDependency bool + wantSame bool + }{ + { + name: "wrapped fetch failure is retryable dependency", + err: fmt.Errorf("reset checkout: %w", gitexec.NewCommandError("fetch", "connection reset", fixtures.exitError)), + wantRetryable: true, + wantDependency: true, + }, + { + name: "wrapped cherry-pick process failure is retryable locally", + err: fmt.Errorf("apply change: %w", gitexec.NewCommandError("cherry-pick", "process killed", fixtures.exitError)), + wantRetryable: true, + }, + { + name: "configuration failure stays non-retryable", + err: fmt.Errorf("prepare checkout: %w", gitexec.NewCommandError("config", "invalid key", fixtures.exitError)), + wantSame: true, + }, + { + name: "authentication failure stays non-retryable dependency", + err: fmt.Errorf("fetch target: %w", gitexec.NewCommandError("fetch", "fatal: Authentication failed", fixtures.exitError)), + wantDependency: true, + wantSame: false, + }, + { + name: "unknown error stays non-retryable", + err: errors.New("anything"), + wantSame: true, + }, + } + + processor := errs.NewClassifierProcessor(Classifier) + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + got := processor.Process(tt.err) + assert.Equal(t, tt.wantRetryable, errs.IsRetryable(got)) + assert.Equal(t, tt.wantDependency, errs.IsDependencyError(got)) + if tt.wantSame { + assert.Same(t, tt.err, got) + } + }) + } +} diff --git a/platform/git/exec/BUILD.bazel b/platform/git/exec/BUILD.bazel index f5c575674..a834acc9b 100644 --- a/platform/git/exec/BUILD.bazel +++ b/platform/git/exec/BUILD.bazel @@ -2,7 +2,10 @@ load("@rules_go//go:def.bzl", "go_library", "go_test") go_library( name = "go_default_library", - srcs = ["gitexec.go"], + srcs = [ + "command_error.go", + "gitexec.go", + ], importpath = "github.com/uber/submitqueue/platform/git/exec", visibility = ["//visibility:public"], ) diff --git a/platform/git/exec/command_error.go b/platform/git/exec/command_error.go new file mode 100644 index 000000000..c0d2709fc --- /dev/null +++ b/platform/git/exec/command_error.go @@ -0,0 +1,74 @@ +// Copyright (c) 2026 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 gitexec + +import "os/exec" + +// CommandError preserves the failed Git operation and its process error for +// backend-specific classification after callers add contextual wrapping. +type CommandError struct { + operation string + message string + cause error +} + +// NewCommandError records a failed Git operation without assigning retry +// policy. Callers supply the rendered diagnostic they want Error to expose. +func NewCommandError(operation, message string, cause error) *CommandError { + return &CommandError{ + operation: operation, + message: message, + cause: cause, + } +} + +// Error returns the command diagnostic supplied by the execution boundary. +func (e *CommandError) Error() string { + if e.message != "" { + return e.message + } + if e.cause != nil { + return e.cause.Error() + } + return "git command failed" +} + +// Unwrap returns the process error reported by os/exec. +func (e *CommandError) Unwrap() error { + return e.cause +} + +// Operation returns the Git subcommand, such as fetch or cherry-pick. +func (e *CommandError) Operation() string { + return e.operation +} + +// Diagnostic returns Git's rendered failure output. +func (e *CommandError) Diagnostic() string { + return e.message +} + +// ProcessExited reports whether Git started and returned a non-zero exit. +func (e *CommandError) ProcessExited() bool { + _, ok := e.cause.(*exec.ExitError) + return ok +} + +func commandOperation(args []string) string { + if len(args) == 0 { + return "" + } + return args[0] +} diff --git a/platform/git/exec/gitexec.go b/platform/git/exec/gitexec.go index 0f7a6d8ad..706be1f05 100644 --- a/platform/git/exec/gitexec.go +++ b/platform/git/exec/gitexec.go @@ -161,7 +161,7 @@ func Output(ctx context.Context, git, dir string, args ...string) (string, error if message == "" { message = err.Error() } - return "", fmt.Errorf("git %s: %s", strings.Join(args, " "), message) + return "", fmt.Errorf("git %s: %w", strings.Join(args, " "), NewCommandError(commandOperation(args), message, err)) } return strings.TrimSpace(string(out)), nil } diff --git a/platform/git/exec/gitexec_test.go b/platform/git/exec/gitexec_test.go index 1658b06d9..60af792c5 100644 --- a/platform/git/exec/gitexec_test.go +++ b/platform/git/exec/gitexec_test.go @@ -15,7 +15,10 @@ package gitexec import ( + "context" + "errors" "os" + "os/exec" "strings" "testing" @@ -95,3 +98,129 @@ func TestEnv_PassthroughDeduplicatesWithTransport(t *testing.T) { func TestEnv_HomeNotInSharedTransportList(t *testing.T) { assert.NotContains(t, transportEnvNames, "HOME") } + +func TestCommandError(t *testing.T) { + cause := errors.New("exit status 128") + tests := []struct { + name string + err *CommandError + wantMessage string + wantCause error + }{ + { + name: "supplied diagnostic is rendered", + err: NewCommandError("fetch", "connection reset", cause), + wantMessage: "connection reset", + wantCause: cause, + }, + { + name: "cause is rendered when diagnostic is empty", + err: NewCommandError("reset", "", cause), + wantMessage: cause.Error(), + wantCause: cause, + }, + { + name: "fallback is rendered without diagnostic or cause", + err: NewCommandError("unknown", "", nil), + wantMessage: "git command failed", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, tt.wantMessage, tt.err.Error()) + assert.Equal(t, tt.err.message, tt.err.Diagnostic()) + assert.Equal(t, tt.err.operation, tt.err.Operation()) + if tt.wantCause == nil { + assert.NoError(t, tt.err.Unwrap()) + } else { + assert.ErrorIs(t, tt.err, tt.wantCause) + } + assert.False(t, tt.err.ProcessExited()) + }) + } +} + +func TestCommandError_ProcessExited(t *testing.T) { + err := exec.Command(os.Args[0], "-test.run=[").Run() + require.Error(t, err) + var exitErr *exec.ExitError + require.ErrorAs(t, err, &exitErr) + + tests := []struct { + name string + err *CommandError + want bool + }{ + { + name: "non-zero process exit is reported", + err: NewCommandError("fetch", "failed", exitErr), + want: true, + }, + { + name: "ordinary cause did not exit a process", + err: NewCommandError("fetch", "failed", errors.New("start failure")), + want: false, + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, tt.want, tt.err.ProcessExited()) + }) + } +} + +func TestOutput_PreservesCommandFailure(t *testing.T) { + tests := []struct { + name string + executable string + args []string + wantOperation string + }{ + { + name: "non-zero process exit retains command provenance and cause", + executable: os.Args[0], + args: []string{"-test.run=["}, + wantOperation: "-test.run=[", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + _, err := Output(context.Background(), tt.executable, "", tt.args...) + require.Error(t, err) + + var commandErr *CommandError + require.ErrorAs(t, err, &commandErr) + assert.Equal(t, tt.wantOperation, commandErr.Operation()) + + var exitErr *exec.ExitError + assert.ErrorAs(t, err, &exitErr) + }) + } +} + +func TestCommandOperation(t *testing.T) { + tests := []struct { + name string + args []string + want string + }{ + { + name: "first argument is the operation", + args: []string{"fetch", "origin"}, + want: "fetch", + }, + { + name: "empty arguments have no operation", + want: "", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + assert.Equal(t, tt.want, commandOperation(tt.args)) + }) + } +} diff --git a/runway/extension/merger/git/BUILD.bazel b/runway/extension/merger/git/BUILD.bazel index 9dbc3feb6..fcaf7ce67 100644 --- a/runway/extension/merger/git/BUILD.bazel +++ b/runway/extension/merger/git/BUILD.bazel @@ -49,6 +49,7 @@ go_test( "//api/base/mergestrategy/protopb:go_default_library", "//api/runway/messagequeue:go_default_library", "//api/runway/messagequeue/protopb:go_default_library", + "//platform/git/exec:go_default_library", "//platform/git/exectest:go_default_library", "//runway/extension/merger:go_default_library", "@com_github_stretchr_testify//assert:go_default_library", diff --git a/runway/extension/merger/git/git_merger.go b/runway/extension/merger/git/git_merger.go index f0d6575f6..8793982af 100644 --- a/runway/extension/merger/git/git_merger.go +++ b/runway/extension/merger/git/git_merger.go @@ -694,7 +694,7 @@ func (m *gitMerger) applyMerge(ctx context.Context, rs resolvedStep) (applied, e // does not establish that anything collided. conflicted := m.hasUnmergedPaths(ctx) _, _ = m.run(ctx, nil, "merge", "--abort") - return applied{}, m.classifyMergeFailure(ref, o, conflicted) + return applied{}, m.classifyMergeFailure(ref, o, conflicted, err) } mergeSHA, err := m.headSHA(ctx) if err != nil { @@ -785,7 +785,7 @@ func (m *gitMerger) promote(ctx context.Context, req *runwaymq.MergeRequest, rs // and fixed by configuration rather than by rebasing, and any other way git // can exit non-zero — a missing object, an unreadable repository, a killed // process — which is infrastructure and should be retried, not made terminal. -func (m *gitMerger) classifyMergeFailure(ref changeRef, out []byte, conflicted bool) error { +func (m *gitMerger) classifyMergeFailure(ref changeRef, out []byte, conflicted bool, cause error) error { detail := strings.TrimSpace(string(out)) if strings.Contains(detail, "refusing to merge unrelated histories") { coremetrics.NamedCounter(m.metricsScope, "merge", "unrelated_histories", 1) @@ -794,7 +794,7 @@ func (m *gitMerger) classifyMergeFailure(ref changeRef, out []byte, conflicted b } if !conflicted { coremetrics.NamedCounter(m.metricsScope, "merge", "merge_errors", 1) - return fmt.Errorf("git merge %s: %s", ref.SHA, detail) + return fmt.Errorf("git merge %s: %w", ref.SHA, cause) } coremetrics.NamedCounter(m.metricsScope, "merge", "merge_conflicts", 1) return fmt.Errorf("%w: git merge %s: %s", merger.ErrConflict, ref.SHA, detail) @@ -892,7 +892,7 @@ func (m *gitMerger) cherryPickRange(ctx context.Context, base, head string) erro detail := strings.TrimSpace(string(out)) if !conflicted { coremetrics.NamedCounter(m.metricsScope, "merge", "cherry_pick_errors", 1) - return fmt.Errorf("git cherry-pick %s..%s: %w: %s", base, head, err, detail) + return fmt.Errorf("git cherry-pick %s..%s: %w", base, head, err) } coremetrics.NamedCounter(m.metricsScope, "merge", "cherry_pick_conflicts", 1) return fmt.Errorf("%w: git cherry-pick %s..%s: %s", merger.ErrConflict, base, head, detail) @@ -1000,7 +1000,12 @@ func (m *gitMerger) isAncestor(ctx context.Context, ancestor, descendant string) if errors.As(err, &exitErr) && exitErr.ExitCode() == 1 { return false, nil } - return false, fmt.Errorf("git merge-base --is-ancestor %s %s: %w: %s", ancestor, descendant, err, strings.TrimSpace(stderr.String())) + message := err.Error() + if detail := strings.TrimSpace(stderr.String()); detail != "" { + message += ": " + detail + } + return false, fmt.Errorf("git merge-base --is-ancestor %s %s: %w", + ancestor, descendant, gitexec.NewCommandError("merge-base", message, err)) } // commitTreeSHA returns the tree SHA recorded in the commit object at ref. @@ -1038,7 +1043,11 @@ func (m *gitMerger) runAs(ctx context.Context, author authorIdent, stdin []byte, cmd.Stdout = &stdout cmd.Stderr = &stderr if err := cmd.Run(); err != nil { - return nil, fmt.Errorf("%w: %s", err, strings.TrimSpace(stderr.String())) + message := err.Error() + if detail := strings.TrimSpace(stderr.String()); detail != "" { + message += ": " + detail + } + return nil, gitexec.NewCommandError(args[0], message, err) } return stdout.Bytes(), nil } @@ -1056,7 +1065,15 @@ func (m *gitMerger) runCombinedAs(ctx context.Context, author authorIdent, stdin if stdin != nil { cmd.Stdin = bytes.NewReader(stdin) } - return cmd.CombinedOutput() + out, err := cmd.CombinedOutput() + if err != nil { + message := err.Error() + if detail := strings.TrimSpace(string(out)); detail != "" { + message += ": " + detail + } + return out, gitexec.NewCommandError(args[0], message, err) + } + return out, nil } // command builds a git command with the committer identity injected via -c diff --git a/runway/extension/merger/git/git_merger_test.go b/runway/extension/merger/git/git_merger_test.go index b1e6e2224..05a84edf7 100644 --- a/runway/extension/merger/git/git_merger_test.go +++ b/runway/extension/merger/git/git_merger_test.go @@ -35,6 +35,7 @@ import ( mergestrategypb "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" + gitexec "github.com/uber/submitqueue/platform/git/exec" gitexectest "github.com/uber/submitqueue/platform/git/exectest" "github.com/uber/submitqueue/runway/extension/merger" ) @@ -210,6 +211,8 @@ func TestCherryPickRange_NonConflictFailureIsRetryable(t *testing.T) { require.Error(t, err) assert.False(t, errors.Is(err, merger.ErrConflict), "a non-conflict failure must stay retryable") assert.False(t, errors.Is(err, merger.ErrInvalidRequest)) + var commandErr *gitexec.CommandError + assert.ErrorAs(t, err, &commandErr) } func TestCherryPickRange_RealConflictIsErrConflict(t *testing.T) { @@ -623,10 +626,14 @@ func TestClassifyMergeFailure(t *testing.T) { for _, tt := range tests { t.Run(tt.name, func(t *testing.T) { - err := m.classifyMergeFailure(ref, []byte(tt.out), tt.conflicted) + cause := errors.New("git exited") + err := m.classifyMergeFailure(ref, []byte(tt.out), tt.conflicted, cause) require.Error(t, err) assert.Equal(t, tt.wantConflict, errors.Is(err, merger.ErrConflict)) assert.Equal(t, tt.wantInvalid, errors.Is(err, merger.ErrInvalidRequest)) + if !tt.wantConflict && !tt.wantInvalid { + assert.ErrorIs(t, err, cause) + } }) } } diff --git a/service/runway/server/BUILD.bazel b/service/runway/server/BUILD.bazel index 113cd6e23..85fa29808 100644 --- a/service/runway/server/BUILD.bazel +++ b/service/runway/server/BUILD.bazel @@ -22,6 +22,7 @@ go_library( "//platform/consumer:go_default_library", "//platform/errs:go_default_library", "//platform/errs/generic:go_default_library", + "//platform/errs/git:go_default_library", "//platform/errs/mysql:go_default_library", "//platform/extension/consumergate:go_default_library", "//platform/extension/consumergate/file:go_default_library", @@ -79,6 +80,7 @@ go_test( srcs = [ "checkout_test.go", "config_test.go", + "main_test.go", ], # Checkout provisioning runs real git, so the test uses the same pinned # runtime the merger does rather than whatever git the host happens to have. @@ -97,10 +99,19 @@ go_test( }, deps = [ "//api/base/mergestrategy/protopb:go_default_library", + "//platform/base/failure:go_default_library", + "//platform/base/messagequeue:go_default_library", + "//platform/consumer:go_default_library", + "//platform/extension/consumergate/noop:go_default_library", + "//platform/extension/messagequeue:go_default_library", + "//platform/extension/messagequeue/mock:go_default_library", + "//platform/git/exec:go_default_library", "//platform/git/exectest:go_default_library", "//runway/extension/merger/git:go_default_library", "@com_github_stretchr_testify//assert:go_default_library", "@com_github_stretchr_testify//require:go_default_library", + "@com_github_uber_go_tally//:go_default_library", + "@org_uber_go_mock//gomock:go_default_library", "@org_uber_go_zap//zaptest:go_default_library", ], ) diff --git a/service/runway/server/main.go b/service/runway/server/main.go index dc9985001..39d1ac95c 100644 --- a/service/runway/server/main.go +++ b/service/runway/server/main.go @@ -37,6 +37,7 @@ import ( "github.com/uber/submitqueue/platform/consumer" "github.com/uber/submitqueue/platform/errs" genericerrs "github.com/uber/submitqueue/platform/errs/generic" + giterrs "github.com/uber/submitqueue/platform/errs/git" mysqlerrs "github.com/uber/submitqueue/platform/errs/mysql" "github.com/uber/submitqueue/platform/extension/consumergate" consumergatefile "github.com/uber/submitqueue/platform/extension/consumergate/file" @@ -163,13 +164,7 @@ func run() error { // group name just like a primary stage. gate := newConsumerGate(logger) - primaryConsumer := consumer.New(logger.Sugar(), scope.SubScope("consumer"), registry, - errs.NewClassifierProcessor( - genericerrs.Classifier, - mysqlerrs.Classifier, - ), - gate, - ) + primaryConsumer := consumer.New(logger.Sugar(), scope.SubScope("consumer"), registry, newPrimaryErrorProcessor(), gate) mergerFactory, err := newMergerFactory(ctx, logger, scope.SubScope("merger")) if err != nil { @@ -305,6 +300,14 @@ func run() error { return err } +func newPrimaryErrorProcessor() errs.ErrorProcessor { + return errs.NewClassifierProcessor( + genericerrs.Classifier, + giterrs.Classifier, + mysqlerrs.Classifier, + ) +} + // newMergerFactory builds the mergers for the server. // // MERGER pins every queue to one implementation explicitly, which is how a test diff --git a/service/runway/server/main_test.go b/service/runway/server/main_test.go new file mode 100644 index 000000000..e46b8b54b --- /dev/null +++ b/service/runway/server/main_test.go @@ -0,0 +1,142 @@ +// Copyright (c) 2026 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 main + +import ( + "context" + "errors" + "os" + "os/exec" + "testing" + + "github.com/stretchr/testify/require" + "github.com/uber-go/tally" + "github.com/uber/submitqueue/platform/base/failure" + entityqueue "github.com/uber/submitqueue/platform/base/messagequeue" + "github.com/uber/submitqueue/platform/consumer" + consumergatenoop "github.com/uber/submitqueue/platform/extension/consumergate/noop" + extqueue "github.com/uber/submitqueue/platform/extension/messagequeue" + queuemock "github.com/uber/submitqueue/platform/extension/messagequeue/mock" + gitexec "github.com/uber/submitqueue/platform/git/exec" + "go.uber.org/mock/gomock" + "go.uber.org/zap/zaptest" +) + +const ( + testGitTopicKey consumer.TopicKey = "git-test" + testGitGroup = "git-test-group" +) + +type errorController struct { + err error +} + +func (c errorController) Process(context.Context, consumer.Delivery) error { + return c.err +} + +func (errorController) Name() string { + return "git-test" +} + +func (errorController) TopicKey() consumer.TopicKey { + return testGitTopicKey +} + +func (errorController) ConsumerGroup() string { + return testGitGroup +} + +func gitExitError(t *testing.T) error { + t.Helper() + err := exec.Command(os.Args[0], "-test.run=[").Run() + require.Error(t, err) + var exitErr *exec.ExitError + require.ErrorAs(t, err, &exitErr) + return exitErr +} + +func TestPrimaryConsumer_GitFailureDisposition(t *testing.T) { + exitErr := gitExitError(t) + tests := []struct { + name string + controller error + wantOutcome string + }{ + { + name: "temporary remote fetch failure is nacked for retry", + controller: gitexec.NewCommandError("fetch", "remote temporarily unavailable", exitErr), + wantOutcome: "nack", + }, + { + name: "unknown error is rejected to dead letter", + controller: errors.New("unknown failure"), + wantOutcome: "reject", + }, + } + + for _, tt := range tests { + t.Run(tt.name, func(t *testing.T) { + ctrl := gomock.NewController(t) + deliveryChannel := make(chan extqueue.Delivery, 1) + subscriber := queuemock.NewMockSubscriber(ctrl) + subscriber.EXPECT().Subscribe(gomock.Any(), gomock.Any(), gomock.Any()).Return(deliveryChannel, nil) + queue := queuemock.NewMockQueue(ctrl) + queue.EXPECT().Subscriber().Return(subscriber) + + registry, err := consumer.NewTopicRegistry([]consumer.TopicConfig{{ + Key: testGitTopicKey, + Name: "git-test", + Queue: queue, + Subscription: extqueue.DefaultSubscriptionConfig( + "git-test-worker", + testGitGroup, + ), + }}) + require.NoError(t, err) + + serviceConsumer := consumer.New( + zaptest.NewLogger(t).Sugar(), + tally.NoopScope, + registry, + newPrimaryErrorProcessor(), + consumergatenoop.New(), + ) + require.NoError(t, serviceConsumer.Register(errorController{err: tt.controller})) + require.NoError(t, serviceConsumer.Start(context.Background())) + + message := entityqueue.NewMessage("git-test-message", []byte("payload"), "partition", nil) + delivery := queuemock.NewMockDelivery(ctrl) + delivery.EXPECT().Message().Return(message).AnyTimes() + delivery.EXPECT().Attempt().Return(1).AnyTimes() + done := make(chan struct{}) + if tt.wantOutcome == "nack" { + delivery.EXPECT().Nack(gomock.Any(), gomock.Any()).DoAndReturn(func(context.Context, failure.Failure) error { + close(done) + return nil + }) + } else { + delivery.EXPECT().Reject(gomock.Any(), gomock.Any()).DoAndReturn(func(context.Context, failure.Failure) error { + close(done) + return nil + }) + } + + deliveryChannel <- delivery + <-done + require.NoError(t, serviceConsumer.Stop(30000)) + }) + } +}