Skip to content

Commit 9bc9a5f

Browse files
committed
fix(orchestrator): add Validated state transition and log in mergeconflictsignal
The mergeconflictsignal controller was publishing directly to the batch topic on success without transitioning the request to Validated or recording a log entry. This adds the missing CAS Started→Validated state transition and publishes a RequestStatusValidated log entry before forwarding to batch.
1 parent 4ddcdb9 commit 9bc9a5f

2 files changed

Lines changed: 38 additions & 11 deletions

File tree

submitqueue/orchestrator/controller/mergeconflictsignal/mergeconflictsignal.go

Lines changed: 16 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -121,21 +121,34 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) (r
121121
"request_id", request.ID,
122122
"reason", result.Reason,
123123
)
124-
// Expected terminal outcome, not a failure: mark the request Error inline
125-
// and ack.
126124
if err := c.failRequest(ctx, request, result.Reason); err != nil {
127125
metrics.NamedCounter(c.metricsScope, opName, "fail_errors", 1)
128126
return fmt.Errorf("failed to fail request %s: %w", request.ID, err)
129127
}
130128
return nil
131129
}
132130

131+
// Advance the request to Validated now that the merge-conflict check passed.
132+
newVersion := request.Version + 1
133+
if err := c.store.GetRequestStore().UpdateState(ctx, request.ID, request.Version, newVersion, entity.RequestStateValidated); err != nil {
134+
metrics.NamedCounter(c.metricsScope, opName, "state_errors", 1)
135+
return fmt.Errorf("failed to update request %s state to validated: %w", request.ID, err)
136+
}
137+
request.Version = newVersion
138+
request.State = entity.RequestStateValidated
139+
140+
logEntry := entity.NewRequestLog(request.ID, entity.RequestStatusValidated, request.Version, "", nil)
141+
if err := corerequest.PublishLog(ctx, c.registry, logEntry, request.ID); err != nil {
142+
metrics.NamedCounter(c.metricsScope, opName, "log_errors", 1)
143+
return fmt.Errorf("failed to publish request log for %s: %w", request.ID, err)
144+
}
145+
133146
if err := c.publishRequestID(ctx, topickey.TopicKeyBatch, request.ID, request.Queue); err != nil {
134147
metrics.NamedCounter(c.metricsScope, opName, "publish_errors", 1)
135148
return fmt.Errorf("failed to publish to batch: %w", err)
136149
}
137150

138-
c.logger.Infow("published request to batch",
151+
c.logger.Infow("request validated and published to batch",
139152
"request_id", request.ID,
140153
"topic_key", topickey.TopicKeyBatch,
141154
)

submitqueue/orchestrator/controller/mergeconflictsignal/mergeconflictsignal_test.go

Lines changed: 22 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -57,24 +57,28 @@ func TestProcess_MergeablePublishesToBatch(t *testing.T) {
5757
reqStore := storagemock.NewMockRequestStore(ctrl)
5858
reqStore.EXPECT().Get(gomock.Any(), testRequestID).Return(
5959
entity.Request{ID: testRequestID, Queue: testQueue, State: entity.RequestStateStarted, Version: 1}, nil)
60+
reqStore.EXPECT().UpdateState(gomock.Any(), testRequestID, int32(1), int32(2), entity.RequestStateValidated).Return(nil)
6061

6162
store := storagemock.NewMockStorage(ctrl)
6263
store.EXPECT().GetRequestStore().Return(reqStore).AnyTimes()
6364

64-
var gotTopic string
65-
var gotPayload []byte
65+
var gotTopics []string
66+
var gotPayloads [][]byte
6667
pub := queuemock.NewMockPublisher(ctrl)
6768
pub.EXPECT().Publish(gomock.Any(), gomock.Any(), gomock.Any()).DoAndReturn(
6869
func(_ context.Context, topic string, msg entityqueue.Message) error {
69-
gotTopic = topic
70-
gotPayload = msg.Payload
70+
gotTopics = append(gotTopics, topic)
71+
gotPayloads = append(gotPayloads, msg.Payload)
7172
return nil
7273
},
73-
)
74+
).Times(2)
7475
q := queuemock.NewMockQueue(ctrl)
7576
q.EXPECT().Publisher().Return(pub).AnyTimes()
7677
registry, err := consumer.NewTopicRegistry(
77-
[]consumer.TopicConfig{{Key: topickey.TopicKeyBatch, Name: "batch", Queue: q}},
78+
[]consumer.TopicConfig{
79+
{Key: topickey.TopicKeyBatch, Name: "batch", Queue: q},
80+
{Key: topickey.TopicKeyLog, Name: "log", Queue: q},
81+
},
7882
)
7983
require.NoError(t, err)
8084

@@ -85,8 +89,18 @@ func TestProcess_MergeablePublishesToBatch(t *testing.T) {
8589
msg := entityqueue.NewMessage(testRequestID, resultPayload(t, res), testQueue, nil)
8690
require.NoError(t, controller.Process(context.Background(), newDelivery(ctrl, msg)))
8791

88-
assert.Equal(t, "batch", gotTopic)
89-
rid, err := entity.RequestIDFromBytes(gotPayload)
92+
require.Len(t, gotTopics, 2)
93+
94+
// First publish: validated log entry.
95+
assert.Equal(t, "log", gotTopics[0])
96+
logEntry, err := entity.RequestLogFromBytes(gotPayloads[0])
97+
require.NoError(t, err)
98+
assert.Equal(t, entity.RequestStatusValidated, logEntry.Status)
99+
assert.Equal(t, int32(2), logEntry.RequestVersion)
100+
101+
// Second publish: request ID to batch topic.
102+
assert.Equal(t, "batch", gotTopics[1])
103+
rid, err := entity.RequestIDFromBytes(gotPayloads[1])
90104
require.NoError(t, err)
91105
assert.Equal(t, testRequestID, rid.ID)
92106
}

0 commit comments

Comments
 (0)