Skip to content

Commit c8ad4fe

Browse files
committed
fix(consumer): derive controllerCtx from Background, not caller context
consumer.subscribe created controllerCtx from the caller's ctx, which inherits any deadline (e.g. Fx OnStart's 15-second timeout). The consume loop selects on controllerCtx.Done() and exits when the deadline fires, silently dropping all subsequent messages. The subscriber already uses context.Background(); this aligns the consumer to match.
1 parent aab7e64 commit c8ad4fe

2 files changed

Lines changed: 51 additions & 2 deletions

File tree

platform/consumer/consumer.go

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -191,8 +191,8 @@ func (m *consumer) subscribe(ctx context.Context, controller Controller) error {
191191
return fmt.Errorf("subscribe failed: %w", err)
192192
}
193193

194-
// Create cancellable context for this controller
195-
controllerCtx, cancel := context.WithCancel(ctx)
194+
// Background-rooted so a short-lived caller deadline (e.g. Fx OnStart) cannot kill the loop.
195+
controllerCtx, cancel := context.WithCancel(context.Background())
196196

197197
// Track active subscription
198198
done := make(chan struct{})

platform/consumer/consumer_test.go

Lines changed: 49 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -811,3 +811,52 @@ func TestConsumer_PartitionWorkerCleanup(t *testing.T) {
811811
err = c.Stop(30000)
812812
require.NoError(t, err)
813813
}
814+
815+
func TestConsumer_ConsumeLoopSurvivesCallerDeadline(t *testing.T) {
816+
ctrl := gomock.NewController(t)
817+
logger := zaptest.NewLogger(t).Sugar()
818+
819+
deliveryChan := make(chan extqueue.Delivery, 1)
820+
mockSub := queuemock.NewMockSubscriber(ctrl)
821+
mockSub.EXPECT().Subscribe(gomock.Any(), gomock.Any(), gomock.Any()).Return(deliveryChan, nil)
822+
823+
mockQ := queuemock.NewMockQueue(ctrl)
824+
mockQ.EXPECT().Subscriber().Return(mockSub)
825+
826+
reg := newRegistry(t, mockQ, topickey.TopicKeyStart, "test-group")
827+
828+
c := consumer.New(logger, tally.NoopScope, reg, errs.NewClassifierProcessor())
829+
830+
processed := make(chan string, 1)
831+
handler := consumermock.NewMockController(ctrl)
832+
setupController(handler, "test-handler", topickey.TopicKeyStart, "test-group",
833+
func(ctx context.Context, delivery consumer.Delivery) error {
834+
processed <- delivery.Message().ID
835+
return nil
836+
},
837+
)
838+
839+
err := c.Register(handler)
840+
require.NoError(t, err)
841+
842+
// Start with a context that expires quickly, simulating an Fx OnStart hook.
843+
startCtx, startCancel := context.WithTimeout(context.Background(), 50*time.Millisecond)
844+
defer startCancel()
845+
846+
err = c.Start(startCtx)
847+
require.NoError(t, err)
848+
849+
<-startCtx.Done()
850+
851+
msg := entityqueue.NewMessage("after-deadline", []byte("payload"), "partition1", nil)
852+
mockDel := queuemock.NewMockDelivery(ctrl)
853+
done := setupDelivery(mockDel, msg, nil, nil)
854+
855+
deliveryChan <- mockDel
856+
<-done
857+
858+
assert.Equal(t, "after-deadline", <-processed)
859+
860+
err = c.Stop(30000)
861+
require.NoError(t, err)
862+
}

0 commit comments

Comments
 (0)