@@ -33,6 +33,7 @@ import (
3333 "encoding/hex"
3434 "errors"
3535 "fmt"
36+ "math"
3637 "strconv"
3738 "time"
3839
@@ -47,6 +48,7 @@ import (
4748 stovepipemq "github.com/uber/submitqueue/stovepipe/core/messagequeue"
4849 "github.com/uber/submitqueue/stovepipe/core/requestlog"
4950 "github.com/uber/submitqueue/stovepipe/entity"
51+ "github.com/uber/submitqueue/stovepipe/extension/projectresult"
5052 "github.com/uber/submitqueue/stovepipe/extension/sourcecontrol"
5153 "github.com/uber/submitqueue/stovepipe/extension/storage"
5254 "go.uber.org/zap"
@@ -60,6 +62,7 @@ type Controller struct {
6062 metricsScope tally.Scope
6163 stores storage.Factory
6264 materializer requestlog.Materializer
65+ projectResult projectresult.Factory
6366 sourceControl sourcecontrol.Factory
6467 registry consumer.TopicRegistry
6568 topicKey consumer.TopicKey
@@ -73,8 +76,7 @@ var _ consumer.Controller = (*Controller)(nil)
7376const _opName = "record"
7477
7578// wholeRepositoryProject is the project component of a fact covering the whole
76- // repository rather than one project within it. Per-project facts need target-graph
77- // attribution that this stage does not do, so every fact it writes is whole-repository.
79+ // repository rather than one project within it.
7880const wholeRepositoryProject = ""
7981
8082// NewController creates a new record controller.
@@ -83,6 +85,7 @@ func NewController(
8385 scope tally.Scope ,
8486 stores storage.Factory ,
8587 materializer requestlog.Materializer ,
88+ projectResult projectresult.Factory ,
8689 sourceControl sourcecontrol.Factory ,
8790 registry consumer.TopicRegistry ,
8891 topicKey consumer.TopicKey ,
@@ -94,6 +97,7 @@ func NewController(
9497 metricsScope : scope .SubScope (name ),
9598 stores : stores ,
9699 materializer : materializer ,
100+ projectResult : projectResult ,
97101 sourceControl : sourceControl ,
98102 registry : registry ,
99103 topicKey : topicKey ,
@@ -150,6 +154,9 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
150154 if err := c .persistValidationFactRecordedLog (ctx , store , request , fact ); err != nil {
151155 return err
152156 }
157+ if err := c .recordProjectFacts (ctx , store , request ); err != nil {
158+ return err
159+ }
153160 if err := c .applyFactToDerivedCaches (ctx , store , request , fact , created ); err != nil {
154161 return err
155162 }
@@ -176,6 +183,46 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er
176183 }
177184}
178185
186+ func (c * Controller ) recordProjectFacts (ctx context.Context , store storage.Storage , request entity.Request ) error {
187+ resolver , err := c .projectResult .For (projectresult.Config {QueueName : request .Queue })
188+ if err != nil {
189+ return fmt .Errorf ("failed to resolve project result resolver for queue %q: %w" , request .Queue , err )
190+ }
191+ results , err := resolver .Resolve (ctx , request )
192+ if err != nil {
193+ return fmt .Errorf ("failed to resolve project results for request %q: %w" , request .ID , err )
194+ }
195+
196+ seen := make (map [string ]struct {}, len (results ))
197+ for _ , result := range results {
198+ if result .Project == "" {
199+ return fmt .Errorf ("project result for request %q has an empty project" , request .ID )
200+ }
201+ if _ , ok := seen [result .Project ]; ok {
202+ return fmt .Errorf ("project result for request %q contains duplicate project %q" , request .ID , result .Project )
203+ }
204+ seen [result .Project ] = struct {}{}
205+ if math .IsNaN (result .Degree ) || result .Degree < entity .DegreeGreen || result .Degree > entity .DegreeBroken {
206+ return fmt .Errorf ("project result for request %q and project %q has invalid degree %v" , request .ID , result .Project , result .Degree )
207+ }
208+
209+ fact , _ , err := c .recordValidationFact (ctx , store .GetValidationFactStore (), entity.ValidationFact {
210+ URI : request .URI ,
211+ Project : result .Project ,
212+ Degree : result .Degree ,
213+ RequestID : request .ID ,
214+ CreatedAt : time .Now ().UnixMilli (),
215+ })
216+ if err != nil {
217+ return err
218+ }
219+ if err := c .persistValidationFactRecordedLog (ctx , store , request , fact ); err != nil {
220+ return err
221+ }
222+ }
223+ return nil
224+ }
225+
179226func (c * Controller ) persistValidationFactRecordedLog (
180227 ctx context.Context ,
181228 store storage.Storage ,
@@ -241,48 +288,49 @@ func (c *Controller) applyFactToDerivedCaches(
241288// second return reports whether this call is the one that wrote the fact, which is
242289// how a caller tells the original delivery from a redelivery.
243290func (c * Controller ) recordFact (ctx context.Context , store storage.Storage , request entity.Request ) (entity.ValidationFact , bool , error ) {
244- factStore := store .GetValidationFactStore ()
245-
246- fact := entity.ValidationFact {
291+ return c .recordValidationFact (ctx , store .GetValidationFactStore (), entity.ValidationFact {
247292 URI : request .URI ,
248293 Project : wholeRepositoryProject ,
249294 Degree : degreeFor (request .State ),
250295 RequestID : request .ID ,
251296 CreatedAt : time .Now ().UnixMilli (),
252- }
297+ })
298+ }
299+
300+ func (c * Controller ) recordValidationFact (ctx context.Context , factStore storage.ValidationFactStore , fact entity.ValidationFact ) (entity.ValidationFact , bool , error ) {
253301
254302 err := factStore .Create (ctx , fact )
255303 switch {
256304 case err == nil :
257305 metrics .NamedCounter (c .metricsScope , _opName , "fact_created" , 1 , metrics .TagsFromContext (ctx )... )
258306 c .logger .Infow ("recorded validation fact" ,
259- "queue " , request . Queue ,
260- "request_id " , request . ID ,
261- "uri " , request . URI ,
307+ "request_id " , fact . RequestID ,
308+ "uri " , fact . URI ,
309+ "project " , fact . Project ,
262310 "degree" , fact .Degree ,
263311 )
264312 return fact , true , nil
265313
266314 case errors .Is (err , storage .ErrAlreadyExists ):
267- stored , getErr := factStore .Get (ctx , request .URI , wholeRepositoryProject )
315+ stored , getErr := factStore .Get (ctx , fact .URI , fact . Project )
268316 if getErr != nil {
269317 metrics .NamedCounter (c .metricsScope , _opName , "storage_errors" , 1 , metrics .TagsFromContext (ctx )... )
270- return entity.ValidationFact {}, false , fmt .Errorf ("failed to load the existing fact for uri %s: %w" , request .URI , getErr )
318+ return entity.ValidationFact {}, false , fmt .Errorf ("failed to load the existing fact for uri %s and project %q : %w" , fact .URI , fact . Project , getErr )
271319 }
272- if stored .RequestID != request . ID {
320+ if stored .RequestID != fact . RequestID {
273321 // Two requests validating one URI would break the dedup ingest
274322 // enforces, so this is a broken invariant rather than a race to
275323 // resolve. Non-retryable: the stored fact is immutable.
276324 metrics .NamedCounter (c .metricsScope , _opName , "invariant_errors" , 1 , metrics .TagsFromContext (ctx )... )
277325 return entity.ValidationFact {}, false , fmt .Errorf (
278- "fact for uri %s is owned by request %s, not %s" , request .URI , stored .RequestID , request . ID )
326+ "fact for uri %s and project %q is owned by request %s, not %s" , fact .URI , fact . Project , stored .RequestID , fact . RequestID )
279327 }
280328 metrics .NamedCounter (c .metricsScope , _opName , "fact_exists" , 1 , metrics .TagsFromContext (ctx )... )
281329 return stored , false , nil
282330
283331 default :
284332 metrics .NamedCounter (c .metricsScope , _opName , "storage_errors" , 1 , metrics .TagsFromContext (ctx )... )
285- return entity.ValidationFact {}, false , fmt .Errorf ("failed to create the fact for uri %s: %w" , request .URI , err )
333+ return entity.ValidationFact {}, false , fmt .Errorf ("failed to create the fact for uri %s and project %q : %w" , fact .URI , fact . Project , err )
286334 }
287335}
288336
0 commit comments