diff --git a/CHANGELOG.md b/CHANGELOG.md index f76d25529085..0c5c7aa890a9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,8 @@ # Main (unreleased) +- [ENHANCEMENT] opentelemetry trace exporters can now be configured to support Oauth utilizing + the opentelemetry-collector-contrib oauth2clientauthextension. (@canuteson) + - [ENHANCEMENT] Strengthen readiness check for metrics instances. (@tpaschalis) # v0.23.0 (2022-01-13) diff --git a/docs/user/configuration/traces-config.md b/docs/user/configuration/traces-config.md index 83557549a591..f4d8845c5ce6 100644 --- a/docs/user/configuration/traces-config.md +++ b/docs/user/configuration/traces-config.md @@ -72,6 +72,31 @@ remote_write: # the latter take precedence. [ insecure_skip_verify: | default = false ] + # Configures opentelemetry exporters to use the OpenTelemetry auth extension `oauth2clientauthextension`. + # Can not be used in combination with `basic_auth`. + # See https://github.com/open-telemetry/opentelemetry-collector-contrib/blob/main/extension/oauth2clientauthextension/README.md + oauth2: + # Configures the TLS settings specific to the oauth2 client + # The client identifier issued to the oauth client + [client_id: ] + # The secret string associated with the oauth client + [client_secret: ] + # The resource server's token endpoint URL + [token_url: ] + # Optional, requested permissions associated with the oauth client + [scopes: []] + # Optional, specifies the timeout fetching tokens from the token_url. Default: no timeout + [timeout: ] + tls: + # Disable validation of the server certificate. + [ insecure: | default = false ] + # Path to the CA cert. For a client this verifies the server certificate. If empty uses system root CA. + [ca_file: ] + # Path to the TLS cert to use for TLS required connections + [cert_file: ] + # Path to the TLS key to use for TLS required connections + [key_file: ] + # Controls TLS settings of the exporter's client. See https://github.com/open-telemetry/opentelemetry-collector/blob/v0.21.0/config/configtls/README.md # This should be used only if `insecure` is set to false tls_config: diff --git a/go.mod b/go.mod index 39630d25fc78..5d5a5a96450c 100644 --- a/go.mod +++ b/go.mod @@ -37,6 +37,7 @@ require ( github.com/open-telemetry/opentelemetry-collector-contrib/exporter/jaegerexporter v0.40.0 github.com/open-telemetry/opentelemetry-collector-contrib/exporter/loadbalancingexporter v0.40.0 github.com/open-telemetry/opentelemetry-collector-contrib/exporter/prometheusexporter v0.40.0 + github.com/open-telemetry/opentelemetry-collector-contrib/extension/oauth2clientauthextension v0.40.0 github.com/open-telemetry/opentelemetry-collector-contrib/processor/attributesprocessor v0.40.0 github.com/open-telemetry/opentelemetry-collector-contrib/processor/spanmetricsprocessor v0.40.0 github.com/open-telemetry/opentelemetry-collector-contrib/processor/tailsamplingprocessor v0.40.0 @@ -421,5 +422,3 @@ replace github.com/jaegertracing/jaeger => github.com/jaegertracing/jaeger v1.27 // Replacement necessary for windows_exporter so that we can use gokit logging and not the old prometheus logging replace github.com/leoluk/perflib_exporter v0.1.0 => github.com/grafana/perflib_exporter v0.1.1-0.20211013152516-e37e14fb8b0a - - diff --git a/go.sum b/go.sum index d72bdaa9dae9..5a64699a32d8 100644 --- a/go.sum +++ b/go.sum @@ -214,8 +214,6 @@ github.com/StackExchange/wmi v1.2.1 h1:VIkavFPXSjcnS+O8yTq7NI32k0R5Aj+v39y29VYDO github.com/StackExchange/wmi v1.2.1/go.mod h1:rcmrprowKIVzvc+NUiLncP2uuArMWLCbu9SBzvHz7e8= github.com/VividCortex/gohistogram v1.0.0 h1:6+hBz+qvs0JOrrNhhmR7lFxo5sINxBCGXrdtl/UvroE= github.com/VividCortex/gohistogram v1.0.0/go.mod h1:Pf5mBqqDxYaXu3hDrrU+w6nw50o/4+TcAqDqk/vUH7g= -github.com/Workiva/go-datastructures v1.0.53 h1:J6Y/52yX10Xc5JjXmGtWoSSxs3mZnGSaq37xZZh7Yig= -github.com/Workiva/go-datastructures v1.0.53/go.mod h1:1yZL+zfsztete+ePzZz/Zb1/t5BnDuE2Ya2MMGhzP6A= github.com/abdullin/seq v0.0.0-20160510034733-d5467c17e7af/go.mod h1:5Jv4cbFiHJMsVxt52+i0Ha45fjshj6wxYr1r19tB9bw= github.com/aerospike/aerospike-client-go v1.27.0/go.mod h1:zj8LBEnWBDOVEIJt8LvaRvDG5ARAoa5dBeHaB472NRc= github.com/afex/hystrix-go v0.0.0-20180502004556-fa1af6a1f4f5/go.mod h1:SkGFH1ia65gfNATL8TAiHDNxPzPdmEL5uirI2Uyuz6c= @@ -1165,10 +1163,6 @@ github.com/grafana/statsd_exporter v0.18.1-0.20211118164740-8e806158da0b h1:eFIc github.com/grafana/statsd_exporter v0.18.1-0.20211118164740-8e806158da0b/go.mod h1:N4Z1+iSqc9rnxlT1N8Qn3l65Vzb5t4Uq0jpg8nxyhio= github.com/grafana/tail v0.0.0-20201004203643-7aa4e4a91f03 h1:fGgFrAraMB0BaPfYumu+iulfDXwHm+GFyHA4xEtBqI8= github.com/grafana/tail v0.0.0-20201004203643-7aa4e4a91f03/go.mod h1:GIMXMPB/lRAllP5rVDvcGif87ryO2hgD7tCtHMdHrho= -github.com/grafana/windows_exporter v0.15.1-0.20211019183116-592dfa92f9fd h1:jQ9JCvwdRW32X/LP3ezXT4EnJCirq9bb3l5svQs8j7k= -github.com/grafana/windows_exporter v0.15.1-0.20211019183116-592dfa92f9fd/go.mod h1:zWjLDqyEy3ZEy1LNlR6iUJJgYCoUDJTyUbrjeLUp3ZE= -github.com/grafana/windows_exporter v0.15.1-0.20220202204425-17b8026ed2f5 h1:nGyxfTz81TvfGAZnlKdLp02MVBd/pEfYQq8RwgQQr5Y= -github.com/grafana/windows_exporter v0.15.1-0.20220202204425-17b8026ed2f5/go.mod h1:zWjLDqyEy3ZEy1LNlR6iUJJgYCoUDJTyUbrjeLUp3ZE= github.com/grafana/windows_exporter v0.15.1-0.20220202211901-871715ba0b43 h1:gb+wDKb+9r4n3QbMzfudcHDtE9lI+kKjx+98bYphqK4= github.com/grafana/windows_exporter v0.15.1-0.20220202211901-871715ba0b43/go.mod h1:zWjLDqyEy3ZEy1LNlR6iUJJgYCoUDJTyUbrjeLUp3ZE= github.com/gregjones/httpcache v0.0.0-20180305231024-9cad4c3443a7/go.mod h1:FecbI9+v66THATjSRHfNgh1IVFe/9kFxbXtjV0ctIMA= @@ -1782,6 +1776,8 @@ github.com/open-telemetry/opentelemetry-collector-contrib/exporter/loadbalancing github.com/open-telemetry/opentelemetry-collector-contrib/exporter/loadbalancingexporter v0.40.0/go.mod h1:8gCz0iEj986dJMmKwi/tYo8pdQj/mnBRiRkemi97XGw= github.com/open-telemetry/opentelemetry-collector-contrib/exporter/prometheusexporter v0.40.0 h1:KCRIWJ8cqooPisXKhhoFPmq6Oo5nhb2X0XU/9uW6yAU= github.com/open-telemetry/opentelemetry-collector-contrib/exporter/prometheusexporter v0.40.0/go.mod h1:kbjb5xSL0+VLPOaroPGXV9/aZKfHDx1NEzRGjX55avA= +github.com/open-telemetry/opentelemetry-collector-contrib/extension/oauth2clientauthextension v0.40.0 h1:LAj9r9orM1tLKh+erijFOpmj5D7vk2RiIh73PZtY1Ko= +github.com/open-telemetry/opentelemetry-collector-contrib/extension/oauth2clientauthextension v0.40.0/go.mod h1:KWUGnn6Ud3OS6AotkxSQt0m4U8Hm3QWQdI28cevlnQ8= github.com/open-telemetry/opentelemetry-collector-contrib/internal/coreinternal v0.40.0 h1:FDoxyvSRJumeWNMMtwmyS6qz+5vDogkOAViz7EUvF6s= github.com/open-telemetry/opentelemetry-collector-contrib/internal/coreinternal v0.40.0/go.mod h1:a56dESln9qTQfvLlG4iVrDSVsLzWaFS5363i8sAxqDo= github.com/open-telemetry/opentelemetry-collector-contrib/internal/sharedcomponent v0.40.0 h1:IICNKhsUNFUw9qLmOVj4EkoTS+MYrGqAij+LgEM/f/c= diff --git a/pkg/traces/config.go b/pkg/traces/config.go index f0af4401c8ec..87848eeb366f 100644 --- a/pkg/traces/config.go +++ b/pkg/traces/config.go @@ -10,15 +10,11 @@ import ( "strings" "time" - "github.com/grafana/agent/pkg/logs" - "github.com/grafana/agent/pkg/traces/automaticloggingprocessor" - "github.com/grafana/agent/pkg/traces/noopreceiver" - "github.com/grafana/agent/pkg/traces/promsdprocessor" - "github.com/grafana/agent/pkg/traces/remotewriteexporter" - "github.com/grafana/agent/pkg/traces/servicegraphprocessor" + "github.com/mitchellh/mapstructure" "github.com/open-telemetry/opentelemetry-collector-contrib/exporter/jaegerexporter" "github.com/open-telemetry/opentelemetry-collector-contrib/exporter/loadbalancingexporter" "github.com/open-telemetry/opentelemetry-collector-contrib/exporter/prometheusexporter" + "github.com/open-telemetry/opentelemetry-collector-contrib/extension/oauth2clientauthextension" "github.com/open-telemetry/opentelemetry-collector-contrib/processor/attributesprocessor" "github.com/open-telemetry/opentelemetry-collector-contrib/processor/spanmetricsprocessor" "github.com/open-telemetry/opentelemetry-collector-contrib/processor/tailsamplingprocessor" @@ -37,6 +33,14 @@ import ( "go.opentelemetry.io/collector/processor/batchprocessor" "go.opentelemetry.io/collector/receiver/otlpreceiver" "go.uber.org/multierr" + + "github.com/grafana/agent/pkg/logs" + "github.com/grafana/agent/pkg/traces/automaticloggingprocessor" + "github.com/grafana/agent/pkg/traces/noopreceiver" + "github.com/grafana/agent/pkg/traces/promsdprocessor" + "github.com/grafana/agent/pkg/traces/remotewriteexporter" + "github.com/grafana/agent/pkg/traces/servicegraphprocessor" + "github.com/grafana/agent/pkg/util" ) const ( @@ -160,6 +164,50 @@ var DefaultRemoteWriteConfig = RemoteWriteConfig{ Format: formatOtlp, } +// TLSClientSetting configures the oauth2client extension TLS; compatible with configtls.TLSClientSetting +type TLSClientSetting struct { + CAFile string `yaml:"ca_file,omitempty"` + CertFile string `yaml:"cert_file,omitempty"` + KeyFile string `yaml:"key_file,omitempty"` + MinVersion string `yaml:"min_version,omitempty"` + MaxVersion string `yaml:"max_version,omitempty"` + Insecure bool `yaml:"insecure"` + InsecureSkipVerify bool `yaml:"insecure_skip_verify"` + ServerNameOverride string `yaml:"server_name_override,omitempty"` +} + +// OAuth2Config configures the oauth2client extension for a remote_write exporter +// compatible with oauth2clientauthextension.Config +type OAuth2Config struct { + ClientID string `yaml:"client_id"` + ClientSecret string `yaml:"client_secret"` + TokenURL string `yaml:"token_url"` + Scopes []string `yaml:"scopes,omitempty"` + TLS TLSClientSetting `yaml:"tls,omitempty"` + Timeout time.Duration `yaml:"timeout,omitempty"` +} + +// Agent uses standard YAML unmarshalling, while the oauth2clientauthextension relies on +// mapstructure without providing YAML labels. `toOtelConfig` marshals `Oauth2Config` to configuration type expected by +// the oauth2clientauthextension Extension Factory +func (c OAuth2Config) toOtelConfig() (*oauth2clientauthextension.Config, error) { + var result *oauth2clientauthextension.Config + decoderConfig := &mapstructure.DecoderConfig{ + MatchName: func(s, t string) bool { return util.CamelToSnake(s) == t }, + Result: &result, + WeaklyTypedInput: true, + DecodeHook: mapstructure.ComposeDecodeHookFunc( + mapstructure.StringToSliceHookFunc(","), + mapstructure.StringToTimeDurationHookFunc(), + ), + } + decoder, _ := mapstructure.NewDecoder(decoderConfig) + if err := decoder.Decode(c); err != nil { + return nil, err + } + return result, nil +} + // RemoteWriteConfig controls the configuration of an exporter type RemoteWriteConfig struct { Endpoint string `yaml:"endpoint,omitempty"` @@ -171,6 +219,7 @@ type RemoteWriteConfig struct { InsecureSkipVerify bool `yaml:"insecure_skip_verify,omitempty"` TLSConfig *prom_config.TLSConfig `yaml:"tls_config,omitempty"` BasicAuth *prom_config.BasicAuth `yaml:"basic_auth,omitempty"` + Oauth2 *OAuth2Config `yaml:"oauth2,omitempty"` Headers map[string]string `yaml:"headers,omitempty"` SendingQueue map[string]interface{} `yaml:"sending_queue,omitempty"` // https://github.com/open-telemetry/opentelemetry-collector/blob/7d7ae2eb34b5d387627875c498d7f43619f37ee3/exporter/exporterhelper/queued_retry.go#L30 RetryOnFailure map[string]interface{} `yaml:"retry_on_failure,omitempty"` // https://github.com/open-telemetry/opentelemetry-collector/blob/7d7ae2eb34b5d387627875c498d7f43619f37ee3/exporter/exporterhelper/queued_retry.go#L54 @@ -181,6 +230,7 @@ func (c *RemoteWriteConfig) UnmarshalYAML(unmarshal func(interface{}) error) err *c = DefaultRemoteWriteConfig type plain RemoteWriteConfig + if err := unmarshal((*plain)(c)); err != nil { return err } @@ -253,6 +303,10 @@ func exporter(rwCfg RemoteWriteConfig) (map[string]interface{}, error) { headers = rwCfg.Headers } + if rwCfg.BasicAuth != nil && rwCfg.Oauth2 != nil { + return nil, fmt.Errorf("Only one auth type may be configured per exporter (basic_auth or oauth2)") + } + if rwCfg.BasicAuth != nil { password := string(rwCfg.BasicAuth.Password) @@ -318,21 +372,21 @@ func exporter(rwCfg RemoteWriteConfig) (map[string]interface{}, error) { return exporter, nil } -func getExporterName(protocol string, format string) (string, error) { +func getExporterName(index int, protocol string, format string) (string, error) { switch format { case formatOtlp: switch protocol { case protocolGRPC: - return "otlp", nil + return fmt.Sprintf("otlp/%d", index), nil case protocolHTTP: - return "otlphttp", nil + return fmt.Sprintf("otlphttp/%d", index), nil default: return "", errors.New("unknown protocol, expected either 'http' or 'grpc'") } case formatJaeger: switch protocol { case protocolGRPC: - return "jaeger", nil + return fmt.Sprintf("jaeger/%d", index), nil default: return "", errors.New("unknown protocol, expected 'grpc'") } @@ -349,16 +403,42 @@ func (c *InstanceConfig) exporters() (map[string]interface{}, error) { if err != nil { return nil, err } - exporterName, err := getExporterName(remoteWriteConfig.Protocol, remoteWriteConfig.Format) + exporterName, err := getExporterName(i, remoteWriteConfig.Protocol, remoteWriteConfig.Format) if err != nil { return nil, err } - exporterName = fmt.Sprintf("%s/%d", exporterName, i) + if remoteWriteConfig.Oauth2 != nil { + exporter["auth"] = map[string]string{"authenticator": getAuthExtensionName(exporterName)} + } exporters[exporterName] = exporter } return exporters, nil } +func getAuthExtensionName(exporterName string) string { + return fmt.Sprintf("oauth2client/%s", strings.Replace(exporterName, "/", "", -1)) +} + +// builds oauth2clientauth extensions required to support RemoteWriteConfigurations. +func (c *InstanceConfig) extensions() (map[string]interface{}, error) { + extensions := map[string]interface{}{} + for i, remoteWriteConfig := range c.RemoteWrite { + if remoteWriteConfig.Oauth2 == nil { + continue + } + exporterName, err := getExporterName(i, remoteWriteConfig.Protocol, remoteWriteConfig.Format) + if err != nil { + return nil, err + } + oauthConfig, err := remoteWriteConfig.Oauth2.toOtelConfig() + if err != nil { + return nil, err + } + extensions[getAuthExtensionName(exporterName)] = oauthConfig + } + return extensions, nil +} + func resolver(config map[string]interface{}) (map[string]interface{}, error) { if len(config) == 0 { return nil, fmt.Errorf("must configure one resolver (dns or static)") @@ -435,6 +515,15 @@ func (c *InstanceConfig) otelConfig() (*config.Config, error) { return nil, errors.New("must have at least one configured receiver") } + extensions, err := c.extensions() + if err != nil { + return nil, err + } + extensionsNames := make([]string, 0, len(extensions)) + for name := range extensions { + extensionsNames = append(extensionsNames, name) + } + exporters, err := c.exporters() if err != nil { return nil, err @@ -603,14 +692,19 @@ func (c *InstanceConfig) otelConfig() (*config.Config, error) { receiversMap := map[string]interface{}(c.Receivers) + otelMapStructure["extensions"] = extensions otelMapStructure["exporters"] = exporters otelMapStructure["processors"] = processors otelMapStructure["receivers"] = receiversMap // pipelines - otelMapStructure["service"] = map[string]interface{}{ + serviceMap := map[string]interface{}{ "pipelines": pipelines, } + if len(extensionsNames) > 0 { + serviceMap["extensions"] = extensionsNames + } + otelMapStructure["service"] = serviceMap factories, err := tracingFactories() if err != nil { @@ -633,7 +727,9 @@ func (c *InstanceConfig) otelConfig() (*config.Config, error) { // tracingFactories() only creates the needed factories. if we decide to add support for a new // processor, exporter, receiver we need to add it here func tracingFactories() (component.Factories, error) { - extensions, err := component.MakeExtensionFactoryMap() + extensions, err := component.MakeExtensionFactoryMap( + oauth2clientauthextension.NewFactory(), + ) if err != nil { return component.Factories{}, err } diff --git a/pkg/traces/config_test.go b/pkg/traces/config_test.go index 07bd508c6b9d..0fd2928234d3 100644 --- a/pkg/traces/config_test.go +++ b/pkg/traces/config_test.go @@ -1059,6 +1059,246 @@ service: exporters: ["jaeger/0", "otlp/1"] processors: [] receivers: ["jaeger"] +`, + }, + { + name: "one exporter with oauth2 and basic auth", + cfg: ` +receivers: + jaeger: + protocols: + grpc: +remote_write: + - endpoint: example.com:12345 + basic_auth: + username: test + password: blerg + oauth2: + client_id: somecclient + client_secret: someclientsecret +`, + expectedError: true, + }, + { + name: "simple oauth2 config", + cfg: ` +receivers: + jaeger: + protocols: + grpc: +remote_write: + - endpoint: example.com:12345 + protocol: http + oauth2: + client_id: someclientid + client_secret: someclientsecret + token_url: https://example.com/oauth2/default/v1/token + scopes: ["api.metrics"] + timeout: 2s +`, + expectedConfig: ` +receivers: + jaeger: + protocols: + grpc: +extensions: + oauth2client/otlphttp0: + client_id: someclientid + client_secret: someclientsecret + token_url: https://example.com/oauth2/default/v1/token + scopes: ["api.metrics"] + timeout: 2s +exporters: + otlphttp/0: + endpoint: example.com:12345 + compression: gzip + retry_on_failure: + max_elapsed_time: 60s + auth: + authenticator: oauth2client/otlphttp0 +service: + extensions: ["oauth2client/otlphttp0"] + pipelines: + traces: + exporters: ["otlphttp/0"] + processors: [] + receivers: ["jaeger"] +`, + }, + { + name: "oauth2 TLS", + cfg: ` +receivers: + jaeger: + protocols: + grpc: +remote_write: + - endpoint: example.com:12345 + protocol: http + oauth2: + client_id: someclientid + client_secret: someclientsecret + token_url: https://example.com/oauth2/default/v1/token + scopes: ["api.metrics"] + timeout: 2s + tls: + insecure: true + ca_file: /var/lib/mycert.pem + cert_file: certfile + key_file: keyfile +`, + expectedConfig: ` +receivers: + jaeger: + protocols: + grpc: +extensions: + oauth2client/otlphttp0: + client_id: someclientid + client_secret: someclientsecret + token_url: https://example.com/oauth2/default/v1/token + scopes: ["api.metrics"] + timeout: 2s + tls: + insecure: true + ca_file: /var/lib/mycert.pem + cert_file: certfile + key_file: keyfile +exporters: + otlphttp/0: + endpoint: example.com:12345 + compression: gzip + retry_on_failure: + max_elapsed_time: 60s + auth: + authenticator: oauth2client/otlphttp0 +service: + extensions: ["oauth2client/otlphttp0"] + pipelines: + traces: + exporters: ["otlphttp/0"] + processors: [] + receivers: ["jaeger"] +`, + }, + { + name: "2 exporters different auth", + cfg: ` +receivers: + jaeger: + protocols: + grpc: +remote_write: + - endpoint: example.com:12345 + protocol: http + oauth2: + client_id: someclientid + client_secret: someclientsecret + token_url: https://example.com/oauth2/default/v1/token + scopes: ["api.metrics"] + timeout: 2s + - endpoint: example.com:12345 + protocol: grpc + oauth2: + client_id: anotherclientid + client_secret: anotherclientsecret + token_url: https://example.com/oauth2/default/v1/token + scopes: ["api.metrics"] + timeout: 2s +`, + expectedConfig: ` +receivers: + jaeger: + protocols: + grpc: +extensions: + oauth2client/otlphttp0: + client_id: someclientid + client_secret: someclientsecret + token_url: https://example.com/oauth2/default/v1/token + scopes: ["api.metrics"] + timeout: 2s + oauth2client/otlp1: + client_id: anotherclientid + client_secret: anotherclientsecret + token_url: https://example.com/oauth2/default/v1/token + scopes: ["api.metrics"] + timeout: 2s +exporters: + otlphttp/0: + endpoint: example.com:12345 + compression: gzip + retry_on_failure: + max_elapsed_time: 60s + auth: + authenticator: oauth2client/otlphttp0 + otlp/1: + endpoint: example.com:12345 + compression: gzip + retry_on_failure: + max_elapsed_time: 60s + auth: + authenticator: oauth2client/otlp1 +service: + extensions: ["oauth2client/otlphttp0", "oauth2client/otlp1"] + pipelines: + traces: + exporters: ["otlphttp/0", "otlp/1"] + processors: [] + receivers: ["jaeger"] +`, + }, + { + name: "exporter with insecure oauth", + cfg: ` +receivers: + jaeger: + protocols: + grpc: +remote_write: + - endpoint: http://example.com:12345 + insecure: true + protocol: http + oauth2: + client_id: someclientid + client_secret: someclientsecret + token_url: https://example.com/oauth2/default/v1/token + scopes: ["api.metrics"] + timeout: 2s + tls: + insecure: true +`, + expectedConfig: ` +receivers: + jaeger: + protocols: + grpc: +extensions: + oauth2client/otlphttp0: + client_id: someclientid + client_secret: someclientsecret + token_url: https://example.com/oauth2/default/v1/token + scopes: ["api.metrics"] + timeout: 2s + tls: + insecure: true +exporters: + otlphttp/0: + endpoint: http://example.com:12345 + tls: + insecure: true + compression: gzip + retry_on_failure: + max_elapsed_time: 60s + auth: + authenticator: oauth2client/otlphttp0 +service: + extensions: ["oauth2client/otlphttp0"] + pipelines: + traces: + exporters: ["otlphttp/0"] + processors: [] + receivers: ["jaeger"] `, }, } @@ -1068,7 +1308,6 @@ service: var cfg InstanceConfig err := yaml.Unmarshal([]byte(tc.cfg), &cfg) require.NoError(t, err) - // check error actualConfig, err := cfg.otelConfig() if tc.expectedError { @@ -1448,7 +1687,9 @@ func sortPipelines(cfg *config.Config) { var ( exp = tracePipeline.Exporters recv = tracePipeline.Receivers + ext = cfg.Service.Extensions ) sort.Slice(exp, func(i, j int) bool { return exp[i].String() > exp[j].String() }) sort.Slice(recv, func(i, j int) bool { return recv[i].String() > recv[j].String() }) + sort.Slice(ext, func(i, j int) bool { return ext[i].String() > ext[j].String() }) } diff --git a/pkg/traces/instance.go b/pkg/traces/instance.go index 36c43b241bc4..4940c0129100 100644 --- a/pkg/traces/instance.go +++ b/pkg/traces/instance.go @@ -6,20 +6,22 @@ import ( "sync" "time" - "github.com/grafana/agent/pkg/build" - "github.com/grafana/agent/pkg/logs" - "github.com/grafana/agent/pkg/metrics/instance" - "github.com/grafana/agent/pkg/traces/automaticloggingprocessor" - "github.com/grafana/agent/pkg/traces/contextkeys" - "github.com/grafana/agent/pkg/util" "github.com/prometheus/client_golang/prometheus" "go.opencensus.io/stats/view" "go.opentelemetry.io/collector/component" "go.opentelemetry.io/collector/config" "go.opentelemetry.io/collector/service/external/builder" + "go.opentelemetry.io/collector/service/external/extensions" "go.opentelemetry.io/otel/metric" "go.opentelemetry.io/otel/trace" "go.uber.org/zap" + + "github.com/grafana/agent/pkg/build" + "github.com/grafana/agent/pkg/logs" + "github.com/grafana/agent/pkg/metrics/instance" + "github.com/grafana/agent/pkg/traces/automaticloggingprocessor" + "github.com/grafana/agent/pkg/traces/contextkeys" + "github.com/grafana/agent/pkg/util" ) // Instance wraps the OpenTelemetry collector to enable tracing pipelines @@ -29,9 +31,10 @@ type Instance struct { logger *zap.Logger metricViews []*view.View - exporter builder.Exporters - pipelines builder.BuiltPipelines - receivers builder.Receivers + extensions extensions.Extensions + exporter builder.Exporters + pipelines builder.BuiltPipelines + receivers builder.Receivers } // NewInstance creates and starts an instance of tracing pipelines. @@ -86,6 +89,11 @@ func (i *Instance) stop() { shutdownCtx, cancel := context.WithTimeout(context.Background(), 30*time.Second) defer cancel() + err := i.extensions.NotifyPipelineNotReady() + if err != nil { + i.logger.Error("failed to notify extension of pipeline shutdown", zap.Error(err)) + } + dependencies := []struct { name string shutdown func() error @@ -117,6 +125,15 @@ func (i *Instance) stop() { return i.exporter.ShutdownAll(shutdownCtx) }, }, + { + name: "extensions", + shutdown: func() error { + if i.extensions == nil { + return nil + } + return i.extensions.ShutdownAll(shutdownCtx) + }, + }, } for _, dep := range dependencies { @@ -129,6 +146,7 @@ func (i *Instance) stop() { i.receivers = nil i.pipelines = nil i.exporter = nil + i.extensions = nil } func (i *Instance) buildAndStartPipeline(ctx context.Context, cfg InstanceConfig, logs *logs.Logs, instManager instance.Manager, reg prometheus.Registerer) error { @@ -180,14 +198,26 @@ func (i *Instance) buildAndStartPipeline(ctx context.Context, cfg InstanceConfig MeterProvider: metric.NewNoopMeterProvider(), } + // start extensions + i.extensions, err = extensions.Build(settings, appinfo, otelConfig, factories.Extensions) + if err != nil { + i.logger.Error(fmt.Sprintf("failed to build extensions:%s", err.Error())) + return fmt.Errorf("failed to create extensions builder: %w", err) + } + err = i.extensions.StartAll(ctx, i) + if err != nil { + i.logger.Error(fmt.Sprintf("failed to start extensions:%s", err.Error())) + return fmt.Errorf("failed to start extensions: %w", err) + } + // start exporter i.exporter, err = builder.BuildExporters(settings, appinfo, otelConfig, factories.Exporters) if err != nil { return fmt.Errorf("failed to create exporters builder: %w", err) } - err = i.exporter.StartAll(ctx, i) if err != nil { + i.logger.Error(fmt.Sprintf("failed to start exporter:%s", err.Error())) return fmt.Errorf("failed to start exporters: %w", err) } @@ -213,7 +243,7 @@ func (i *Instance) buildAndStartPipeline(ctx context.Context, cfg InstanceConfig return fmt.Errorf("failed to start receivers: %w", err) } - return nil + return i.extensions.NotifyPipelineReady() } // ReportFatalError implements component.Host @@ -228,7 +258,7 @@ func (i *Instance) GetFactory(component.Kind, config.Type) component.Factory { // GetExtensions implements component.Host func (i *Instance) GetExtensions() map[config.ComponentID]component.Extension { - return nil + return i.extensions.ToMap() } // GetExporters implements component.Host diff --git a/pkg/util/strings.go b/pkg/util/strings.go new file mode 100644 index 000000000000..0caaadcecbd0 --- /dev/null +++ b/pkg/util/strings.go @@ -0,0 +1,15 @@ +package util + +import ( + "regexp" + "strings" +) + +// CamelToSnake is a helper function for converting CamelCase to Snake Case +func CamelToSnake(str string) string { + var matchFirstCap = regexp.MustCompile("(.)([A-Z][a-z]+)") + var matchAllCap = regexp.MustCompile("([a-z0-9])([A-Z])") + snake := matchFirstCap.ReplaceAllString(str, "${1}_${2}") + snake = matchAllCap.ReplaceAllString(snake, "${1}_${2}") + return strings.ToLower(snake) +}