@@ -61,6 +61,9 @@ import (
6161 "github.com/uber/submitqueue/submitqueue/extension/scorer/composite"
6262 scorerfake "github.com/uber/submitqueue/submitqueue/extension/scorer/fake"
6363 "github.com/uber/submitqueue/submitqueue/extension/scorer/heuristic"
64+ prioritizationlimitstatic "github.com/uber/submitqueue/submitqueue/extension/speculation/prioritizationlimit/static"
65+ "github.com/uber/submitqueue/submitqueue/extension/speculation/prioritizer"
66+ "github.com/uber/submitqueue/submitqueue/extension/speculation/prioritizer/sticky"
6467 "github.com/uber/submitqueue/submitqueue/extension/storage"
6568 mysqlstorage "github.com/uber/submitqueue/submitqueue/extension/storage/mysql"
6669 validatorfake "github.com/uber/submitqueue/submitqueue/extension/validator/fake"
@@ -74,6 +77,7 @@ import (
7477 "github.com/uber/submitqueue/submitqueue/orchestrator/controller/merge"
7578 "github.com/uber/submitqueue/submitqueue/orchestrator/controller/mergeconflictsignal"
7679 "github.com/uber/submitqueue/submitqueue/orchestrator/controller/mergesignal"
80+ "github.com/uber/submitqueue/submitqueue/orchestrator/controller/prioritize"
7781 "github.com/uber/submitqueue/submitqueue/orchestrator/controller/score"
7882 "github.com/uber/submitqueue/submitqueue/orchestrator/controller/speculate"
7983 "github.com/uber/submitqueue/submitqueue/orchestrator/controller/start"
@@ -244,13 +248,14 @@ func run() error {
244248 brf := buildRunnerFactory {queues }
245249 scf := scorerFactory {queues }
246250 cof := analyzerFactory {queues }
251+ prf := prioritizerFactory {queues }
247252
248253 // Register controllers
249- primaryCount , err := registerPrimaryControllers (primaryConsumer , logger .Sugar (), scope , registry , cpf , brf , scf , cof , cnt , store )
254+ primaryCount , err := registerPrimaryControllers (primaryConsumer , logger .Sugar (), scope , registry , cpf , brf , scf , cof , prf , cnt , store )
250255 if err != nil {
251256 return err
252257 }
253- dlqCount , err := registerDLQControllers (dlqConsumer , logger .Sugar (), scope , store )
258+ dlqCount , err := registerDLQControllers (dlqConsumer , logger .Sugar (), scope , registry , store )
254259 if err != nil {
255260 return err
256261 }
@@ -380,9 +385,10 @@ func newTopicRegistry(q extqueue.Queue, subscriberName string) (consumer.TopicRe
380385 {topickey .TopicKeyBatch , "batch" , "orchestrator-batch" },
381386 {topickey .TopicKeyScore , "score" , "orchestrator-score" },
382387 {topickey .TopicKeySpeculate , "speculate" , "orchestrator-speculate" },
388+ {topickey .TopicKeyPrioritize , "prioritize" , "orchestrator-prioritize" },
383389 {topickey .TopicKeyBuild , "build" , "orchestrator-build" },
384390 {topickey .TopicKeyBuildSignal , "buildsignal" , "orchestrator-buildsignal" },
385- {topickey .TopicKeyMerge , "merge" , "orchestrator-merge" },
391+ {topickey .TopicKeyMerge , "submitqueue- merge" , "orchestrator-merge" },
386392 {runwaymq .TopicKeyMergeSignal , "merge-signal" , "orchestrator-mergesignal" },
387393 {topickey .TopicKeyConclude , "conclude" , "orchestrator-conclude" },
388394 }
@@ -450,7 +456,7 @@ func newTopicRegistry(q extqueue.Queue, subscriberName string) (consumer.TopicRe
450456 // consumed primary topic above.
451457 configs = append (configs , consumer.TopicConfig {
452458 Key : runwaymq .TopicKeyMerge ,
453- Name : "merge" ,
459+ Name : "runway- merge" ,
454460 Queue : q ,
455461 })
456462
@@ -471,6 +477,13 @@ func newTopicRegistry(q extqueue.Queue, subscriberName string) (consumer.TopicRe
471477// merge-conflict-check queue (⇢); runway performs the merge attempt and
472478// publishes the result to merge-conflict-check-signal, which mergeconflictsignal
473479// consumes before fanning the request out to batch.
480+ //
481+ // prioritize sits alongside this per-batch flow rather than in its line: it
482+ // is queue-wide, not batch-scoped. Its message carries only a queue name; on
483+ // each invocation it loads every Speculating batch's speculation tree for
484+ // that queue, ranks the queue-wide candidate paths against the queue's build
485+ // budget, applies the resulting decisions, and republishes to build for any
486+ // path newly (or still) cleared to run.
474487
475488// TODO(wiring abstraction): queueExtensions + queueRegistry currently live here
476489// as example-local wiring. Evaluate promoting them into a defined abstraction in
@@ -493,6 +506,7 @@ type queueExtensions struct {
493506 buildRunner buildrunner.BuildRunner
494507 scorer scorer.Scorer
495508 analyzer conflict.Analyzer
509+ prioritizer prioritizer.Prioritizer
496510}
497511
498512// queueRegistry maps a queue name to its extensions, falling back to a default
@@ -538,7 +552,13 @@ func (f analyzerFactory) For(cfg conflict.Config) (conflict.Analyzer, error) {
538552 return f .reg .get (cfg .QueueName ).analyzer , nil
539553}
540554
541- func registerPrimaryControllers (c consumer.Consumer , logger * zap.SugaredLogger , scope tally.Scope , registry consumer.TopicRegistry , cpf changeprovider.Factory , brf buildrunner.Factory , scf scorer.Factory , cof conflict.Factory , cnt counter.Counter , store storage.Storage ) (int , error ) {
555+ type prioritizerFactory struct { reg queueRegistry }
556+
557+ func (f prioritizerFactory ) For (cfg prioritizer.Config ) (prioritizer.Prioritizer , error ) {
558+ return f .reg .get (cfg .QueueName ).prioritizer , nil
559+ }
560+
561+ func registerPrimaryControllers (c consumer.Consumer , logger * zap.SugaredLogger , scope tally.Scope , registry consumer.TopicRegistry , cpf changeprovider.Factory , brf buildrunner.Factory , scf scorer.Factory , cof conflict.Factory , prf prioritizer.Factory , cnt counter.Counter , store storage.Storage ) (int , error ) {
542562 var count int
543563 requestController := start .NewController (
544564 logger ,
@@ -637,6 +657,20 @@ func registerPrimaryControllers(c consumer.Consumer, logger *zap.SugaredLogger,
637657 }
638658 count ++
639659
660+ prioritizeController := prioritize .NewController (
661+ logger ,
662+ scope ,
663+ store ,
664+ prf ,
665+ registry ,
666+ topickey .TopicKeyPrioritize ,
667+ "orchestrator-prioritize" ,
668+ )
669+ if err := c .Register (prioritizeController ); err != nil {
670+ return count , fmt .Errorf ("failed to register prioritize controller: %w" , err )
671+ }
672+ count ++
673+
640674 buildController := build .NewController (
641675 logger ,
642676 scope ,
@@ -712,7 +746,7 @@ func registerPrimaryControllers(c consumer.Consumer, logger *zap.SugaredLogger,
712746// registers them with the DLQ consumer. Each reconciler drives the affected
713747// request or batch into a terminal Error/Failed state so the gateway stops
714748// reporting it as stuck-in-progress.
715- func registerDLQControllers (c consumer.Consumer , logger * zap.SugaredLogger , scope tally.Scope , store storage.Storage ) (int , error ) {
749+ func registerDLQControllers (c consumer.Consumer , logger * zap.SugaredLogger , scope tally.Scope , registry consumer. TopicRegistry , store storage.Storage ) (int , error ) {
716750 dlqScope := scope .SubScope ("dlq" )
717751 dlqRegs := []struct {
718752 name string
@@ -725,6 +759,7 @@ func registerDLQControllers(c consumer.Consumer, logger *zap.SugaredLogger, scop
725759 {"batch_dlq" , dlq .NewDLQRequestController (logger , dlqScope , store , dlq .DecodeRequestID , dlq .TopicKey (topickey .TopicKeyBatch ), "orchestrator-batch-dlq" )},
726760 {"score_dlq" , dlq .NewDLQBatchController (logger , dlqScope , store , dlq .TopicKey (topickey .TopicKeyScore ), "orchestrator-score-dlq" )},
727761 {"speculate_dlq" , dlq .NewDLQBatchController (logger , dlqScope , store , dlq .TopicKey (topickey .TopicKeySpeculate ), "orchestrator-speculate-dlq" )},
762+ {"prioritize_dlq" , dlq .NewDLQQueueController (logger , dlqScope , registry , dlq .TopicKey (topickey .TopicKeyPrioritize ), "orchestrator-prioritize-dlq" )},
728763 {"build_dlq" , dlq .NewDLQBatchController (logger , dlqScope , store , dlq .TopicKey (topickey .TopicKeyBuild ), "orchestrator-build-dlq" )},
729764 {"buildsignal_dlq" , dlq .NewDLQBuildSignalController (logger , dlqScope , store , dlq .TopicKey (topickey .TopicKeyBuildSignal ), "orchestrator-buildsignal-dlq" )},
730765 {"merge_dlq" , dlq .NewDLQBatchController (logger , dlqScope , store , dlq .TopicKey (topickey .TopicKeyMerge ), "orchestrator-merge-dlq" )},
@@ -859,6 +894,11 @@ func newPhabChangeProvider(logger *zap.Logger, scope tally.Scope) (changeprovide
859894 }), nil
860895}
861896
897+ // defaultPrioritizationLimit is the baseline queue-wide concurrent-build
898+ // budget handed to the sticky prioritizer. It is a parity default —
899+ // effectively admit-all — until per-queue budgets are configured.
900+ const defaultPrioritizationLimit = 1000
901+
862902// newQueueRegistry builds the per-queue extension profiles for the example.
863903// Edge integrations (change provider) and the build
864904// runner form a shared baseline; each per-queue profile starts from that
@@ -880,16 +920,19 @@ func newQueueRegistry(logger *zap.Logger, scope tally.Scope, resolver changeset.
880920
881921 // Baseline profile: shared edge integrations + a fake build runner (every
882922 // build succeeds unless a head URI carries a failure marker), plus permissive
883- // defaults for scorer and conflict. The build runner instance is shared by
884- // the build and buildsignal controllers (same profile, same instance) so a
885- // build's recorded outcome survives across their separate factory lookups.
923+ // defaults for scorer, conflict, and prioritization. The build runner
924+ // instance is shared by the build and buildsignal controllers (same
925+ // profile, same instance) so a build's recorded outcome survives across
926+ // their separate factory lookups.
886927 //
887928 // The scorer is wrapped by scorerfake so a change URI carrying
888929 // "sq-fake=score-error" forces a scoring error end-to-end; it is a pure
889930 // passthrough otherwise. The analyzer is wrapped by conflictfake with a nil
890931 // predicate (passthrough) — swap the predicate (e.g. conflictfake.FailAlways)
891932 // on a queue to exercise the analyzer error path, as e2e-conflict-error-queue
892- // below does.
933+ // below does. The prioritizer is sticky over a static budget: it never
934+ // preempts a running build and admits Selected candidates by score until
935+ // defaultPrioritizationLimit concurrent builds are in flight.
893936 base := queueExtensions {
894937 changeProvider : cp ,
895938 buildRunner : buildfake .New (resolver ),
@@ -900,7 +943,8 @@ func newQueueRegistry(logger *zap.Logger, scope tally.Scope, resolver changeset.
900943 )),
901944 // TODO: replace the delegate with a real analyzer (e.g. Tango target
902945 // analysis). "all" serializes the queue conservatively.
903- analyzer : conflictfake .New (all .New (), nil ),
946+ analyzer : conflictfake .New (all .New (), nil ),
947+ prioritizer : sticky .New (prioritizationlimitstatic .New (defaultPrioritizationLimit )),
904948 }
905949
906950 // test-queue: bucketed heuristic scorer; conservative (serialized) conflicts
0 commit comments