Skip to content
Draft
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
6 changes: 6 additions & 0 deletions service/stovepipe/server/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,9 @@ go_library(
"@com_github_go_sql_driver_mysql//:go_default_library",
"@com_github_uber_go_tally//:go_default_library",
"@org_golang_google_grpc//:go_default_library",
"@org_golang_google_grpc//codes:go_default_library",
"@org_golang_google_grpc//reflection:go_default_library",
"@org_golang_google_grpc//status:go_default_library",
"@org_uber_go_zap//:go_default_library",
],
)
Expand Down Expand Up @@ -78,10 +80,14 @@ go_test(
deps = [
"//api/base/hook:go_default_library",
"//platform/consumer:go_default_library",
"//platform/errs:go_default_library",
"//stovepipe/controller:go_default_library",
"//stovepipe/controller/dlq: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_golang_google_grpc//codes:go_default_library",
"@org_golang_google_grpc//status:go_default_library",
"@org_uber_go_zap//zaptest:go_default_library",
],
)
53 changes: 47 additions & 6 deletions service/stovepipe/server/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -59,14 +59,17 @@ import (
storageMySQL "github.com/uber/submitqueue/stovepipe/extension/storage/mysql"
"go.uber.org/zap"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/reflection"
"google.golang.org/grpc/status"
)

// StovepipeServer wraps the controllers and implements the gRPC service interface.
type StovepipeServer struct {
pb.UnimplementedStovepipeServer
pingController *controller.PingController
ingestController *controller.IngestController
pingController *controller.PingController
ingestController *controller.IngestController
getProjectStatusByURIController *controller.GetProjectStatusByURIController
}

// Ping delegates to the controller.
Expand All @@ -84,6 +87,34 @@ func (s *StovepipeServer) Ingest(ctx context.Context, req *pb.IngestRequest) (*p
return mapper.IngestResultToProto(result), nil
}

// GetProjectStatusByURI returns current repository validation for an exact commit URI.
func (s *StovepipeServer) GetProjectStatusByURI(ctx context.Context, req *pb.GetProjectStatusByURIRequest) (*pb.GetProjectStatusByURIResponse, error) {
result, err := s.getProjectStatusByURIController.GetProjectStatusByURI(ctx, mapper.ProtoToGetProjectStatusByURIRequest(req))
if err != nil {
return nil, err
}
return mapper.GetProjectStatusByURIResultToProto(result), nil
}

func stovepipeStatusError(err error) error {
switch {
case errors.Is(err, context.Canceled):
return status.Error(codes.Canceled, err.Error())
case errors.Is(err, context.DeadlineExceeded):
return status.Error(codes.DeadlineExceeded, err.Error())
case controller.IsProjectStatusNotFound(err):
return status.Error(codes.NotFound, err.Error())
case controller.IsProjectStatusConsistency(err):
return status.Error(codes.Internal, err.Error())
case controller.IsInvalidRequest(err):
return status.Error(codes.InvalidArgument, err.Error())
case errs.IsRetryable(err):
return status.Error(codes.Unavailable, err.Error())
default:
return err
}
}

// inMemoryCounter is a minimal, process-local counter.Counter used to wire the example
// server. It is not durable; a real deployment supplies a persistent implementation
// (e.g. platform/extension/counter/mysql).
Expand Down Expand Up @@ -323,8 +354,16 @@ func run() error {
}
logger.Info("consumers started")

// Create gRPC server
grpcServer := grpc.NewServer()
// Create gRPC server with stable transport codes for controller outcomes.
grpcServer := grpc.NewServer(grpc.UnaryInterceptor(
func(ctx context.Context, req interface{}, info *grpc.UnaryServerInfo, handler grpc.UnaryHandler) (interface{}, error) {
resp, err := handler(ctx, req)
if err != nil {
return nil, stovepipeStatusError(err)
}
return resp, nil
},
))

