Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions platform/consumer/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -41,6 +41,8 @@ go_test(
"@com_github_stretchr_testify//require:go_default_library",
"@com_github_uber_go_tally//:go_default_library",
"@org_uber_go_mock//gomock:go_default_library",
"@org_uber_go_zap//:go_default_library",
"@org_uber_go_zap//zaptest:go_default_library",
"@org_uber_go_zap//zaptest/observer:go_default_library",
],
)
92 changes: 60 additions & 32 deletions platform/consumer/consumer.go
Original file line number Diff line number Diff line change
Expand Up @@ -414,13 +414,16 @@ func (m *consumer) processDelivery(ctx context.Context, controller Controller, d
}

topicKey := controller.TopicKey()
ownerLogFields := deliveryOwnerLogFields(delivery)

m.logger.Debugw("processing delivery",
"controller", controller.Name(),
"topic_key", topicKey,
"message_id", msg.ID,
"partition_key", msg.PartitionKey,
"attempt", delivery.Attempt(),
append([]any{
"controller", controller.Name(),
"topic_key", topicKey,
"message_id", msg.ID,
"partition_key", msg.PartitionKey,
"attempt", delivery.Attempt(),
}, ownerLogFields...)...,
)

// Wrap delivery to hide Ack/Nack from controller
Expand Down Expand Up @@ -453,10 +456,12 @@ func (m *consumer) processDelivery(ctx context.Context, controller Controller, d
if wrapped.held {
metrics.NamedCounter(controllerScope, opName, "hold_ignored", 1, metrics.TagsFromContext(ctx)...)
m.logger.Warnw("hold recorded but controller returned error, failure outcome wins",
"controller", controller.Name(),
"topic_key", topicKey,
"message_id", msg.ID,
"partition_key", msg.PartitionKey,
append([]any{
"controller", controller.Name(),
"topic_key", topicKey,
"message_id", msg.ID,
"partition_key", msg.PartitionKey,
}, ownerLogFields...)...,
)
}

Expand All @@ -473,13 +478,15 @@ func (m *consumer) processDelivery(ctx context.Context, controller Controller, d
// Check if the error is non-retryable (poison pill message)
if !errs.IsRetryable(err) {
m.logger.Errorw("non-retryable controller error, rejecting message",
"controller", controller.Name(),
"topic_key", controller.TopicKey(),
"message_id", msg.ID,
"partition_key", msg.PartitionKey,
"attempt", delivery.Attempt(),
"error", err,
"elapsed_ms", elapsed.Milliseconds(),
append([]any{
"controller", controller.Name(),
"topic_key", controller.TopicKey(),
"message_id", msg.ID,
"partition_key", msg.PartitionKey,
"attempt", delivery.Attempt(),
"error", err,
"elapsed_ms", elapsed.Milliseconds(),
}, ownerLogFields...)...,
)

// Reject moves to DLQ (or acks if DLQ disabled)
Expand All @@ -488,10 +495,12 @@ func (m *consumer) processDelivery(ctx context.Context, controller Controller, d
rejectOp.Complete(rejectErr)
if rejectErr != nil {
m.logger.Errorw("failed to reject non-retryable message",
"controller", controller.Name(),
"topic_key", controller.TopicKey(),
"message_id", msg.ID,
"error", rejectErr,
append([]any{
"controller", controller.Name(),
"topic_key", controller.TopicKey(),
"message_id", msg.ID,
"error", rejectErr,
}, ownerLogFields...)...,
)
}
return
Expand All @@ -504,25 +513,29 @@ func (m *consumer) processDelivery(ctx context.Context, controller Controller, d
what = "cancel"
}
m.logger.Errorw("controller error or cancel, nacking message",
"what", what,
"controller", controller.Name(),
"topic_key", topicKey,
"message_id", msg.ID,
"partition_key", msg.PartitionKey,
"attempt", delivery.Attempt(),
"error", err,
"elapsed_ms", elapsed.Milliseconds(),
append([]any{
"what", what,
"controller", controller.Name(),
"topic_key", topicKey,
"message_id", msg.ID,
"partition_key", msg.PartitionKey,
"attempt", delivery.Attempt(),
"error", err,
"elapsed_ms", elapsed.Milliseconds(),
}, ownerLogFields...)...,
)

nackOp := metrics.Begin(controllerScope, "nack", metrics.StorageLatencyBuckets, metrics.TagsFromContext(ctx)...)
nackErr := delivery.Nack(ctx, controllerFailure)
nackOp.Complete(nackErr)
if nackErr != nil {
m.logger.Errorw("failed to nack message",
"controller", controller.Name(),
"topic_key", topicKey,
"message_id", msg.ID,
"error", nackErr,
append([]any{
"controller", controller.Name(),
"topic_key", topicKey,
"message_id", msg.ID,
"error", nackErr,
}, ownerLogFields...)...,
)
}
return
Expand Down Expand Up @@ -764,3 +777,18 @@ func (m *consumer) unsubscribeAll(timeoutMs int64) error {
m.logger.Debugw("all controllers stopped gracefully")
return nil
}

func deliveryOwnerLogFields(delivery extqueue.Delivery) []any {
meta := delivery.Metadata()
if len(meta) == 0 {
return nil
}
var fields []any
if v := meta["leased_by"]; v != "" {
fields = append(fields, "leased_by", v)
}
if v := meta["consumer_group"]; v != "" {
fields = append(fields, "consumer_group", v)
}
return fields
}
121 changes: 121 additions & 0 deletions platform/consumer/consumer_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,9 @@ import (
queuemock "github.com/uber/submitqueue/platform/extension/messagequeue/mock"
"github.com/uber/submitqueue/platform/metrics"
"go.uber.org/mock/gomock"
"go.uber.org/zap"
"go.uber.org/zap/zaptest"
"go.uber.org/zap/zaptest/observer"
)

const (
Expand Down Expand Up @@ -627,6 +629,123 @@ func TestConsumer_ProcessDelivery_NonRetryableError(t *testing.T) {
require.NoError(t, err)
}

func TestConsumer_ProcessDelivery_LogsLeasedBy(t *testing.T) {
ctrl := gomock.NewController(t)
core, logs := observer.New(zap.DebugLevel)
logger := zap.New(core).Sugar()

deliveryChan := make(chan extqueue.Delivery, 1)
mockSub := queuemock.NewMockSubscriber(ctrl)
mockSub.EXPECT().Subscribe(gomock.Any(), gomock.Any(), gomock.Any()).Return(deliveryChan, nil)

mockQ := queuemock.NewMockQueue(ctrl)
mockQ.EXPECT().Subscriber().Return(mockSub)

reg := newRegistry(t, mockQ, testTopicKeyStart, "test-group")
c := New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())

handler := &testController{}
setupController(handler, "test-handler", testTopicKeyStart, "test-group",
func(ctx context.Context, delivery Delivery) error {
return fmt.Errorf("bad payload")
},
)
require.NoError(t, c.Register(handler))

ctx, cancel := context.WithCancel(context.Background())
defer cancel()
require.NoError(t, c.Start(ctx))

msg := entityqueue.NewMessage("poison-msg", []byte("bad"), "partition1", nil)
msg.Tenant = testTenant
done := make(chan struct{})
mockDel := queuemock.NewMockDelivery(ctrl)
mockDel.EXPECT().Message().Return(msg).AnyTimes()
mockDel.EXPECT().Attempt().Return(1).AnyTimes()
mockDel.EXPECT().ReceivedAt().Return(time.Now().UnixMilli()).AnyTimes()
mockDel.EXPECT().Metadata().Return(map[string]string{
"leased_by": "host-1",
"consumer_group": "test-group",
}).AnyTimes()
mockDel.EXPECT().DeliveryID().Return(msg.ID).AnyTimes()
mockDel.EXPECT().Reject(gomock.Any(), gomock.Any()).DoAndReturn(func(ctx context.Context, _ failure.Failure) error {
close(done)
return nil
})

deliveryChan <- mockDel
<-done
require.NoError(t, c.Stop(30000))

processLogs := logs.FilterMessage("processing delivery").All()
require.NotEmpty(t, processLogs)
assert.Equal(t, "host-1", processLogs[0].ContextMap()["leased_by"])
assert.Equal(t, "test-group", processLogs[0].ContextMap()["consumer_group"])

rejectLogs := logs.FilterMessage("non-retryable controller error, rejecting message").All()
require.NotEmpty(t, rejectLogs)
assert.Equal(t, "host-1", rejectLogs[0].ContextMap()["leased_by"])
assert.Equal(t, "partition1", rejectLogs[0].ContextMap()["partition_key"])
}

func TestConsumer_ProcessDelivery_HoldIgnoredLogsLeasedBy(t *testing.T) {
ctrl := gomock.NewController(t)
core, logs := observer.New(zap.DebugLevel)
logger := zap.New(core).Sugar()

deliveryChan := make(chan extqueue.Delivery, 1)
mockSub := queuemock.NewMockSubscriber(ctrl)
mockSub.EXPECT().Subscribe(gomock.Any(), gomock.Any(), gomock.Any()).Return(deliveryChan, nil)

mockQ := queuemock.NewMockQueue(ctrl)
mockQ.EXPECT().Subscriber().Return(mockSub)

reg := newRegistry(t, mockQ, testTopicKeyStart, "test-group")
c := New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor(), consumergatenoop.New())

deliveryMetadata := map[string]string{
"leased_by": "host-1",
"consumer_group": "test-group",
}
handler := &testController{}
setupController(handler, "test-handler", testTopicKeyStart, "test-group",
func(ctx context.Context, delivery Delivery) error {
delivery.Metadata()["leased_by"] = "mutated-host"
delivery.Metadata()["consumer_group"] = "mutated-group"
delivery.Hold(5000)
return errs.NewRetryableError(fmt.Errorf("processing failed"))
},
)
require.NoError(t, c.Register(handler))

ctx, cancel := context.WithCancel(context.Background())
defer cancel()
require.NoError(t, c.Start(ctx))

msg := entityqueue.NewMessage("held-msg", []byte("payload"), "partition1", nil)
msg.Tenant = testTenant
done := make(chan struct{})
mockDel := queuemock.NewMockDelivery(ctrl)
mockDel.EXPECT().Message().Return(msg).AnyTimes()
mockDel.EXPECT().Attempt().Return(1).AnyTimes()
mockDel.EXPECT().ReceivedAt().Return(time.Now().UnixMilli()).AnyTimes()
mockDel.EXPECT().Metadata().Return(deliveryMetadata).AnyTimes()
mockDel.EXPECT().DeliveryID().Return(msg.ID).AnyTimes()
mockDel.EXPECT().Nack(gomock.Any(), gomock.Any()).DoAndReturn(func(ctx context.Context, _ failure.Failure) error {
close(done)
return nil
})

deliveryChan <- mockDel
<-done
require.NoError(t, c.Stop(30000))

holdLogs := logs.FilterMessage("hold recorded but controller returned error, failure outcome wins").All()
require.Len(t, holdLogs, 1)
assert.Equal(t, "host-1", holdLogs[0].ContextMap()["leased_by"])
assert.Equal(t, "test-group", holdLogs[0].ContextMap()["consumer_group"])
}

// The failure handed to the queue is built from whatever the controller
// attributed, and a controller that attributes nothing must still produce
// exactly what callers sent before failures carried structure: the error text
Expand Down Expand Up @@ -1176,6 +1295,7 @@ func TestConsumer_SamePartitionKeyAcrossTenantsProcessesIndependently(t *testing
delA := queuemock.NewMockDelivery(ctrl)
delA.EXPECT().Message().Return(msgA).AnyTimes()
delA.EXPECT().Attempt().Return(1).AnyTimes()
delA.EXPECT().Metadata().Return(nil).AnyTimes()
delA.EXPECT().Ack(gomock.Any()).Return(nil).MaxTimes(1)
deliveryChan <- delA
<-tenantABlocked
Expand All @@ -1185,6 +1305,7 @@ func TestConsumer_SamePartitionKeyAcrossTenantsProcessesIndependently(t *testing
delB := queuemock.NewMockDelivery(ctrl)
delB.EXPECT().Message().Return(msgB).AnyTimes()
delB.EXPECT().Attempt().Return(1).AnyTimes()
delB.EXPECT().Metadata().Return(nil).AnyTimes()
delB.EXPECT().Ack(gomock.Any()).Return(nil).MaxTimes(1)
deliveryChan <- delB
<-tenantBProcessed
Expand Down
1 change: 1 addition & 0 deletions platform/extension/messagequeue/mysql/BUILD.bazel
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ go_test(
"@com_github_uber_go_tally//:go_default_library",
"@org_uber_go_mock//gomock:go_default_library",
"@org_uber_go_zap//:go_default_library",
"@org_uber_go_zap//zapcore:go_default_library",
"@org_uber_go_zap//zaptest:go_default_library",
"@org_uber_go_zap//zaptest/observer:go_default_library",
],
Expand Down
16 changes: 11 additions & 5 deletions platform/extension/messagequeue/mysql/constants.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,9 +18,15 @@ package mysql

const (
// Common log field names (used extensively across all stores)
logTenant = "tenant"
logTopic = "topic"
logPartitionKey = "partition_key"
logMessageID = "message_id"
logError = "error"
logTenant = "tenant"
logTopic = "topic"
logPartitionKey = "partition_key"
logMessageID = "message_id"
logError = "error"
logLeasedBy = "leased_by"
logConsumerGroup = "consumer_group"
logPreviousOwner = "previous_owner"
logReason = "reason"
logOwnedPartitions = "owned_partitions"
logRenewed = "renewed"
)
21 changes: 12 additions & 9 deletions platform/extension/messagequeue/mysql/mock_stores.go

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading
Loading