File size: 3,125 Bytes
1f10f31 | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 | package common
import (
"context"
"fmt"
"log/slog"
"github.com/ThreeDotsLabs/watermill/message"
"github.com/confluentinc/confluent-kafka-go/v2/kafka"
"go.opentelemetry.io/otel/metric"
"go.opentelemetry.io/otel/trace"
"github.com/openmeterio/openmeter/app/config"
"github.com/openmeterio/openmeter/openmeter/ingest"
"github.com/openmeterio/openmeter/openmeter/ingest/ingestadapter"
"github.com/openmeterio/openmeter/openmeter/ingest/kafkaingest"
"github.com/openmeterio/openmeter/openmeter/ingest/kafkaingest/serializer"
"github.com/openmeterio/openmeter/openmeter/ingest/kafkaingest/topicresolver"
watermillkafka "github.com/openmeterio/openmeter/openmeter/watermill/driver/kafka"
pkgkafka "github.com/openmeterio/openmeter/pkg/kafka"
)
func NewKafkaIngestCollector(
config config.KafkaIngestConfiguration,
producer *kafka.Producer,
topicResolver topicresolver.Resolver,
topicProvisioner pkgkafka.TopicProvisioner,
logger *slog.Logger,
tracer trace.Tracer,
) (*kafkaingest.Collector, error) {
collector, err := kafkaingest.NewCollector(
producer,
serializer.NewJSONSerializer(),
topicResolver,
topicProvisioner,
config.Partitions,
logger,
tracer,
)
if err != nil {
return nil, fmt.Errorf("failed to initialize kafka ingest: %w", err)
}
return collector, nil
}
func NewIngestCollector(
dedupeConfig config.DedupeConfiguration,
kafkaCollector *kafkaingest.Collector,
logger *slog.Logger,
meter metric.Meter,
tracer trace.Tracer,
) (ingest.Collector, func(), error) {
collector, err := ingestadapter.WithTelemetry(kafkaCollector, meter, tracer)
if err != nil {
return nil, nil, fmt.Errorf("init kafka ingest: %w", err)
}
if dedupeConfig.Enabled {
deduplicator, err := dedupeConfig.NewDeduplicator()
if err != nil {
return nil, nil, fmt.Errorf("failed to initialize deduplicator: %w", err)
}
return ingest.DeduplicatingCollector{
Collector: collector,
Deduplicator: deduplicator,
}, func() {
collector.Close()
logger.Info("closing deduplicator")
err := deduplicator.Close()
if err != nil {
logger.Error("failed to close deduplicator", "error", err)
}
}, nil
}
// Note: closing function is called by dedupe as well
return collector, func() { collector.Close() }, nil
}
// TODO: create a separate file or package for each application instead
func NewServerPublisher(
ctx context.Context,
options watermillkafka.PublisherOptions,
logger *slog.Logger,
) (message.Publisher, func(), error) {
return NewPublisher(ctx, options, logger)
}
func ServerProvisionTopics(conf config.EventsConfiguration) []pkgkafka.TopicConfig {
var provisionTopics []pkgkafka.TopicConfig
if conf.SystemEvents.AutoProvision.Enabled {
provisionTopics = append(provisionTopics, pkgkafka.TopicConfig{
Name: conf.SystemEvents.Topic,
Partitions: conf.SystemEvents.AutoProvision.Partitions,
})
}
return provisionTopics
}
func NewIngestService(
collector ingest.Collector,
logger *slog.Logger,
) (ingest.Service, error) {
return ingest.NewService(ingest.Config{
Collector: collector,
Logger: logger,
})
}
|