@@ -23,6 +23,7 @@ import (
2323 "github.com/uber-go/tally"
2424 runwaymq "github.com/uber/submitqueue/api/runway/messagequeue"
2525 runwaypb "github.com/uber/submitqueue/api/runway/messagequeue/protopb"
26+ "github.com/uber/submitqueue/platform/base/failure"
2627 entityqueue "github.com/uber/submitqueue/platform/base/messagequeue"
2728 "github.com/uber/submitqueue/platform/consumer"
2829 consumermock "github.com/uber/submitqueue/platform/consumer/mock"
@@ -43,14 +44,22 @@ type publishedMsg struct {
4344 msg entityqueue.Message
4445}
4546
46- func newDelivery (t * testing.T , ctrl * gomock.Controller , payload []byte , meta map [string ]string ) * consumermock.MockDelivery {
47+ func newDelivery (
48+ t * testing.T ,
49+ ctrl * gomock.Controller ,
50+ payload []byte ,
51+ meta map [string ]string ,
52+ recordedFailure failure.Failure ,
53+ failed bool ,
54+ ) * consumermock.MockDelivery {
4755 t .Helper ()
4856 msg := entityqueue .NewMessage (testID , payload , testPartitionKey , nil )
4957 msg .Tenant = testQueue
5058 d := consumermock .NewMockDelivery (ctrl )
5159 d .EXPECT ().Message ().Return (msg ).AnyTimes ()
5260 d .EXPECT ().Metadata ().Return (meta ).AnyTimes ()
5361 d .EXPECT ().Attempt ().Return (1 ).AnyTimes ()
62+ d .EXPECT ().Failure ().Return (recordedFailure , failed ).AnyTimes ()
5463 return d
5564}
5665
@@ -89,45 +98,93 @@ func newController(t *testing.T, registry consumer.TopicRegistry) *Controller {
8998 })
9099}
91100
92- func TestProcess_DecodableRepublishesFailure (t * testing.T ) {
93- ctrl := gomock .NewController (t )
94- registry , published := newRegistry (t , ctrl )
95- controller := newController (t , registry )
96-
101+ func mergeRequestPayload (t * testing.T ) []byte {
102+ t .Helper ()
97103 req := & runwaymq.MergeRequest {
98104 Id : testID ,
99105 QueueName : testQueue ,
100106 Steps : []* runwaymq.MergeStep {{StepId : "step-1" }},
101107 }
102108 payload , err := runwaymq .Marshal (req )
103109 require .NoError (t , err )
110+ return payload
111+ }
104112
105- meta := map [string ]string {
106- "dlq.last_error" : "boom: connection refused" ,
107- "dlq.original_topic" : "runway-merge" ,
113+ func TestProcess_DecodableRepublishesFailure (t * testing.T ) {
114+ payload := mergeRequestPayload (t )
115+ tests := []struct {
116+ name string
117+ recordedFailure failure.Failure
118+ failed bool
119+ metadata map [string ]string
120+ wantReason string
121+ }{
122+ {
123+ name : "portable failure overrides conflicting legacy metadata" ,
124+ recordedFailure : failure .New ("git provider rejected the push" ),
125+ failed : true ,
126+ metadata : map [string ]string {
127+ "dlq.last_error" : "legacy metadata reason" ,
128+ "dlq.original_topic" : "runway-merge" ,
129+ },
130+ wantReason : "dead-lettered: git provider rejected the push" ,
131+ },
132+ {
133+ name : "legacy metadata is used when portable failure is absent" ,
134+ failed : false ,
135+ metadata : map [string ]string {
136+ "dlq.last_error" : "boom: connection refused" ,
137+ "dlq.original_topic" : "runway-merge" ,
138+ },
139+ wantReason : "dead-lettered: boom: connection refused" ,
140+ },
141+ {
142+ name : "default is used when no failure reason is available" ,
143+ failed : false ,
144+ metadata : map [string ]string {"dlq.original_topic" : "runway-merge" },
145+ wantReason : "dead-lettered: runway failed to process the merge request" ,
146+ },
108147 }
109- delivery := newDelivery (t , ctrl , payload , meta )
110148
111- require .NoError (t , controller .Process (context .Background (), delivery ))
112-
113- require .Len (t , * published , 1 )
114- got := (* published )[0 ]
115- assert .Equal (t , "merge-signal" , got .topic )
116- assert .Equal (t , testQueue , got .msg .Tenant )
117-
118- result := & runwaymq.MergeResult {}
119- require .NoError (t , runwaymq .Unmarshal (got .msg .Payload , result ))
120- assert .Equal (t , testID , result .Id )
121- assert .Equal (t , runwaypb .Outcome_FAILED , result .Outcome )
122- assert .Contains (t , result .Reason , "boom: connection refused" )
149+ for _ , tt := range tests {
150+ t .Run (tt .name , func (t * testing.T ) {
151+ ctrl := gomock .NewController (t )
152+ registry , published := newRegistry (t , ctrl )
153+ controller := newController (t , registry )
154+ delivery := newDelivery (t , ctrl , payload , tt .metadata , tt .recordedFailure , tt .failed )
155+
156+ require .NoError (t , controller .Process (context .Background (), delivery ))
157+
158+ require .Len (t , * published , 1 )
159+ got := (* published )[0 ]
160+ assert .Equal (t , "merge-signal" , got .topic )
161+ assert .Equal (t , testID + "/dlq" , got .msg .ID )
162+ assert .Equal (t , testPartitionKey , got .msg .PartitionKey )
163+ assert .Equal (t , testQueue , got .msg .Tenant )
164+
165+ result := & runwaymq.MergeResult {}
166+ require .NoError (t , runwaymq .Unmarshal (got .msg .Payload , result ))
167+ assert .Equal (t , testID , result .Id )
168+ assert .Equal (t , testQueue , result .QueueName )
169+ assert .Equal (t , runwaypb .Outcome_FAILED , result .Outcome )
170+ assert .Equal (t , tt .wantReason , result .Reason )
171+ })
172+ }
123173}
124174
125175func TestProcess_UndecodableAcksAndPublishesNothing (t * testing.T ) {
126176 ctrl := gomock .NewController (t )
127177 registry , published := newRegistry (t , ctrl )
128178 controller := newController (t , registry )
129179
130- delivery := newDelivery (t , ctrl , []byte ("{bad" ), map [string ]string {"dlq.original_topic" : "runway-merge" })
180+ delivery := newDelivery (
181+ t ,
182+ ctrl ,
183+ []byte ("{bad" ),
184+ map [string ]string {"dlq.original_topic" : "runway-merge" },
185+ failure.Failure {},
186+ false ,
187+ )
131188
132189 require .NoError (t , controller .Process (context .Background (), delivery ))
133190 assert .Empty (t , * published )
0 commit comments