File size: 2,160 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
package common

import (
	"context"
	"fmt"
	"log/slog"

	"github.com/ClickHouse/clickhouse-go/v2"
	"github.com/google/wire"

	"github.com/openmeterio/openmeter/app/config"
	"github.com/openmeterio/openmeter/openmeter/namespace"
	"github.com/openmeterio/openmeter/openmeter/progressmanager"
	"github.com/openmeterio/openmeter/openmeter/streaming"
	clickhouseconnector "github.com/openmeterio/openmeter/openmeter/streaming/clickhouse"
	streamingretry "github.com/openmeterio/openmeter/openmeter/streaming/retry"
)

var Streaming = wire.NewSet(
	NewStreamingConnector,
)

func NewStreamingConnector(
	ctx context.Context,
	conf config.AggregationConfiguration,
	clickHouse clickhouse.Conn,
	logger *slog.Logger,
	progressmanager progressmanager.Service,
	namespaceManager *namespace.Manager,
) (streaming.Connector, error) {
	var connector streaming.Connector
	var err error

	connector, err = clickhouseconnector.New(ctx, clickhouseconnector.Config{
		ClickHouse:             clickHouse,
		Database:               conf.ClickHouse.Database,
		EventsTableName:        conf.EventsTableName,
		Logger:                 logger,
		AsyncInsert:            conf.AsyncInsert,
		AsyncInsertWait:        conf.AsyncInsertWait,
		InsertQuerySettings:    conf.InsertQuerySettings,
		MeterQuerySettings:     conf.MeterQuerySettings,
		EnablePrewhere:         conf.EnablePrewhere,
		EnableDecimalPrecision: conf.EnableDecimalPrecision,
		ProgressManager:        progressmanager,
	})
	if err != nil {
		return nil, fmt.Errorf("init clickhouse connector: %w", err)
	}

	if conf.ClickHouse.Retry.Enabled {
		connector, err = streamingretry.New(streamingretry.Config{
			DownstreamConnector: connector,
			Logger:              logger,
			RetryWaitDuration:   conf.ClickHouse.Retry.RetryWaitDuration,
			MaxTries:            conf.ClickHouse.Retry.MaxTries,
			MaxDelay:            conf.ClickHouse.Retry.MaxDelay,
		})
		if err != nil {
			return nil, fmt.Errorf("init retry connector: %w", err)
		}
	}

	err = namespaceManager.RegisterHandler(connector)
	if err != nil {
		return nil, fmt.Errorf("failed to register streaming namespace handler: %w", err)
	}

	return connector, nil
}