From cf8e1962487f61b0c3a1d89ed93002f28969c95f Mon Sep 17 00:00:00 2001 From: windyboy Date: Wed, 24 Dec 2025 14:53:14 +0800 Subject: [PATCH] =?UTF-8?q?=F0=9F=94=A7=20Add=20nil=20logger=20handling=20?= =?UTF-8?q?in=20NewConsumerManager=20and=20NewStreamManager;=20update=20in?= =?UTF-8?q?tegration=20test=20to=20pass=20config?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- internal/infra/nats/consumer_manager.go | 32 ++++++++++++++++++- internal/infra/nats/stream_manager.go | 3 ++ .../jetstream_to_timescale_test.go | 2 +- 3 files changed, 35 insertions(+), 2 deletions(-) diff --git a/internal/infra/nats/consumer_manager.go b/internal/infra/nats/consumer_manager.go index 29f06e2..f7bec78 100644 --- a/internal/infra/nats/consumer_manager.go +++ b/internal/infra/nats/consumer_manager.go @@ -19,6 +19,9 @@ type ConsumerManager struct { // NewConsumerManager creates a new consumer manager func NewConsumerManager(js nats.JetStreamContext, streamName, consumerName, subject string, logger *zap.Logger) *ConsumerManager { + if logger == nil { + logger = zap.NewNop() + } return &ConsumerManager{ js: js, streamName: streamName, @@ -66,5 +69,32 @@ func (cm *ConsumerManager) CreatePullSubscription() (*nats.Subscription, error) // CreatePullSubscriptionWithRecovery creates a pull subscription func (cm *ConsumerManager) CreatePullSubscriptionWithRecovery(streamManager *StreamManager, consumerConfig *nats.ConsumerConfig) (*nats.Subscription, error) { - return cm.CreatePullSubscription() + sub, err := cm.CreatePullSubscription() + if err == nil { + return sub, nil + } + + // Attempt recovery when the consumer or stream is missing. + if !errors.Is(err, nats.ErrConsumerNotFound) && !errors.Is(err, nats.ErrStreamNotFound) { + return nil, fmt.Errorf("create pull subscription: %w", err) + } + + if streamManager != nil { + if streamErr := streamManager.EnsureStream(nil); streamErr != nil { + return nil, fmt.Errorf("recover stream %s: %w", cm.streamName, streamErr) + } + } + + if consumerConfig == nil { + return nil, fmt.Errorf("consumer config is required for recovery") + } + if err := cm.EnsureConsumer(consumerConfig); err != nil { + return nil, fmt.Errorf("recover consumer %s: %w", cm.consumerName, err) + } + + sub, err = cm.CreatePullSubscription() + if err != nil { + return nil, fmt.Errorf("create pull subscription after recovery: %w", err) + } + return sub, nil } diff --git a/internal/infra/nats/stream_manager.go b/internal/infra/nats/stream_manager.go index de0c588..244b0d9 100644 --- a/internal/infra/nats/stream_manager.go +++ b/internal/infra/nats/stream_manager.go @@ -29,6 +29,9 @@ type StreamConfig struct { // NewStreamManager creates a new stream manager func NewStreamManager(js nats.JetStreamContext, streamName string, subjects []string, logger *zap.Logger) *StreamManager { + if logger == nil { + logger = zap.NewNop() + } return &StreamManager{ js: js, streamName: streamName, diff --git a/test/integration/jetstream_to_timescale_test.go b/test/integration/jetstream_to_timescale_test.go index e90a6f6..4e4aeb9 100644 --- a/test/integration/jetstream_to_timescale_test.go +++ b/test/integration/jetstream_to_timescale_test.go @@ -82,7 +82,7 @@ func TestJetStreamToTimescaleFlow(t *testing.T) { } telemetryRecorder := telemetryinfra.NewNoop() - proc := app.NewMessageProcessor(parser.ProvideParser(), repo, publisher, telemetryRecorder, logger) + proc := app.NewMessageProcessor(parser.ProvideParser(), repo, publisher, telemetryRecorder, logger, cfg) consumer, err := natsinfra.ProvideConsumer(conn, js, proc, cfg, telemetryRecorder, logger) if err != nil { t.Fatalf("failed to init consumer: %v", err)