// Create controllers and wrap them for gRPC
pingController := controller.NewPingController(logger, scope)
Expand All @@ -336,9 +375,11 @@ func run() error {
storageFty,
registry,
)
getProjectStatusByURIController := controller.NewGetProjectStatusByURIController(logger.Sugar(), scope, storageFty)
srv := &StovepipeServer{
pingController: pingController,
ingestController: ingestController,
pingController: pingController,
ingestController: ingestController,
getProjectStatusByURIController: getProjectStatusByURIController,
}
pb.RegisterStovepipeServer(grpcServer, srv)

Expand Down
29 changes: 29 additions & 0 deletions service/stovepipe/server/main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ package main

import (
"context"
"errors"
"strings"
"testing"

Expand All @@ -24,10 +25,38 @@ import (
"github.com/uber-go/tally"
basehook "github.com/uber/submitqueue/api/base/hook"
"github.com/uber/submitqueue/platform/consumer"
"github.com/uber/submitqueue/platform/errs"
"github.com/uber/submitqueue/stovepipe/controller"
"github.com/uber/submitqueue/stovepipe/controller/dlq"
"go.uber.org/zap/zaptest"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/status"
)

func TestStovepipeStatusError(t *testing.T) {
tests := []struct {
name string
err error
code codes.Code
}{
{name: "invalid request", err: controller.ErrInvalidRequest, code: codes.InvalidArgument},
{name: "not found", err: &controller.ProjectStatusNotFoundError{Queue: "q", ChangeURI: "uri"}, code: codes.NotFound},
{name: "inconsistent records", err: &controller.ProjectStatusConsistencyError{Message: "inconsistent"}, code: codes.Internal},
{name: "retryable", err: errs.NewRetryableError(errors.New("try again")), code: codes.Unavailable},
{name: "canceled", err: errors.Join(errs.NewRetryableError(errors.New("try again")), context.Canceled), code: codes.Canceled},
{name: "deadline exceeded", err: context.DeadlineExceeded, code: codes.DeadlineExceeded},
}

for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
assert.Equal(t, tt.code, status.Code(stovepipeStatusError(tt.err)))
})
}

infrastructureErr := errors.New("storage unavailable")
assert.Equal(t, infrastructureErr, stovepipeStatusError(infrastructureErr))
}

