diff --git a/runway/README.md b/runway/README.md index ed4882b7..5113f738 100644 --- a/runway/README.md +++ b/runway/README.md @@ -11,5 +11,8 @@ Runway is a single service (the domain *is* the service); its controllers live d - `merge-conflict-check` — dry-run check that an ordered sequence of merge steps applies cleanly, without committing. - `merge` — committing merge: apply and commit the ordered steps. -Both controllers currently deserialize the `MergeRequest` off the queue and log it; performing the -merge and publishing a `MergeResult` to the corresponding signal queue is not wired yet. +Each controller deserializes the `MergeRequest`, obtains a `Merger` for the request's queue from the [`merger`](extension/merger) extension, applies the ordered steps, and publishes a `MergeResult` to the corresponding signal queue (`merge-conflict-check-signal` / `merge-signal`). `merge` commits and reports the produced revisions; `merge-conflict-check` is a dry run that reports mergeability with empty outputs. + +## Failure handling + +A merge outcome the controller can name is published as a `FAILED` result and acked, not retried: a merge conflict (`merger.ErrConflict`) or an invalid request (`merger.ErrInvalidRequest` — unknown strategy, malformed change URI, invalid PROMOTE composition). The `merger.IsTerminal` helper draws that line. Any other error is an infrastructure fault and is nacked for retry. diff --git a/runway/controller/merge/merge.go b/runway/controller/merge/merge.go index f89a9f0a..1d78c9d1 100644 --- a/runway/controller/merge/merge.go +++ b/runway/controller/merge/merge.go @@ -103,15 +103,24 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er result, err := m.Merge(ctx, request) if err != nil { - if !errors.Is(err, merger.ErrConflict) { + if !merger.IsTerminal(err) { metrics.NamedCounter(c.metricsScope, opName, "merge_errors", 1) return fmt.Errorf("failed to merge for %s: %w", request.GetId(), err) } - metrics.NamedCounter(c.metricsScope, opName, "merge_conflicts", 1) - c.logger.Infow("merge conflict detected", - "id", request.GetId(), - "queue_name", request.GetQueueName(), - ) + if errors.Is(err, merger.ErrInvalidRequest) { + metrics.NamedCounter(c.metricsScope, opName, "invalid_requests", 1) + c.logger.Infow("invalid merge request", + "id", request.GetId(), + "queue_name", request.GetQueueName(), + "err", err, + ) + } else { + metrics.NamedCounter(c.metricsScope, opName, "merge_conflicts", 1) + c.logger.Infow("merge conflict detected", + "id", request.GetId(), + "queue_name", request.GetQueueName(), + ) + } result = &runwaymq.MergeResult{ Id: request.GetId(), Outcome: runwaypb.Outcome_FAILED, diff --git a/runway/controller/merge/merge_test.go b/runway/controller/merge/merge_test.go index 0b6ce6e6..bab87891 100644 --- a/runway/controller/merge/merge_test.go +++ b/runway/controller/merge/merge_test.go @@ -198,6 +198,51 @@ func TestProcess_MergeConflict(t *testing.T) { assert.NotEmpty(t, result.Reason) } +func TestProcess_InvalidRequest(t *testing.T) { + ctrl := gomock.NewController(t) + + m := mergermock.NewMockMerger(ctrl) + m.EXPECT().Merge(gomock.Any(), gomock.Any()).Return(nil, fmt.Errorf("bad strategy: %w", merger.ErrInvalidRequest)) + + factory := mergermock.NewMockFactory(ctrl) + factory.EXPECT().For(merger.Config{QueueName: testQueue}).Return(m, nil) + + var gotTopic string + var gotPayload []byte + pub := queuemock.NewMockPublisher(ctrl) + pub.EXPECT().Publish(gomock.Any(), gomock.Any(), gomock.Any()).DoAndReturn( + func(_ context.Context, topic string, msg entityqueue.Message) error { + gotTopic = topic + gotPayload = msg.Payload + return nil + }, + ) + q := queuemock.NewMockQueue(ctrl) + q.EXPECT().Publisher().Return(pub).AnyTimes() + registry, err := consumer.NewTopicRegistry([]consumer.TopicConfig{ + {Key: runwaymq.TopicKeyMergeSignal, Name: "merge-signal", Queue: q}, + }) + require.NoError(t, err) + + controller := newController(t, factory, registry) + + req := &runwaymq.MergeRequest{ + Id: testID, + QueueName: testQueue, + Steps: []*runwaymq.MergeStep{{StepId: "step-1"}}, + } + delivery := newDelivery(t, ctrl, requestPayload(t, req)) + + require.NoError(t, controller.Process(context.Background(), delivery)) + + assert.Equal(t, "merge-signal", gotTopic) + result := &runwaymq.MergeResult{} + require.NoError(t, runwaymq.Unmarshal(gotPayload, result)) + assert.Equal(t, testID, result.Id) + assert.Equal(t, runwaypb.Outcome_FAILED, result.Outcome) + assert.NotEmpty(t, result.Reason) +} + func TestProcess_MergerInfraError(t *testing.T) { ctrl := gomock.NewController(t) diff --git a/runway/controller/mergeconflictcheck/mergeconflictcheck.go b/runway/controller/mergeconflictcheck/mergeconflictcheck.go index 4cdd2436..8ccfa017 100644 --- a/runway/controller/mergeconflictcheck/mergeconflictcheck.go +++ b/runway/controller/mergeconflictcheck/mergeconflictcheck.go @@ -103,15 +103,24 @@ func (c *Controller) Process(ctx context.Context, delivery consumer.Delivery) er result, err := m.CheckMergeability(ctx, request) if err != nil { - if !errors.Is(err, merger.ErrConflict) { + if !merger.IsTerminal(err) { metrics.NamedCounter(c.metricsScope, opName, "check_errors", 1) return fmt.Errorf("failed to check mergeability for %s: %w", request.GetId(), err) } - metrics.NamedCounter(c.metricsScope, opName, "merge_conflicts", 1) - c.logger.Infow("merge conflict detected", - "id", request.GetId(), - "queue_name", request.GetQueueName(), - ) + if errors.Is(err, merger.ErrInvalidRequest) { + metrics.NamedCounter(c.metricsScope, opName, "invalid_requests", 1) + c.logger.Infow("invalid merge request", + "id", request.GetId(), + "queue_name", request.GetQueueName(), + "err", err, + ) + } else { + metrics.NamedCounter(c.metricsScope, opName, "merge_conflicts", 1) + c.logger.Infow("merge conflict detected", + "id", request.GetId(), + "queue_name", request.GetQueueName(), + ) + } result = &runwaymq.MergeResult{ Id: request.GetId(), Outcome: runwaypb.Outcome_FAILED, diff --git a/runway/controller/mergeconflictcheck/mergeconflictcheck_test.go b/runway/controller/mergeconflictcheck/mergeconflictcheck_test.go index 31e1c239..31d9d497 100644 --- a/runway/controller/mergeconflictcheck/mergeconflictcheck_test.go +++ b/runway/controller/mergeconflictcheck/mergeconflictcheck_test.go @@ -195,6 +195,51 @@ func TestProcess_MergeConflict(t *testing.T) { assert.NotEmpty(t, result.Reason) } +func TestProcess_InvalidRequest(t *testing.T) { + ctrl := gomock.NewController(t) + + m := mergermock.NewMockMerger(ctrl) + m.EXPECT().CheckMergeability(gomock.Any(), gomock.Any()).Return(nil, fmt.Errorf("bad strategy: %w", merger.ErrInvalidRequest)) + + factory := mergermock.NewMockFactory(ctrl) + factory.EXPECT().For(merger.Config{QueueName: testQueue}).Return(m, nil) + + var gotTopic string + var gotPayload []byte + pub := queuemock.NewMockPublisher(ctrl) + pub.EXPECT().Publish(gomock.Any(), gomock.Any(), gomock.Any()).DoAndReturn( + func(_ context.Context, topic string, msg entityqueue.Message) error { + gotTopic = topic + gotPayload = msg.Payload + return nil + }, + ) + q := queuemock.NewMockQueue(ctrl) + q.EXPECT().Publisher().Return(pub).AnyTimes() + registry, err := consumer.NewTopicRegistry([]consumer.TopicConfig{ + {Key: runwaymq.TopicKeyMergeConflictCheckSignal, Name: "merge-conflict-check-signal", Queue: q}, + }) + require.NoError(t, err) + + controller := newController(t, factory, registry) + + req := &runwaymq.MergeRequest{ + Id: testID, + QueueName: testQueue, + Steps: []*runwaymq.MergeStep{{StepId: "step-1"}}, + } + delivery := newDelivery(t, ctrl, requestPayload(t, req)) + + require.NoError(t, controller.Process(context.Background(), delivery)) + + assert.Equal(t, "merge-conflict-check-signal", gotTopic) + result := &runwaymq.MergeResult{} + require.NoError(t, runwaymq.Unmarshal(gotPayload, result)) + assert.Equal(t, testID, result.Id) + assert.Equal(t, runwaypb.Outcome_FAILED, result.Outcome) + assert.NotEmpty(t, result.Reason) +} + func TestProcess_MergerInfraError(t *testing.T) { ctrl := gomock.NewController(t) diff --git a/runway/extension/merger/merger.go b/runway/extension/merger/merger.go index 9a77963a..8371e0f7 100644 --- a/runway/extension/merger/merger.go +++ b/runway/extension/merger/merger.go @@ -32,6 +32,20 @@ import ( // result), not an infrastructure error. var ErrConflict = errors.New("merge conflict") +// ErrInvalidRequest signals that the request can never be applied as written — +// an unknown merge strategy, a malformed change URI, or an invalid strategy +// composition (e.g. PROMOTE mixed with other steps). Like ErrConflict it is a +// terminal outcome: retrying never succeeds, so controllers ack and publish a +// failure result rather than nacking into an infinite retry / dead-letter. +var ErrInvalidRequest = errors.New("invalid merge request") + +// IsTerminal reports whether err is a terminal merge outcome — one the client +// must be told about via a FAILED result rather than retried. Controllers use +// it to decide between publishing a FAILED result (ack) and nacking for retry. +func IsTerminal(err error) bool { + return errors.Is(err, ErrConflict) || errors.Is(err, ErrInvalidRequest) +} + // Merger performs version-control operations against a single merge target. // Both methods accept the same MergeRequest payload; the behavioral difference // is whether the result is committed to the remote.