Skip to content

Commit 74a2a52

Browse files
authored
refactor(stovepipe): Move RPC handlers (#729)
## Summary **What**: - Move gRPC request handling into a dedicated server component. - Keep process startup focused on wiring and serving. **Why**: - Keep the service entrypoint manageable as endpoints grow. - Preserve existing behavior while improving separation. ## Test Plan - [x] Run unit tests. ## Revert Plan - Revert this PR to move the RPC handlers and their tests back into the Stovepipe server entrypoint; request behavior remains unchanged. ## Issues - [CODEM-449](https://linear.app/uber/issue/CODEM-449)
1 parent 08e5ffd commit 74a2a52

6 files changed

Lines changed: 246 additions & 167 deletions

File tree

‎service/stovepipe/server/BUILD.bazel‎

Lines changed: 1 addition & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -22,7 +22,7 @@ go_library(
2222
"//platform/extension/messagequeue/mysql:go_default_library",
2323
"//platform/hook:go_default_library",
2424
"//service/messagequeue:go_default_library",
25-
"//service/stovepipe/server/mapper:go_default_library",
25+
"//service/stovepipe/server/handler:go_default_library",
2626
"//stovepipe/controller:go_default_library",
2727
"//stovepipe/controller/build:go_default_library",
2828
"//stovepipe/controller/buildsignal:go_default_library",
@@ -80,12 +80,9 @@ go_test(
8080
embed = [":go_default_library"],
8181
deps = [
8282
"//api/base/hook:go_default_library",
83-
"//api/stovepipe/protopb:go_default_library",
8483
"//platform/consumer:go_default_library",
85-
"//stovepipe/controller:go_default_library",
8684
"//stovepipe/controller/dlq:go_default_library",
8785
"//stovepipe/core/requestlog:go_default_library",
88-
"//stovepipe/entity:go_default_library",
8986
"@com_github_stretchr_testify//assert:go_default_library",
9087
"@com_github_stretchr_testify//require:go_default_library",
9188
"@com_github_uber_go_tally//:go_default_library",
Lines changed: 26 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,26 @@
1+
load("@rules_go//go:def.bzl", "go_library", "go_test")
2+
3+
go_library(
4+
name = "go_default_library",
5+
srcs = ["server.go"],
6+
importpath = "github.com/uber/submitqueue/service/stovepipe/server/handler",
7+
visibility = ["//visibility:public"],
8+
deps = [
9+
"//api/stovepipe/protopb:go_default_library",
10+
"//service/stovepipe/server/mapper:go_default_library",
11+
"//stovepipe/controller:go_default_library",
12+
],
13+
)
14+
15+
go_test(
16+
name = "go_default_test",
17+
srcs = ["server_test.go"],
18+
embed = [":go_default_library"],
19+
deps = [
20+
"//api/stovepipe/protopb:go_default_library",
21+
"//stovepipe/controller:go_default_library",
22+
"//stovepipe/entity:go_default_library",
23+
"@com_github_stretchr_testify//assert:go_default_library",
24+
"@com_github_stretchr_testify//require:go_default_library",
25+
],
26+
)
Lines changed: 78 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,78 @@
1+
// Copyright (c) 2026 Uber Technologies, Inc.
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
// Package handler implements Stovepipe's gRPC service handlers.
16+
package handler
17+
18+
import (
19+
"context"
20+
21+
pb "github.com/uber/submitqueue/api/stovepipe/protopb"
22+
"github.com/uber/submitqueue/service/stovepipe/server/mapper"
23+
"github.com/uber/submitqueue/stovepipe/controller"
24+
)
25+
26+
// StovepipeServer wraps the controllers and implements the gRPC service interface.
27+
type StovepipeServer struct {
28+
pb.UnimplementedStovepipeServer
29+
pingController *controller.PingController
30+
ingestController *controller.IngestController
31+
requestHistoryController controller.RequestHistoryController
32+
}
33+
34+
// NewStovepipeServer creates a gRPC service handler from Stovepipe controllers.
35+
func NewStovepipeServer(
36+
pingController *controller.PingController,
37+
ingestController *controller.IngestController,
38+
requestHistoryController controller.RequestHistoryController,
39+
) *StovepipeServer {
40+
return &StovepipeServer{
41+
pingController: pingController,
42+
ingestController: ingestController,
43+
requestHistoryController: requestHistoryController,
44+
}
45+
}
46+
47+
// Ping delegates to the controller.
48+
func (s *StovepipeServer) Ping(ctx context.Context, req *pb.PingRequest) (*pb.PingResponse, error) {
49+
return s.pingController.Ping(ctx, req)
50+
}
51+
52+
// Ingest maps the wire request to an entity, delegates to the controller, and maps
53+
// the result back to the wire response.
54+
func (s *StovepipeServer) Ingest(ctx context.Context, req *pb.IngestRequest) (*pb.IngestResponse, error) {
55+
result, err := s.ingestController.Ingest(ctx, mapper.ProtoToIngestRequest(req))
56+
if err != nil {
57+
return nil, err
58+
}
59+
return mapper.IngestResultToProto(result), nil
60+
}
61+
62+
// GetRequestHistoryByID returns retained history for one request ID.
63+
func (s *StovepipeServer) GetRequestHistoryByID(ctx context.Context, req *pb.GetRequestHistoryByIDRequest) (*pb.GetRequestHistoryByIDResponse, error) {
64+
events, err := s.requestHistoryController.GetRequestHistoryByID(ctx, mapper.ProtoToGetRequestHistoryByIDRequest(req))
65+
if err != nil {
66+
return nil, err
67+
}
68+
return &pb.GetRequestHistoryByIDResponse{Events: mapper.HistoryEventsToProto(events)}, nil
69+
}
70+
71+
// GetRequestHistoryByURI returns retained histories for one commit URI.
72+
func (s *StovepipeServer) GetRequestHistoryByURI(ctx context.Context, req *pb.GetRequestHistoryByURIRequest) (*pb.GetRequestHistoryByURIResponse, error) {
73+
histories, err := s.requestHistoryController.GetRequestHistoryByURI(ctx, mapper.ProtoToGetRequestHistoryByURIRequest(req))
74+
if err != nil {
75+
return nil, err
76+
}
77+
return &pb.GetRequestHistoryByURIResponse{Histories: mapper.RequestHistoriesToProto(histories)}, nil
78+
}
Lines changed: 139 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,139 @@
1+
// Copyright (c) 2026 Uber Technologies, Inc.
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package handler
16+
17+
import (
18+
"context"
19+
"errors"
20+
"testing"
21+
22+
"github.com/stretchr/testify/assert"
23+
"github.com/stretchr/testify/require"
24+
pb "github.com/uber/submitqueue/api/stovepipe/protopb"
25+
"github.com/uber/submitqueue/stovepipe/controller"
26+
"github.com/uber/submitqueue/stovepipe/entity"
27+
)
28+
29+
type fakeRequestHistoryController struct {
30+
getByID func(context.Context, entity.GetRequestHistoryByIDRequest) ([]entity.RequestLog, error)
31+
getByURI func(context.Context, entity.GetRequestHistoryByURIRequest) ([]entity.RequestHistory, error)
32+
}
33+
34+
var _ controller.RequestHistoryController = (*fakeRequestHistoryController)(nil)
35+
36+
func (f *fakeRequestHistoryController) GetRequestHistoryByID(ctx context.Context, req entity.GetRequestHistoryByIDRequest) ([]entity.RequestLog, error) {
37+
return f.getByID(ctx, req)
38+
}
39+
40+
func (f *fakeRequestHistoryController) GetRequestHistoryByURI(ctx context.Context, req entity.GetRequestHistoryByURIRequest) ([]entity.RequestHistory, error) {
41+
return f.getByURI(ctx, req)
42+
}
43+
44+
func TestGetRequestHistoryByID(t *testing.T) {
45+
controllerErr := errors.New("controller failed")
46+
logs := []entity.RequestLog{
47+
{ID: "occurrence/1", State: entity.RequestStateAccepted, TimestampMs: 1000},
48+
{ID: "occurrence/2", Event: entity.RequestEventBuildTriggered, TimestampMs: 2000},
49+
}
50+
tests := []struct {
51+
name string
52+
logs []entity.RequestLog
53+
err error
54+
wantLogs int
55+
}{
56+
{name: "maps successful result", logs: logs, wantLogs: 2},
57+
{name: "maps empty result", logs: nil, wantLogs: 0},
58+
{name: "returns controller error unchanged", err: controllerErr},
59+
}
60+
61+
for _, tt := range tests {
62+
t.Run(tt.name, func(t *testing.T) {
63+
var gotReq entity.GetRequestHistoryByIDRequest
64+
fake := &fakeRequestHistoryController{
65+
getByID: func(_ context.Context, req entity.GetRequestHistoryByIDRequest) ([]entity.RequestLog, error) {
66+
gotReq = req
67+
return tt.logs, tt.err
68+
},
69+
}
70+
srv := &StovepipeServer{requestHistoryController: fake}
71+
72+
resp, err := srv.GetRequestHistoryByID(context.Background(), &pb.GetRequestHistoryByIDRequest{
73+
Queue: "monorepo/main", RequestId: "request/1",
74+
})
75+
76+
assert.Equal(t, entity.GetRequestHistoryByIDRequest{Queue: "monorepo/main", ID: "request/1"}, gotReq)
77+
if tt.err != nil {
78+
require.ErrorIs(t, err, tt.err)
79+
assert.Nil(t, resp)
80+
return
81+
}
82+
require.NoError(t, err)
83+
require.Len(t, resp.Events, tt.wantLogs)
84+
if tt.wantLogs > 0 {
85+
assert.Equal(t, "accepted", resp.Events[0].GetRequestState())
86+
assert.Equal(t, "build_triggered", resp.Events[1].GetEvent())
87+
}
88+
})
89+
}
90+
}
91+
92+
func TestGetRequestHistoryByURI(t *testing.T) {
93+
controllerErr := errors.New("controller failed")
94+
histories := []entity.RequestHistory{{
95+
RequestID: "request/1",
96+
Events: []entity.RequestLog{{ID: "occurrence/1", State: entity.RequestStateAccepted}},
97+
}}
98+
tests := []struct {
99+
name string
100+
histories []entity.RequestHistory
101+
err error
102+
wantHistories int
103+
}{
104+
{name: "maps successful result", histories: histories, wantHistories: 1},
105+
{name: "maps empty result", histories: nil, wantHistories: 0},
106+
{name: "returns controller error unchanged", err: controllerErr},
107+
}
108+
109+
for _, tt := range tests {
110+
t.Run(tt.name, func(t *testing.T) {
111+
var gotReq entity.GetRequestHistoryByURIRequest
112+
fake := &fakeRequestHistoryController{
113+
getByURI: func(_ context.Context, req entity.GetRequestHistoryByURIRequest) ([]entity.RequestHistory, error) {
114+
gotReq = req
115+
return tt.histories, tt.err
116+
},
117+
}
118+
srv := &StovepipeServer{requestHistoryController: fake}
119+
120+
resp, err := srv.GetRequestHistoryByURI(context.Background(), &pb.GetRequestHistoryByURIRequest{
121+
Queue: "monorepo/main", Uri: "git://monorepo/abc",
122+
})
123+
124+
assert.Equal(t, entity.GetRequestHistoryByURIRequest{Queue: "monorepo/main", URI: "git://monorepo/abc"}, gotReq)
125+
if tt.err != nil {
126+
require.ErrorIs(t, err, tt.err)
127+
assert.Nil(t, resp)
128+
return
129+
}
130+
require.NoError(t, err)
131+
require.Len(t, resp.Histories, tt.wantHistories)
132+
if tt.wantHistories > 0 {
133+
assert.Equal(t, "request/1", resp.Histories[0].RequestId)
134+
require.Len(t, resp.Histories[0].Events, 1)
135+
assert.Equal(t, "accepted", resp.Histories[0].Events[0].GetRequestState())
136+
}
137+
})
138+
}
139+
}

‎service/stovepipe/server/main.go‎

Lines changed: 2 additions & 47 deletions
Original file line numberDiff line numberDiff line change
@@ -44,7 +44,7 @@ import (
4444
queueMySQL "github.com/uber/submitqueue/platform/extension/messagequeue/mysql"
4545
platformhook "github.com/uber/submitqueue/platform/hook"
4646
servicemq "github.com/uber/submitqueue/service/messagequeue"
47-
"github.com/uber/submitqueue/service/stovepipe/server/mapper"
47+
"github.com/uber/submitqueue/service/stovepipe/server/handler"
4848
"github.com/uber/submitqueue/stovepipe/controller"
4949
"github.com/uber/submitqueue/stovepipe/controller/build"
5050
"github.com/uber/submitqueue/stovepipe/controller/buildsignal"
@@ -65,47 +65,6 @@ import (
6565
"google.golang.org/grpc/reflection"
6666
)
6767

68-
// StovepipeServer wraps the controllers and implements the gRPC service interface.
69-
type StovepipeServer struct {
70-
pb.UnimplementedStovepipeServer
71-
pingController *controller.PingController
72-
ingestController *controller.IngestController
73-
requestHistoryController controller.RequestHistoryController
74-
}
75-
76-
// Ping delegates to the controller.
77-
func (s *StovepipeServer) Ping(ctx context.Context, req *pb.PingRequest) (*pb.PingResponse, error) {
78-
return s.pingController.Ping(ctx, req)
79-
}
80-
81-
// Ingest maps the wire request to an entity, delegates to the controller, and maps
82-
// the result back to the wire response.
83-
func (s *StovepipeServer) Ingest(ctx context.Context, req *pb.IngestRequest) (*pb.IngestResponse, error) {
84-
result, err := s.ingestController.Ingest(ctx, mapper.ProtoToIngestRequest(req))
85-
if err != nil {
86-
return nil, err
87-
}
88-
return mapper.IngestResultToProto(result), nil
89-
}
90-
91-
// GetRequestHistoryByID returns retained history for one request ID.
92-
func (s *StovepipeServer) GetRequestHistoryByID(ctx context.Context, req *pb.GetRequestHistoryByIDRequest) (*pb.GetRequestHistoryByIDResponse, error) {
93-
events, err := s.requestHistoryController.GetRequestHistoryByID(ctx, mapper.ProtoToGetRequestHistoryByIDRequest(req))
94-
if err != nil {
95-
return nil, err
96-
}
97-
return &pb.GetRequestHistoryByIDResponse{Events: mapper.HistoryEventsToProto(events)}, nil
98-
}
99-
100-
// GetRequestHistoryByURI returns retained histories for one commit URI.
101-
func (s *StovepipeServer) GetRequestHistoryByURI(ctx context.Context, req *pb.GetRequestHistoryByURIRequest) (*pb.GetRequestHistoryByURIResponse, error) {
102-
histories, err := s.requestHistoryController.GetRequestHistoryByURI(ctx, mapper.ProtoToGetRequestHistoryByURIRequest(req))
103-
if err != nil {
104-
return nil, err
105-
}
106-
return &pb.GetRequestHistoryByURIResponse{Histories: mapper.RequestHistoriesToProto(histories)}, nil
107-
}
108-
10968
// inMemoryCounter is a minimal, process-local counter.Counter used to wire the example
11069
// server. It is not durable; a real deployment supplies a persistent implementation
11170
// (e.g. platform/extension/counter/mysql).
@@ -372,11 +331,7 @@ func run() error {
372331
tenants,
373332
)
374333
requestHistoryController := controller.NewRequestHistoryController(logger.Sugar(), scope, storageFty)
375-
srv := &StovepipeServer{
376-
pingController: pingController,
377-
ingestController: ingestController,
378-
requestHistoryController: requestHistoryController,
379-
}
334+
srv := handler.NewStovepipeServer(pingController, ingestController, requestHistoryController)
380335
pb.RegisterStovepipeServer(grpcServer, srv)
381336

382337
// Register reflection service for debugging with grpcurl

0 commit comments

Comments
 (0)