Skip to content

Commit 85cb46b

Browse files
kevinlnewclaude
andcommitted
feat: widen Pusher interface with QueueTarget and PushItem
Consolidate QueueTarget into the shared entity/ package and update all pushqueue imports accordingly. Widen Pusher.Push to accept entity.QueueTarget and []PushItem (Change + Strategy). Add queueconfig.Store dependency to the merge controller for resolving queue names to VCS targets. Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
1 parent e871d8a commit 85cb46b

30 files changed

Lines changed: 277 additions & 121 deletions

File tree

entity/BUILD.bazel

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -2,7 +2,10 @@ load("@rules_go//go:def.bzl", "go_library")
22

33
go_library(
44
name = "entity",
5-
srcs = ["entity.go"],
5+
srcs = [
6+
"entity.go",
7+
"queue_target.go",
8+
],
69
importpath = "github.com/uber/submitqueue/entity",
710
visibility = ["//visibility:public"],
811
)

entity/queue_target.go

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,25 @@
1+
// Copyright (c) 2025 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 entity
16+
17+
// QueueTarget identifies a landing destination in a version control system.
18+
type QueueTarget struct {
19+
// Name is an optional logical identifier for correlation and config lookup.
20+
Name string
21+
// Address is the VCS repository address (remote URL, depot path).
22+
Address string
23+
// Target is the landing ref (branch name, stream path).
24+
Target string
25+
}

