/
sso_log_exporter.go
83 lines (70 loc) · 2.25 KB
/
sso_log_exporter.go
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
// Copyright The OpenTelemetry Authors
// SPDX-License-Identifier: Apache-2.0
package opensearchexporter // import "github.com/open-telemetry/opentelemetry-collector-contrib/exporter/opensearchexporter"
import (
"context"
"strings"
"github.com/opensearch-project/opensearch-go/v2"
"go.opentelemetry.io/collector/component"
"go.opentelemetry.io/collector/config/confighttp"
"go.opentelemetry.io/collector/exporter"
"go.opentelemetry.io/collector/pdata/plog"
)
type logExporter struct {
client *opensearch.Client
Index string
bulkAction string
model mappingModel
httpSettings confighttp.ClientConfig
telemetry component.TelemetrySettings
}
func newLogExporter(cfg *Config, set exporter.CreateSettings) (*logExporter, error) {
if err := cfg.Validate(); err != nil {
return nil, err
}
model := &encodeModel{
dedup: cfg.Dedup,
dedot: cfg.Dedot,
sso: cfg.MappingsSettings.Mode == MappingSS4O.String(),
flattenAttributes: cfg.MappingsSettings.Mode == MappingFlattenAttributes.String(),
timestampField: cfg.MappingsSettings.TimestampField,
unixTime: cfg.MappingsSettings.UnixTimestamp,
dataset: cfg.Dataset,
namespace: cfg.Namespace,
}
return &logExporter{
telemetry: set.TelemetrySettings,
Index: getIndexName(cfg.Dataset, cfg.Namespace, cfg.LogsIndex),
bulkAction: cfg.BulkAction,
httpSettings: cfg.ClientConfig,
model: model,
}, nil
}
func (l *logExporter) Start(ctx context.Context, host component.Host) error {
httpClient, err := l.httpSettings.ToClient(ctx, host, l.telemetry)
if err != nil {
return err
}
client, err := newOpenSearchClient(l.httpSettings.Endpoint, httpClient, l.telemetry.Logger)
if err != nil {
return err
}
l.client = client
return nil
}
func (l *logExporter) pushLogData(ctx context.Context, ld plog.Logs) error {
indexer := newLogBulkIndexer(l.Index, l.bulkAction, l.model)
startErr := indexer.start(l.client)
if startErr != nil {
return startErr
}
indexer.submit(ctx, ld)
indexer.close(ctx)
return indexer.joinedError()
}
func getIndexName(dataset, namespace, index string) string {
if len(index) != 0 {
return index
}
return strings.Join([]string{"ss4o_logs", dataset, namespace}, "-")
}