package consumer import ( "context" "crypto/tls" "crypto/x509" "fmt" "os" "time" "git.unkin.net/unkin/logarchiver/internal/config" "github.com/nats-io/nats.go" "github.com/nats-io/nats.go/jetstream" ) // Connect dials NATS as the configured user and returns the connection and a // JetStream context. Callers must Close the returned *nats.Conn. func Connect(cfg config.NATSConfig) (*nats.Conn, jetstream.JetStream, error) { opts := []nats.Option{ nats.Name("logarchiver"), nats.MaxReconnects(-1), nats.ReconnectWait(2 * time.Second), } if cfg.User != "" { opts = append(opts, nats.UserInfo(cfg.User, cfg.Password)) } if cfg.CAFile != "" { pool := x509.NewCertPool() pem, err := os.ReadFile(cfg.CAFile) if err != nil { return nil, nil, fmt.Errorf("read nats ca %s: %w", cfg.CAFile, err) } if !pool.AppendCertsFromPEM(pem) { return nil, nil, fmt.Errorf("no certs parsed from nats ca %s", cfg.CAFile) } opts = append(opts, nats.Secure(&tls.Config{RootCAs: pool, MinVersion: tls.VersionTLS12})) } nc, err := nats.Connect(cfg.URL, opts...) if err != nil { return nil, nil, fmt.Errorf("connect nats %s: %w", cfg.URL, err) } js, err := jetstream.New(nc) if err != nil { nc.Close() return nil, nil, fmt.Errorf("jetstream context: %w", err) } return nc, js, nil } // EnsureConsumer creates or updates the durable pull consumer on the stream with // the configured subject filters. Independent offsets and explicit acks give // logarchiver at-least-once delivery decoupled from the transform tier. func EnsureConsumer(ctx context.Context, js jetstream.JetStream, cfg config.NATSConfig) (jetstream.Consumer, error) { ackWait := cfg.AckWait if ackWait <= 0 { ackWait = 2 * time.Minute } consCfg := jetstream.ConsumerConfig{ Durable: cfg.Durable, Name: cfg.Durable, AckPolicy: jetstream.AckExplicitPolicy, DeliverPolicy: jetstream.DeliverAllPolicy, AckWait: ackWait, MaxDeliver: -1, ReplayPolicy: jetstream.ReplayInstantPolicy, } switch len(cfg.Subjects) { case 0: return nil, fmt.Errorf("no subject filters configured") case 1: consCfg.FilterSubject = cfg.Subjects[0] default: consCfg.FilterSubjects = cfg.Subjects } cons, err := js.CreateOrUpdateConsumer(ctx, cfg.Stream, consCfg) if err != nil { return nil, fmt.Errorf("ensure consumer %s on stream %s: %w", cfg.Durable, cfg.Stream, err) } return cons, nil }