example/pushqueue/gateway/server/BUILD.bazel

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ go_library(
66
importpath = "github.com/uber/submitqueue/example/pushqueue/gateway/server",
77
visibility = ["//visibility:private"],
88
deps = [
9+
"//entity",
910
"//pushqueue/entity",
1011
"//pushqueue/extension/landqueue",
1112
"//pushqueue/extension/vcs",

example/pushqueue/gateway/server/main.go

Lines changed: 7 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -26,7 +26,8 @@ import (
2626
"time"
2727

2828
"github.com/uber-go/tally/v4"
29-
"github.com/uber/submitqueue/pushqueue/entity"
29+
"github.com/uber/submitqueue/entity"
30+
pqentity "github.com/uber/submitqueue/pushqueue/entity"
3031
"github.com/uber/submitqueue/pushqueue/extension/landqueue"
3132
"github.com/uber/submitqueue/pushqueue/extension/vcs"
3233
"github.com/uber/submitqueue/pushqueue/gateway/controller"
@@ -162,23 +163,23 @@ func run() error {
162163
// noopVCS is a placeholder VCS that errors on every operation.
163164
type noopVCS struct{}
164165

165-
func (noopVCS) CheckMergeability(_ context.Context, _ entity.QueueTarget, items []entity.LandItem) ([]vcs.MergeabilityResult, error) {
166+
func (noopVCS) CheckMergeability(_ context.Context, _ entity.QueueTarget, items []pqentity.LandItem) ([]vcs.MergeabilityResult, error) {
166167
results := make([]vcs.MergeabilityResult, len(items))
167168
for i := range results {
168169
results[i] = vcs.MergeabilityResult{Mergeable: false, Reason: "noop VCS: not configured"}
169170
}
170171
return results, nil
171172
}
172173

173-
func (noopVCS) Prepare(_ context.Context, _ entity.QueueTarget, _ []entity.LandItem) error {
174+
func (noopVCS) Prepare(_ context.Context, _ entity.QueueTarget, _ []pqentity.LandItem) error {
174175
return fmt.Errorf("noop VCS: not configured")
175176
}
176177

177-
func (noopVCS) Push(_ context.Context, _ entity.QueueTarget, _ []entity.LandItem) (vcs.PushResult, error) {
178+
func (noopVCS) Push(_ context.Context, _ entity.QueueTarget, _ []pqentity.LandItem) (vcs.PushResult, error) {
178179
return vcs.PushResult{}, fmt.Errorf("noop VCS: not configured")
179180
}
180181

181-
func (noopVCS) Finalize(_ context.Context, _ entity.QueueTarget, _ []entity.LandItem) error {
182+
func (noopVCS) Finalize(_ context.Context, _ entity.QueueTarget, _ []pqentity.LandItem) error {
182183
return fmt.Errorf("noop VCS: not configured")
183184
}
184185

@@ -187,7 +188,7 @@ type noopQueue struct{}
187188

188189
var _ landqueue.Queue = noopQueue{}
189190

190-
func (noopQueue) Enqueue(_ context.Context, _ entity.QueueTarget, _ []entity.LandItem) error {
191+
func (noopQueue) Enqueue(_ context.Context, _ entity.QueueTarget, _ []pqentity.LandItem) error {
191192
return nil
192193
}
193194

example/submitqueue/orchestrator/server/BUILD.bazel

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -14,6 +14,7 @@ go_library(
1414
"//core/errs/generic",
1515
"//core/errs/mysql",
1616
"//core/httpclient",
17+
"//entity",
1718
"//extension/queue",
1819
"//extension/queue/mysql",
1920
"//submitqueue/core/consumer",
@@ -31,6 +32,8 @@ go_library(
3132
"//submitqueue/extension/mergechecker/github",
3233
"//submitqueue/extension/pusher",
3334
"//submitqueue/extension/pusher/git",
35+
"//submitqueue/extension/queueconfig",
36+
"//submitqueue/extension/queueconfig/yaml",
3437
"//submitqueue/extension/scorer/heuristic",
3538
"//submitqueue/extension/storage",
3639
"//submitqueue/extension/storage/mysql",

example/submitqueue/orchestrator/server/main.go

Lines changed: 25 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,7 @@ import (
3333
genericerrs "github.com/uber/submitqueue/core/errs/generic"
3434
mysqlerrs "github.com/uber/submitqueue/core/errs/mysql"
3535
"github.com/uber/submitqueue/core/httpclient"
36+
sharedentity "github.com/uber/submitqueue/entity"
3637
extqueue "github.com/uber/submitqueue/extension/queue"
3738
queueMySQL "github.com/uber/submitqueue/extension/queue/mysql"
3839
"github.com/uber/submitqueue/submitqueue/core/consumer"
@@ -50,6 +51,8 @@ import (
5051
githubchecker "github.com/uber/submitqueue/submitqueue/extension/mergechecker/github"
5152
"github.com/uber/submitqueue/submitqueue/extension/pusher"
5253
gitpusher "github.com/uber/submitqueue/submitqueue/extension/pusher/git"
54+
"github.com/uber/submitqueue/submitqueue/extension/queueconfig"
55+
yamlqueueconfig "github.com/uber/submitqueue/submitqueue/extension/queueconfig/yaml"
5356
"github.com/uber/submitqueue/submitqueue/extension/scorer/heuristic"
5457
"github.com/uber/submitqueue/submitqueue/extension/storage"
5558
mysqlstorage "github.com/uber/submitqueue/submitqueue/extension/storage/mysql"
@@ -231,8 +234,14 @@ func run() error {
231234
// (every build immediately succeeds) until a real backend is wired in.
232235
br := buildnoop.New()
233236

237+
// Create queue config store
238+
qcfg, err := newQueueConfigStore(logger)
239+
if err != nil {
240+
return fmt.Errorf("failed to create queue config store: %w", err)
241+
}
242+
234243
// Register controllers
235-
if err := registerControllers(c, logger.Sugar(), scope, registry, mc, cp, psh, br, cnt, store, changeStore); err != nil {
244+
if err := registerControllers(c, logger.Sugar(), scope, registry, mc, cp, psh, br, cnt, store, changeStore, qcfg); err != nil {
236245
return err
237246
}
238247

@@ -425,7 +434,7 @@ func newTopicRegistry(q extqueue.Queue, subscriberName string) (consumer.TopicRe
425434
// │ │ │
426435
// └────────┴───────────────────────┘
427436

428-
func registerControllers(c consumer.Consumer, logger *zap.SugaredLogger, scope tally.Scope, registry consumer.TopicRegistry, mc mergechecker.MergeChecker, cp changeprovider.ChangeProvider, psh pusher.Pusher, br buildrunner.BuildRunner, cnt counter.Counter, store storage.Storage, changeStore changestore.ChangeStore) error {
437+
func registerControllers(c consumer.Consumer, logger *zap.SugaredLogger, scope tally.Scope, registry consumer.TopicRegistry, mc mergechecker.MergeChecker, cp changeprovider.ChangeProvider, psh pusher.Pusher, br buildrunner.BuildRunner, cnt counter.Counter, store storage.Storage, changeStore changestore.ChangeStore, qcfg queueconfig.Store) error {
429438
requestController := start.NewController(
430439
logger,
431440
scope,
@@ -551,6 +560,7 @@ func registerControllers(c consumer.Consumer, logger *zap.SugaredLogger, scope t
551560
store,
552561
registry,
553562
psh,
563+
qcfg,
554564
consumer.TopicKeyMerge,
555565
"orchestrator-merge",
556566
)
@@ -675,6 +685,18 @@ func newPusher(logger *zap.Logger, scope tally.Scope) (pusher.Pusher, error) {
675685
}), nil
676686
}
677687

688+
// newQueueConfigStore loads queue configuration from a YAML file pointed to by
689+
// QUEUE_CONFIG_PATH. If the env var is not set, returns an empty store — the
690+
// merge controller will fail to resolve any queue target.
691+
func newQueueConfigStore(logger *zap.Logger) (queueconfig.Store, error) {
692+
path := os.Getenv("QUEUE_CONFIG_PATH")
693+
if path == "" {
694+
logger.Warn("QUEUE_CONFIG_PATH not set; merge controller will fail to resolve queue targets")
695+
return yamlqueueconfig.Store{}, nil
696+
}
697+
return yamlqueueconfig.NewStore(path)
698+
}
699+
678700
// noopPusher is a fallback Pusher used when PUSHER_CHECKOUT_PATH is not
679701
// configured. It returns an error on every Push so the merge controller
680702
// (which treats non-ErrConflict errors as transient and nacks the message)
@@ -683,6 +705,6 @@ func newPusher(logger *zap.Logger, scope tally.Scope) (pusher.Pusher, error) {
683705
// that don't run the merge step.
684706
type noopPusher struct{}
685707

686-
func (noopPusher) Push(_ context.Context, _ []entity.Change) (pusher.Result, error) {
708+
func (noopPusher) Push(_ context.Context, _ sharedentity.QueueTarget, _ []pusher.PushItem) (pusher.Result, error) {
687709
return pusher.Result{}, fmt.Errorf("pusher not configured: set PUSHER_CHECKOUT_PATH to enable pushing")
688710
}

pushqueue/entity/land.go

Lines changed: 0 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -14,17 +14,6 @@
1414

1515
package entity
1616

17-
// QueueTarget identifies a landing destination in a version control system.
18-
// Defined locally in pushqueue; consolidated into shared entity/ by Chunk 2.
19-
type QueueTarget struct {
20-
// Name is an optional logical identifier for correlation and config lookup.
21-
Name string
22-
// Address is the VCS repository address (remote URL, depot path).
23-
Address string
24-
// Target is the landing ref (branch name, stream path).
25-
Target string
26-
}
27-
2817
// LandStrategy defines the possible landing methods for a code change.
2918
type LandStrategy string
3019

pushqueue/extension/landqueue/BUILD.bazel

Lines changed: 4 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -5,5 +5,8 @@ go_library(
55
srcs = ["landqueue.go"],
66
importpath = "github.com/uber/submitqueue/pushqueue/extension/landqueue",
77
visibility = ["//visibility:public"],
8-
deps = ["//pushqueue/entity"],
8+
deps = [
9+
"//entity",
10+
"//pushqueue/entity",
11+
],
912
)

pushqueue/extension/landqueue/landqueue.go

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -19,15 +19,16 @@ package landqueue
1919
import (
2020
"context"
2121

22-
"github.com/uber/submitqueue/pushqueue/entity"
22+
"github.com/uber/submitqueue/entity"
23+
pqentity "github.com/uber/submitqueue/pushqueue/entity"
2324
)
2425

2526
// Preparer performs pre-push preparation for a set of land items.
2627
// Implementations are typically backed by a VCS extension. The Queue
2728
// receives a Preparer at construction time and may invoke it between
2829
// Enqueue and Wait to pipeline preparation with queue wait time.
2930
type Preparer interface {
30-
Prepare(ctx context.Context, target entity.QueueTarget, items []entity.LandItem) error
31+
Prepare(ctx context.Context, target entity.QueueTarget, items []pqentity.LandItem) error
3132
}
3233

3334
// Queue serializes access to a landing target, ensuring only one request
@@ -44,7 +45,7 @@ type Preparer interface {
4445
// - Implementations decide when to call Preparer.Prepare — during the
4546
// wait (pipelined) or synchronously before Wait returns (simple).
4647
type Queue interface {
47-
Enqueue(ctx context.Context, target entity.QueueTarget, items []entity.LandItem) error
48+
Enqueue(ctx context.Context, target entity.QueueTarget, items []pqentity.LandItem) error
4849
Wait(ctx context.Context, target entity.QueueTarget) error
4950
Dequeue(ctx context.Context, target entity.QueueTarget) error
5051
}

pushqueue/extension/landqueue/mock/BUILD.bazel

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,6 +6,7 @@ go_library(
66
importpath = "github.com/uber/submitqueue/pushqueue/extension/landqueue/mock",
77
visibility = ["//visibility:public"],
88
deps = [
9+
"//entity",
910
"//pushqueue/entity",
1011
"@org_uber_go_mock//gomock",
1112
],

0 commit comments

Comments
 (0)