// recordingConsumer captures what the host registers instead of subscribing.
type recordingConsumer struct {
controllers []consumer.Controller
Expand Down
10 changes: 8 additions & 2 deletions service/stovepipe/server/mapper/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,10 @@ load("@rules_go//go:def.bzl", "go_library", "go_test")

go_library(
name = "go_default_library",
srcs = ["ingest.go"],
srcs = [
"get_project_status_by_uri.go",
"ingest.go",
],
importpath = "github.com/uber/submitqueue/service/stovepipe/server/mapper",
visibility = ["//visibility:public"],
deps = [
Expand All @@ -13,7 +16,10 @@ go_library(

go_test(
name = "go_default_test",
srcs = ["ingest_test.go"],
srcs = [
"get_project_status_by_uri_test.go",
"ingest_test.go",
],
embed = [":go_default_library"],
deps = [
"//api/stovepipe/protopb:go_default_library",
Expand Down
60 changes: 60 additions & 0 deletions service/stovepipe/server/mapper/get_project_status_by_uri.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
// 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 mapper

import (
pb "github.com/uber/submitqueue/api/stovepipe/protopb"
"github.com/uber/submitqueue/stovepipe/entity"
)

// ProtoToGetProjectStatusByURIRequest maps the wire selector to its domain form.
func ProtoToGetProjectStatusByURIRequest(req *pb.GetProjectStatusByURIRequest) entity.GetProjectStatusByURIRequest {
result := entity.GetProjectStatusByURIRequest{
Queue: req.GetQueue(),
ChangeURI: req.GetChangeUri(),
PageSize: req.GetPageSize(),
PageToken: req.GetPageToken(),
}
if req.Project != nil {
result.Project = req.GetProject()
result.HasProject = true
}
return result
}

// GetProjectStatusByURIResultToProto maps a domain projection to the wire response.
func GetProjectStatusByURIResultToProto(result entity.GetProjectStatusByURIResult) *pb.GetProjectStatusByURIResponse {
response := &pb.GetProjectStatusByURIResponse{
RequestId: result.RequestID,
Queue: result.Queue,
ChangeUri: result.ChangeURI,
BaseUri: result.BaseURI,
RequestState: string(result.RequestState),
ProjectResultsComplete: result.ProjectResultsComplete,
Projects: make([]*pb.ProjectValidation, 0, len(result.Projects)),
NextPageToken: result.NextPageToken,
}
if result.HasRepositoryBreakageDegree {
response.RepositoryBreakageDegree = &result.RepositoryBreakageDegree
}
for _, project := range result.Projects {
mapped := &pb.ProjectValidation{Project: project.Project}
if project.HasBreakageDegree {
mapped.BreakageDegree = &project.BreakageDegree
}
response.Projects = append(response.Projects, mapped)
}
return response
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,67 @@
// 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 mapper

import (
"testing"

"github.com/stretchr/testify/assert"
pb "github.com/uber/submitqueue/api/stovepipe/protopb"
"github.com/uber/submitqueue/stovepipe/entity"
)

func TestProtoToGetProjectStatusByURIRequest(t *testing.T) {
project := ""
got := ProtoToGetProjectStatusByURIRequest(&pb.GetProjectStatusByURIRequest{
Queue: "monorepo/main", ChangeUri: "git://commit", Project: &project, PageSize: 10, PageToken: "token",
})

assert.Equal(t, entity.GetProjectStatusByURIRequest{
Queue: "monorepo/main", ChangeURI: "git://commit", Project: "", HasProject: true, PageSize: 10, PageToken: "token",
}, got)

omitted := ProtoToGetProjectStatusByURIRequest(&pb.GetProjectStatusByURIRequest{})
assert.False(t, omitted.HasProject)
}

func TestGetProjectStatusByURIResultToProto(t *testing.T) {
t.Run("preserves optional field presence", func(t *testing.T) {
result := entity.GetProjectStatusByURIResult{
RequestID: "request/monorepo/main/7", Queue: "monorepo/main", ChangeURI: "git://commit",
BaseURI: "git://base", RequestState: entity.ProjectStatusRequestStateSucceeded,
RepositoryBreakageDegree: entity.DegreeGreen, HasRepositoryBreakageDegree: true,
Projects: []entity.ProjectValidation{
{Project: "//green", BreakageDegree: entity.DegreeGreen, HasBreakageDegree: true},
{Project: "//pending"},
},
}

got := GetProjectStatusByURIResultToProto(result)

assert.NotNil(t, got.RepositoryBreakageDegree)
assert.Equal(t, entity.DegreeGreen, got.GetRepositoryBreakageDegree())
assert.Equal(t, "succeeded", got.GetRequestState())
assert.Len(t, got.Projects, 2)
assert.NotNil(t, got.Projects[0].BreakageDegree)
assert.Nil(t, got.Projects[1].BreakageDegree)
})

t.Run("keeps missing repository fact absent", func(t *testing.T) {
got := GetProjectStatusByURIResultToProto(entity.GetProjectStatusByURIResult{})

assert.Nil(t, got.RepositoryBreakageDegree)
assert.Empty(t, got.Projects)
})
}
3 changes: 3 additions & 0 deletions stovepipe/controller/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ load("@rules_go//go:def.bzl", "go_library", "go_test")
go_library(
name = "go_default_library",
srcs = [
"get_project_status_by_uri.go",
"ingest.go",
"ping.go",
],
Expand All @@ -27,13 +28,15 @@ go_library(
go_test(
name = "go_default_test",
srcs = [
"get_project_status_by_uri_test.go",
"ingest_test.go",
"ping_test.go",
],
embed = [":go_default_library"],
deps = [
"//api/stovepipe/protopb:go_default_library",
"//platform/consumer:go_default_library",
"//platform/errs:go_default_library",
"//platform/extension/counter:go_default_library",
"//platform/extension/counter/mock:go_default_library",
"//platform/extension/messagequeue/mock:go_default_library",
Expand Down
Loading