Enhance NATS configuration and error handling. Introduce stream limits and consumer rules in configuration files. Refactor message processing to handle permanent errors. Update README and development configuration to reflect changes. Add tests for new error handling mechanisms.

This commit is contained in:
windyboy
2025-11-14 21:42:04 +08:00
parent a574cfcf27
commit 9cce4610b6
13 changed files with 510 additions and 57 deletions
+110 -14
View File
@@ -6,23 +6,25 @@ import (
"context"
"errors"
"fmt"
"go.uber.org/zap"
"github.com/nats-io/nats.go"
"go.uber.org/zap"
"time"
)
// Consumer handles NATS JetStream message consumption
type Consumer struct {
js nats.JetStreamContext
processor *app.MessageProcessor
cfg *config.Config
logger *zap.Logger
subject string
conn *nats.Conn
js nats.JetStreamContext
processor *app.MessageProcessor
cfg *config.Config
logger *zap.Logger
subject string
consumerName string
}
// ProvideConsumer creates a NATS consumer
func ProvideConsumer(
conn *nats.Conn,
js nats.JetStreamContext,
processor *app.MessageProcessor,
cfg *config.Config,
@@ -39,6 +41,7 @@ func ProvideConsumer(
}
consumer := &Consumer{
conn: conn,
js: js,
processor: processor,
cfg: cfg,
@@ -62,12 +65,21 @@ func (c *Consumer) ensureConsumer() error {
streamName = "TELEGRAM"
}
ackWait := c.cfg.NATS.ConsumerRules.AckWait
if ackWait == 0 {
ackWait = c.cfg.Timeouts.AckWait
}
if ackWait == 0 {
ackWait = 30 * time.Second
}
consumerConfig := &nats.ConsumerConfig{
Durable: c.consumerName,
DeliverPolicy: nats.DeliverAllPolicy,
AckPolicy: nats.AckExplicitPolicy,
AckWait: c.cfg.Timeouts.AckWait,
MaxDeliver: 5, // Maximum number of delivery attempts
AckWait: ackWait,
MaxDeliver: c.cfg.NATS.ConsumerRules.MaxDeliver,
MaxAckPending: c.cfg.NATS.ConsumerRules.MaxAckPending,
FilterSubject: c.subject,
}
@@ -116,6 +128,17 @@ func (c *Consumer) Start(ctx context.Context) error {
batchTimeout = 2 * time.Second
}
c.logger.Info("Consumer pull configuration",
zap.Int("batch_size", batchSize),
zap.Duration("batch_timeout", batchTimeout),
zap.Int("max_deliver", c.cfg.NATS.ConsumerRules.MaxDeliver),
zap.Duration("ack_wait", c.cfg.NATS.ConsumerRules.AckWait),
)
statsCtx, statsCancel := context.WithCancel(ctx)
defer statsCancel()
go c.emitConsumerStats(statsCtx, streamName)
for {
select {
case <-ctx.Done():
@@ -139,24 +162,97 @@ func (c *Consumer) Start(ctx context.Context) error {
// Process each message
for _, msg := range msgs {
if err := c.processMessage(ctx, msg); err != nil {
isPermanent := app.IsPermanent(err)
c.logger.Error("Failed to process message",
zap.String("subject", msg.Subject),
zap.Error(err),
zap.Bool("permanent", isPermanent),
)
// NAK the message to retry
if isPermanent {
if termErr := msg.Term(); termErr != nil {
c.logger.Error("Failed to TERM message", zap.Error(termErr))
}
continue
}
// Transient error: request redelivery
if nakErr := msg.Nak(); nakErr != nil {
c.logger.Error("Failed to NAK message", zap.Error(nakErr))
}
} else {
// ACK the message
if ackErr := msg.Ack(); ackErr != nil {
c.logger.Error("Failed to ACK message", zap.Error(ackErr))
}
continue
}
// ACK the message
if ackErr := msg.Ack(); ackErr != nil {
c.logger.Error("Failed to ACK message", zap.Error(ackErr))
}
}
}
}
func (c *Consumer) emitConsumerStats(ctx context.Context, streamName string) {
interval := c.cfg.App.MonitorInterval
if interval <= 0 {
interval = 30 * time.Second
}
ticker := time.NewTicker(interval)
defer ticker.Stop()
for {
select {
case <-ctx.Done():
return
case <-ticker.C:
info, err := c.js.ConsumerInfo(streamName, c.consumerName)
if err != nil {
c.logger.Warn("Failed to fetch consumer info", zap.Error(err))
continue
}
c.logger.Info("JetStream consumer metrics",
zap.String("stream", streamName),
zap.String("consumer", c.consumerName),
zap.Uint64("num_ack_pending", uint64(info.NumAckPending)),
zap.Uint64("num_redelivered", uint64(info.NumRedelivered)),
zap.Uint64("num_pending", uint64(info.NumPending)),
zap.Uint64("delivered_consumer_seq", uint64(info.Delivered.Consumer)),
zap.Uint64("delivered_stream_seq", uint64(info.Delivered.Stream)),
)
}
}
}
// Shutdown drains the underlying NATS connection gracefully.
func (c *Consumer) Shutdown(ctx context.Context) error {
if c.conn == nil {
return nil
}
timeout := c.cfg.Timeouts.Close
if timeout <= 0 {
timeout = 10 * time.Second
}
closeCtx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()
errCh := make(chan error, 1)
go func() {
errCh <- c.conn.Drain()
}()
select {
case err := <-errCh:
c.conn.Close()
return err
case <-closeCtx.Done():
c.conn.Close()
return fmt.Errorf("nats drain timeout: %w", closeCtx.Err())
}
}
// processMessage processes a single message
func (c *Consumer) processMessage(ctx context.Context, msg *nats.Msg) error {
msgID := msg.Header.Get("Nats-Msg-Id")
+30 -8
View File
@@ -3,14 +3,14 @@ package nats
import (
"caatsm/internal/infra/config"
"fmt"
"go.uber.org/zap"
"strings"
"github.com/nats-io/nats.go"
"time"
"go.uber.org/zap"
)
// ProvideJetStream creates a NATS JetStream connection
func ProvideJetStream(cfg *config.Config, logger *zap.Logger) (nats.JetStreamContext, error) {
// Connect to NATS
// ProvideNATSConn creates a reusable NATS connection.
func ProvideNATSConn(cfg *config.Config, logger *zap.Logger) (*nats.Conn, error) {
nc, err := nats.Connect(
cfg.NATS.URL,
nats.RetryOnFailedConnect(true),
@@ -29,6 +29,11 @@ func ProvideJetStream(cfg *config.Config, logger *zap.Logger) (nats.JetStreamCon
return nil, fmt.Errorf("failed to connect to NATS: %w", err)
}
return nc, nil
}
// ProvideJetStream creates a NATS JetStream context using an existing connection.
func ProvideJetStream(nc *nats.Conn, cfg *config.Config, logger *zap.Logger) (nats.JetStreamContext, error) {
// Get JetStream context
js, err := nc.JetStream()
if err != nil {
@@ -43,13 +48,30 @@ func ProvideJetStream(cfg *config.Config, logger *zap.Logger) (nats.JetStreamCon
subject = "telegram.>"
}
streamLimits := cfg.NATS.StreamLimits
storage := nats.FileStorage
switch strings.ToLower(streamLimits.Storage) {
case "memory":
storage = nats.MemoryStorage
case "file":
storage = nats.FileStorage
}
discard := nats.DiscardOld
if strings.EqualFold(streamLimits.Discard, "new") {
discard = nats.DiscardNew
}
streamConfig := &nats.StreamConfig{
Name: streamName,
Subjects: []string{subject},
Retention: nats.LimitsPolicy,
MaxAge: 24 * time.Hour,
Storage: nats.FileStorage,
Replicas: 1,
MaxMsgs: streamLimits.MaxMsgs,
MaxBytes: streamLimits.MaxBytes,
MaxAge: streamLimits.MaxAge,
Discard: discard,
Storage: storage,
Replicas: streamLimits.Replicas,
}
_, err = js.AddStream(streamConfig)
+19 -2
View File
@@ -2,11 +2,13 @@ package nats
import (
"caatsm/internal/adapter"
"caatsm/internal/domain"
"caatsm/internal/infra/config"
"encoding/json"
"fmt"
"go.uber.org/zap"
"github.com/google/uuid"
"github.com/nats-io/nats.go"
"go.uber.org/zap"
)
// Publisher publishes messages to NATS JetStream
@@ -42,8 +44,23 @@ func (p *Publisher) Publish(message interface{}) error {
return fmt.Errorf("failed to marshal message: %w", err)
}
// Build JetStream message to attach dedup headers
jsMsg := nats.NewMsg(topic)
jsMsg.Data = messageBytes
switch typed := message.(type) {
case *domain.ParsedMessage:
if typed != nil && typed.Uuid != "" {
jsMsg.Header.Set("Nats-Msg-Id", typed.Uuid)
} else {
jsMsg.Header.Set("Nats-Msg-Id", uuid.NewString())
}
default:
jsMsg.Header.Set("Nats-Msg-Id", uuid.NewString())
}
// Publish to JetStream
_, err = p.js.Publish(topic, messageBytes)
_, err = p.js.PublishMsg(jsMsg)
if err != nil {
return fmt.Errorf("failed to publish message: %w", err)
}