diff --git a/CHANGELOG.md b/CHANGELOG.md index 498be5787d8..b6493cfa9ab 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -3,6 +3,7 @@ ## master / unreleased * [FEATURE] Ruler: Add experimental support for federated rule groups. A rule group listing tenants in its `source_tenants` field is evaluated against those tenants while the resulting series and alerts are written to the tenant owning the rule group. Enabled with `-ruler.enable-federated-rules` (requires `-tenant-federation.enabled`), and restricted to selected tenants with `-ruler.allowed-federated-tenants` and `-ruler.disallowed-federated-tenants`. #7828 +* [FEATURE] Distributor: Add experimental OTLP metrics ingestion over gRPC (`opentelemetry.proto.collector.metrics.v1.MetricsService/Export`) on the gRPC server port. Enable it with `-distributor.otlp.grpc-enabled`. The tenant is read from the `X-Scope-OrgID` gRPC metadata, and the request size is limited by `-server.grpc-max-recv-msg-size-bytes`. #7873 ## 1.22.0 in progress * [CHANGE] Ruler: Remove the deprecated `-ruler.evaluation-delay-duration` flag and its `ruler_evaluation_delay_duration` per-tenant limit. Use `-ruler.query-offset` / `ruler_query_offset`, which no longer takes the higher of the two values. Cortex decodes the runtime config strictly, so a leftover `ruler_evaluation_delay_duration` override makes the runtime config fail to load: Cortex **exits at startup** (`module failed`, `module=runtime-config`), and on an already-running process every reload fails, pinning the last good overrides and dropping `cortex_runtime_config_last_reload_successful` to 0. Run `grep -r ruler_evaluation_delay_duration` over your runtime configs before upgrading. #7792 diff --git a/docs/api/_index.md b/docs/api/_index.md index 8e0e96b55d8..cf025506561 100644 --- a/docs/api/_index.md +++ b/docs/api/_index.md @@ -27,6 +27,7 @@ For the sake of clarity, in this document we have grouped API endpoints by servi | [Fgprof](#fgprof) | _All services_ || `GET /debug/fgprof` | | [Remote write](#remote-write) | Distributor || `POST /api/v1/push` | | [OTLP receiver](#otlp-receiver) | Distributor || `POST /api/v1/otlp/v1/metrics` | +| [OTLP receiver (gRPC)](#otlp-receiver-grpc) | Distributor || gRPC `opentelemetry.proto.collector.metrics.v1.MetricsService/Export` | | [Tenants stats](#tenants-stats) | Distributor || `GET /distributor/all_user_stats` | | [HA tracker status](#ha-tracker-status) | Distributor || `GET /distributor/ha_tracker` | | [Flush blocks](#flush-blocks) | Ingester || `GET,POST /ingester/flush` | @@ -232,6 +233,34 @@ This API endpoint accepts a HTTP POST request using [OTLP](https://opentelemetry _Requires [authentication](#authentication)._ +### OTLP Receiver (gRPC) + +``` +gRPC opentelemetry.proto.collector.metrics.v1.MetricsService/Export +``` + +Entrypoint for the OTLP Receiver over gRPC. It is experimental, and it is disabled by default. Enable it with `-distributor.otlp.grpc-enabled=true`. + +This gRPC service accepts the standard [OTLP](https://opentelemetry.io/docs/specs/otlp/) metrics export request on the distributor gRPC server port (`-server.grpc-listen-port`). The conversion to Prometheus series is the same as for the HTTP endpoint, and it uses the same `-distributor.otlp.*` flags. + +- The tenant is read from the `X-Scope-OrgID` gRPC metadata. When `-auth.enabled=true`, a request without it is rejected. +- The maximum request size is set by `-server.grpc-max-recv-msg-size-bytes` (default 4 MiB), not by `-distributor.otlp-max-recv-msg-size`. The server rejects a larger request with `RESOURCE_EXHAUSTED`. Increase the limit, or decrease the batch size in the OpenTelemetry Collector. +- When some metrics cannot be converted (for example delta temporality metrics when `-distributor.otlp.allow-delta-temporality=false`), the other metrics are ingested, and the response is a partial success with an error message. The `rejected_data_points` field is always 0, because the number of dropped data points is not known. + +Errors are returned with status codes that let OTLP clients decide if they retry: + +| Distributor result | gRPC status code | Retried by OTLP clients | +|---|---|---| +| Request deduplicated by the HA tracker | `OK` | No | +| Invalid request (HTTP 4xx) | `INVALID_ARGUMENT` | No | +| Missing or wrong tenant | `UNAUTHENTICATED` / `PERMISSION_DENIED` | No | +| Rate limited (HTTP 429) | `UNAVAILABLE` | Yes | +| Server error (HTTP 5xx) | `UNAVAILABLE` | Yes | +| Client canceled the request | `CANCELED` | No | +| Deadline exceeded | `DEADLINE_EXCEEDED` | Yes | + +_Requires [authentication](#authentication)._ + ### Distributor ring status ``` diff --git a/docs/configuration/config-file-reference.md b/docs/configuration/config-file-reference.md index 6fa991fa3c4..447b866d13a 100644 --- a/docs/configuration/config-file-reference.md +++ b/docs/configuration/config-file-reference.md @@ -3757,6 +3757,14 @@ otlp: # If true, suffixes will be added to the metrics for name normalization. # CLI flag: -distributor.otlp.add-metric-suffixes [add_metric_suffixes: | default = true] + + # EXPERIMENTAL: If true, the distributor accepts OTLP metrics over gRPC + # (opentelemetry.proto.collector.metrics.v1.MetricsService/Export) on the gRPC + # server port. The tenant is read from the X-Scope-OrgID gRPC metadata. The + # maximum request size is set by -server.grpc-max-recv-msg-size-bytes, not by + # -distributor.otlp-max-recv-msg-size. + # CLI flag: -distributor.otlp.grpc-enabled + [grpc_enabled: | default = false] ``` ### `etcd_config` diff --git a/docs/configuration/v1-guarantees.md b/docs/configuration/v1-guarantees.md index b7a84922803..9bba2b68bad 100644 --- a/docs/configuration/v1-guarantees.md +++ b/docs/configuration/v1-guarantees.md @@ -103,6 +103,7 @@ Currently experimental features are: - `alertmanager-sharding-ring.final-sleep` (duration) CLI flag - OTLP Receiver - Ingest delta temporality OTLP metrics (`-distributor.otlp.allow-delta-temporality=true`) + - Ingest OTLP metrics over gRPC (`-distributor.otlp.grpc-enabled=true`) - Persistent tokens in the Ruler Ring: - `-ruler.ring.tokens-file-path` (path) CLI flag - String interning for metrics labels diff --git a/docs/guides/open-telemetry-collector.md b/docs/guides/open-telemetry-collector.md index 0454523eaf9..56a31b39624 100644 --- a/docs/guides/open-telemetry-collector.md +++ b/docs/guides/open-telemetry-collector.md @@ -64,6 +64,35 @@ service: exporters: [otlphttp] ``` +### Push with OTLP over gRPC + +The distributor can also receive OTLP metrics over gRPC. This is experimental, and it is disabled by default. Enable it with `-distributor.otlp.grpc-enabled=true`. Then use the [otlp](https://github.com/open-telemetry/opentelemetry-collector/tree/main/exporter/otlpexporter) exporter with the distributor gRPC server port (`-server.grpc-listen-port`, default 9095): + +``` +exporters: + otlp: + endpoint: :9095 + compression: gzip + tls: + insecure: true + headers: + X-Scope-OrgId: + +... + +service: + pipelines: + metrics: + receivers: [...] + processors: [...] + exporters: [otlp] +``` + +Notes: +- The `otlp` exporter uses TLS by default. Set `tls` to match the Cortex gRPC server configuration. +- The maximum request size is set by `-server.grpc-max-recv-msg-size-bytes` (default 4 MiB). This limit applies to all gRPC traffic of the distributor. If the Collector sends larger batches, increase the limit, or decrease the batch size in the Collector. +- The gRPC port is often only reachable inside the cluster, with no authenticating gateway in front of it. The distributor trusts the `X-Scope-OrgId` metadata as it is sent. Do not expose the port to clients that you do not trust. + ## Cortex configurations for ingesting OTLP metrics You can configure OTLP-related flags in the config file. diff --git a/integration/e2ecortex/client.go b/integration/e2ecortex/client.go index 0cb584f8e29..dac67c09cfb 100644 --- a/integration/e2ecortex/client.go +++ b/integration/e2ecortex/client.go @@ -336,6 +336,12 @@ func otlpWriteRequest(name, unit string, temporality pmetric.AggregationTemporal return pmetricotlp.NewExportRequestFromMetrics(d) } +// OTLPWriteRequest builds an OTLP export request with one counter and one exemplar. Use it +// to push over OTLP gRPC, which the HTTP client does not support. +func OTLPWriteRequest(name, unit string, temporality pmetric.AggregationTemporality, labels ...prompb.Label) pmetricotlp.ExportRequest { + return otlpWriteRequest(name, unit, temporality, labels...) +} + func (c *Client) OTLPPushExemplar(name, unit string, temporality pmetric.AggregationTemporality, labels ...prompb.Label) (*http.Response, error) { data, err := otlpWriteRequest(name, unit, temporality, labels...).MarshalProto() if err != nil { diff --git a/integration/otlp_test.go b/integration/otlp_test.go index f8ac2c322f0..4395b8313a2 100644 --- a/integration/otlp_test.go +++ b/integration/otlp_test.go @@ -12,12 +12,18 @@ import ( "time" "github.com/prometheus/common/model" + "github.com/prometheus/prometheus/model/labels" "github.com/prometheus/prometheus/prompb" "github.com/prometheus/prometheus/tsdb/tsdbutil" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "github.com/thanos-io/objstore/providers/s3" "go.opentelemetry.io/collector/pdata/pmetric" + "go.opentelemetry.io/collector/pdata/pmetric/pmetricotlp" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials/insecure" + "google.golang.org/grpc/encoding/gzip" + "google.golang.org/grpc/metadata" "github.com/cortexproject/cortex/integration/e2e" e2edb "github.com/cortexproject/cortex/integration/e2e/db" @@ -378,3 +384,69 @@ func TestOTLPPushDeltaTemporality(t *testing.T) { require.True(t, ok) require.Equal(t, 1, len(vector)) } + +func TestOTLPGRPC(t *testing.T) { + s, err := e2e.NewScenario(networkName) + require.NoError(t, err) + defer s.Close() + + // Start dependencies. + minio := e2edb.NewMinio(9000, bucketName) + require.NoError(t, s.StartAndWaitReady(minio)) + + require.NoError(t, copyFileToSharedDir(s, "docs/configuration/single-process-config-blocks.yaml", cortexConfigFile)) + + flags := map[string]string{ + "-blocks-storage.s3.access-key-id": e2edb.MinioAccessKey, + "-blocks-storage.s3.secret-access-key": e2edb.MinioSecretKey, + "-blocks-storage.s3.bucket-name": bucketName, + "-blocks-storage.s3.endpoint": fmt.Sprintf("%s-minio-9000:9000", networkName), + "-blocks-storage.s3.insecure": "true", + // The tenant must come from the x-scope-orgid gRPC metadata. + "-auth.enabled": "true", + "-distributor.otlp.grpc-enabled": "true", + // alert manager + "-alertmanager.web.external-url": "http://localhost/alertmanager", + "-alertmanager-storage.backend": "local", + "-alertmanager-storage.local.path": filepath.Join(e2e.ContainerSharedDir, "alertmanager_configs"), + } + require.NoError(t, writeFileToSharedDir(s, "alertmanager_configs", []byte{})) + + cortex := e2ecortex.NewSingleBinaryWithConfigFile("cortex-1", cortexConfigFile, flags, "", 9009, 9095) + require.NoError(t, s.StartAndWaitReady(cortex)) + + // Push with the standard OTLP gRPC client, as the OpenTelemetry Collector otlp exporter does. + conn, err := grpc.NewClient(cortex.GRPCEndpoint(), + grpc.WithTransportCredentials(insecure.NewCredentials()), + grpc.WithDefaultCallOptions(grpc.UseCompressor(gzip.Name)), + ) + require.NoError(t, err) + defer conn.Close() + otlpClient := pmetricotlp.NewGRPCClient(conn) + + ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) + defer cancel() + + // A request without a tenant is rejected. + _, err = otlpClient.Export(ctx, e2ecortex.OTLPWriteRequest("series_grpc", "", pmetric.AggregationTemporalityCumulative)) + require.Error(t, err) + + resp, err := otlpClient.Export(metadata.AppendToOutgoingContext(ctx, "x-scope-orgid", "user-1"), + e2ecortex.OTLPWriteRequest("series_grpc", "", pmetric.AggregationTemporalityCumulative)) + require.NoError(t, err) + require.Empty(t, resp.PartialSuccess().ErrorMessage()) + + require.NoError(t, cortex.WaitSumMetricsWithOptions(e2e.Equals(1), []string{"cortex_distributor_push_requests_total"}, + e2e.WithLabelMatchers(labels.MustNewMatcher(labels.MatchEqual, "type", "otlp_grpc")))) + + // Read the series back for the tenant that pushed it. + c, err := e2ecortex.NewClient(cortex.HTTPEndpoint(), cortex.HTTPEndpoint(), "", "", "user-1") + require.NoError(t, err) + + result, err := c.Query("series_grpc", time.Now()) + require.NoError(t, err) + require.Equal(t, model.ValVector, result.Type()) + vector := result.(model.Vector) + require.Len(t, vector, 1) + assert.Equal(t, model.SampleValue(10), vector[0].Value) +} diff --git a/pkg/api/api.go b/pkg/api/api.go index 08da9f1a11f..b30d718429c 100644 --- a/pkg/api/api.go +++ b/pkg/api/api.go @@ -20,6 +20,7 @@ import ( "github.com/prometheus/prometheus/util/httputil" "github.com/weaveworks/common/middleware" "github.com/weaveworks/common/server" + "go.opentelemetry.io/collector/pdata/pmetric/pmetricotlp" "github.com/cortexproject/cortex/pkg/alertmanager" "github.com/cortexproject/cortex/pkg/alertmanager/alertmanagerpb" @@ -43,6 +44,7 @@ import ( "github.com/cortexproject/cortex/pkg/storegateway" "github.com/cortexproject/cortex/pkg/storegateway/storegatewaypb" "github.com/cortexproject/cortex/pkg/util/flagext" + util_log "github.com/cortexproject/cortex/pkg/util/log" "github.com/cortexproject/cortex/pkg/util/push" "github.com/cortexproject/cortex/pkg/util/validation" ) @@ -296,6 +298,11 @@ func (a *API) RegisterDistributor(d *distributor.Distributor, pushConfig distrib a.RegisterRoute("/api/v1/push", push.Handler(pushConfig.RemoteWriteV2Enabled, pushConfig.AcceptUnknownRemoteWriteContentType, pushConfig.MaxRecvMsgSize, overrides, a.sourceIPs, a.cfg.wrapDistributorPush(d), requestTotal), true, "POST") a.RegisterRoute("/api/v1/otlp/v1/metrics", push.OTLPHandler(pushConfig.OTLPMaxRecvMsgSize, overrides, pushConfig.OTLPConfig, a.sourceIPs, a.cfg.wrapDistributorPush(d), requestTotal), true, "POST") + if pushConfig.OTLPConfig.GRPCEnabled { + util_log.WarnExperimentalUse("OTLP gRPC receiver") + pmetricotlp.RegisterGRPCServer(a.server.GRPC, push.NewOTLPGRPCServer(overrides, pushConfig.OTLPConfig, a.sourceIPs, a.cfg.wrapDistributorPush(d), requestTotal)) + } + a.indexPage.AddLink(SectionAdminEndpoints, "/distributor/ring", "Distributor Ring Status") a.indexPage.AddLink(SectionAdminEndpoints, "/distributor/all_user_stats", "Usage Statistics") a.indexPage.AddLink(SectionAdminEndpoints, "/distributor/ha_tracker", "HA Tracking Status") diff --git a/pkg/api/api_test.go b/pkg/api/api_test.go index f864199ee37..ad35ba2a7e9 100644 --- a/pkg/api/api_test.go +++ b/pkg/api/api_test.go @@ -12,6 +12,8 @@ import ( "github.com/prometheus/prometheus/model/labels" "github.com/stretchr/testify/require" "github.com/weaveworks/common/server" + + "github.com/cortexproject/cortex/pkg/distributor" ) const ( @@ -212,3 +214,30 @@ func Benchmark_Compression(b *testing.B) { }) } } + +func TestRegisterDistributor_OTLPGRPCService(t *testing.T) { + const otlpMetricsService = "opentelemetry.proto.collector.metrics.v1.MetricsService" + + for _, enabled := range []bool{false, true} { + t.Run(fmt.Sprintf("grpc_enabled=%v", enabled), func(t *testing.T) { + serverCfg := server.Config{ + HTTPListenNetwork: server.DefaultNetwork, + GRPCListenNetwork: server.DefaultNetwork, + MetricsNamespace: fmt.Sprintf("otlp_grpc_%v", enabled), + } + srv, err := server.New(serverCfg) + require.NoError(t, err) + t.Cleanup(srv.Shutdown) + + api, err := New(Config{}, serverCfg, srv, &FakeLogger{}) + require.NoError(t, err) + + pushCfg := distributor.Config{} + pushCfg.OTLPConfig.GRPCEnabled = enabled + api.RegisterDistributor(&distributor.Distributor{}, pushCfg, nil, prometheus.NewRegistry()) + + _, registered := srv.GRPC.GetServiceInfo()[otlpMetricsService] + require.Equal(t, enabled, registered) + }) + } +} diff --git a/pkg/cortexpb/codec.go b/pkg/cortexpb/codec.go index 0bf037ac97a..5ff19b822c3 100644 --- a/pkg/cortexpb/codec.go +++ b/pkg/cortexpb/codec.go @@ -28,6 +28,20 @@ type GogoProtoMessage interface { MarshalToSizedBuffer(dAtA []byte) (int, error) } +// otelProtoMessage is implemented by the OpenTelemetry pdata messages (for example the +// OTLP ExportMetricsServiceRequest). They are not proto.Message values, so the codec must +// encode them with their own methods. +// +// pdata registers its own "proto" codec, which wraps the codec that exists when its init +// runs. This codec is also registered as "proto", and in the Cortex binary it is +// registered after the pdata codec, so it replaces it. Without this case, the OTLP gRPC +// receiver cannot decode or encode any request. +type otelProtoMessage interface { + SizeProto() int + MarshalProto(buf []byte) int + UnmarshalProto(buf []byte) error +} + type cortexCodec struct { noOpBufferPool mem.BufferPool defaultBufferPool mem.BufferPool @@ -40,6 +54,12 @@ func (c cortexCodec) Name() string { // Marshal is basically the same as https://github.com/grpc/grpc-go/blob/d2e836604b36400a54fbf04af495d12b38fa1e3a/encoding/proto/proto.go#L43-L67 // but it uses gogo proto methods where applicable. func (c *cortexCodec) Marshal(v any) (data mem.BufferSlice, err error) { + if m, ok := v.(otelProtoMessage); ok { + buf := make([]byte, m.SizeProto()) + n := m.MarshalProto(buf) + return mem.BufferSlice{mem.SliceBuffer(buf[:n])}, nil + } + vv := messageV2Of(v) if vv == nil { return nil, fmt.Errorf("proto: failed to marshal, message is %T, want proto.Message", v) @@ -95,6 +115,13 @@ func (c *cortexCodec) Marshal(v any) (data mem.BufferSlice, err error) { // Unmarshal Copied from https://github.com/grpc/grpc-go/blob/d2e836604b36400a54fbf04af495d12b38fa1e3a/encoding/proto/proto.go#L69-L81 // but without releasing the buffer func (c *cortexCodec) Unmarshal(data mem.BufferSlice, v any) error { + if m, ok := v.(otelProtoMessage); ok { + // Do not use a pooled buffer. The decoded message can keep references into the + // buffer, and nothing releases the buffer after the request. + buf := data.MaterializeToBuffer(c.noOpBufferPool) + return m.UnmarshalProto(buf.ReadOnlyData()) + } + vv := messageV2Of(v) if vv == nil { return fmt.Errorf("failed to unmarshal, message is %T, want proto.Message", v) diff --git a/pkg/cortexpb/codec_test.go b/pkg/cortexpb/codec_test.go index 69038687da4..335c71e260a 100644 --- a/pkg/cortexpb/codec_test.go +++ b/pkg/cortexpb/codec_test.go @@ -86,3 +86,37 @@ func TestNoopBufferWhenNotReleasableMessage(t *testing.T) { }) } } + +// fakeOTelMessage has the same encoding methods as the OpenTelemetry pdata messages, +// and it is not a proto.Message. +type fakeOTelMessage struct { + payload []byte +} + +func (m *fakeOTelMessage) SizeProto() int { return len(m.payload) } + +func (m *fakeOTelMessage) MarshalProto(buf []byte) int { return copy(buf, m.payload) } + +func (m *fakeOTelMessage) UnmarshalProto(buf []byte) error { + m.payload = append([]byte(nil), buf...) + return nil +} + +func TestCodecOTelProtoMessage(t *testing.T) { + codec := &cortexCodec{ + noOpBufferPool: &wrappedBufferPool{inner: mem.NopBufferPool{}}, + defaultBufferPool: &wrappedBufferPool{inner: mem.DefaultBufferPool()}, + } + + in := &fakeOTelMessage{payload: []byte(strings.Repeat("otlp", 5000))} + data, err := codec.Marshal(in) + require.NoError(t, err) + + out := &fakeOTelMessage{} + require.NoError(t, codec.Unmarshal(data, out)) + require.Equal(t, in.payload, out.payload) + + // The decoded message can keep references into the buffer, so the codec must not + // take the buffer from the shared pool. + require.Equal(t, 0, codec.defaultBufferPool.(*wrappedBufferPool).getCount) +} diff --git a/pkg/distributor/distributor.go b/pkg/distributor/distributor.go index a60b70bc1d2..b938247ccee 100644 --- a/pkg/distributor/distributor.go +++ b/pkg/distributor/distributor.go @@ -212,6 +212,7 @@ type OTLPConfig struct { AllowDeltaTemporality bool `yaml:"allow_delta_temporality"` EnableTypeAndUnitLabels bool `yaml:"enable_type_and_unit_labels"` AddMetricSuffixes bool `yaml:"add_metric_suffixes"` + GRPCEnabled bool `yaml:"grpc_enabled"` } // RegisterFlags adds the flags required to config this to the given FlagSet @@ -245,6 +246,7 @@ func (cfg *Config) RegisterFlags(f *flag.FlagSet) { f.BoolVar(&cfg.OTLPConfig.AllowDeltaTemporality, "distributor.otlp.allow-delta-temporality", false, "EXPERIMENTAL: If true, delta temporality otlp metrics to be ingested.") f.BoolVar(&cfg.OTLPConfig.EnableTypeAndUnitLabels, "distributor.otlp.enable-type-and-unit-labels", false, "Deprecated: Use `-distributor.enable-type-and-unit-labels` flag instead.") f.BoolVar(&cfg.OTLPConfig.AddMetricSuffixes, "distributor.otlp.add-metric-suffixes", true, "If true, suffixes will be added to the metrics for name normalization.") + f.BoolVar(&cfg.OTLPConfig.GRPCEnabled, "distributor.otlp.grpc-enabled", false, "EXPERIMENTAL: If true, the distributor accepts OTLP metrics over gRPC (opentelemetry.proto.collector.metrics.v1.MetricsService/Export) on the gRPC server port. The tenant is read from the X-Scope-OrgID gRPC metadata. The maximum request size is set by -server.grpc-max-recv-msg-size-bytes, not by -distributor.otlp-max-recv-msg-size.") } // Validate config and returns error on failure diff --git a/pkg/util/push/otlp.go b/pkg/util/push/otlp.go index 910da6d1783..adb7b324219 100644 --- a/pkg/util/push/otlp.go +++ b/pkg/util/push/otlp.go @@ -66,36 +66,13 @@ func OTLPHandler(maxRecvMsgSize int, overrides *validation.Overrides, cfg distri requestTotal.WithLabelValues(labelValueOTLP).Inc() } - prwReq := cortexpb.WriteRequest{ - Source: cortexpb.API, - Metadata: nil, - SkipLabelNameValidation: false, - } - - // otlp to prompb TimeSeries - promTsList, promMetadata, err := convertToPromTS(r.Context(), req.Metrics(), cfg, overrides, userID, logger) - if err != nil && len(promTsList) == 0 { + prwReq, err := convertOTLPToWriteRequest(r.Context(), req.Metrics(), cfg, overrides, userID, logger) + if err != nil && len(prwReq.Timeseries) == 0 { http.Error(w, err.Error(), http.StatusBadRequest) return } - // convert prompb to cortexpb TimeSeries - tsList := make([]cortexpb.PreallocTimeseries, 0, len(promTsList)) - for _, v := range promTsList { - tsList = append(tsList, cortexpb.PreallocTimeseries{TimeSeries: &cortexpb.TimeSeries{ - Labels: makeLabels(v.Labels), - Samples: makeSamples(v.Samples), - Exemplars: makeExemplars(v.Exemplars), - Histograms: makeHistograms(v.Histograms), - }}) - } - - metadata := makeMetadata(promMetadata) - - prwReq.Timeseries = tsList - prwReq.Metadata = metadata - - if _, err := push(ctx, &prwReq); err != nil { + if _, err := push(ctx, prwReq); err != nil { if errors.Is(err, context.Canceled) { err = httpgrpc.Errorf(util_api.StatusClientClosedRequest, "%s", err.Error()) } @@ -114,6 +91,39 @@ func OTLPHandler(maxRecvMsgSize int, overrides *validation.Overrides, cfg distri }) } +// convertOTLPToWriteRequest converts OTLP metrics into a Cortex WriteRequest. It is +// shared by the HTTP and the gRPC OTLP receivers. +// +// The returned error is non-nil when some metrics could not be converted. The returned +// request is never nil, and it holds every series that did convert. The caller decides +// whether a partial conversion is a failure. +func convertOTLPToWriteRequest(ctx context.Context, md pmetric.Metrics, cfg distributor.OTLPConfig, overrides *validation.Overrides, userID string, logger log.Logger) (*cortexpb.WriteRequest, error) { + prwReq := &cortexpb.WriteRequest{ + Source: cortexpb.API, + Metadata: nil, + SkipLabelNameValidation: false, + } + + // otlp to prompb TimeSeries + promTsList, promMetadata, convErr := convertToPromTS(ctx, md, cfg, overrides, userID, logger) + + // convert prompb to cortexpb TimeSeries + tsList := make([]cortexpb.PreallocTimeseries, 0, len(promTsList)) + for _, v := range promTsList { + tsList = append(tsList, cortexpb.PreallocTimeseries{TimeSeries: &cortexpb.TimeSeries{ + Labels: makeLabels(v.Labels), + Samples: makeSamples(v.Samples), + Exemplars: makeExemplars(v.Exemplars), + Histograms: makeHistograms(v.Histograms), + }}) + } + + prwReq.Timeseries = tsList + prwReq.Metadata = makeMetadata(promMetadata) + + return prwReq, convErr +} + func makeMetadata(promMetadata []prompb.MetricMetadata) []*cortexpb.MetricMetadata { metadata := make([]*cortexpb.MetricMetadata, 0, len(promMetadata)) for _, m := range promMetadata { diff --git a/pkg/util/push/otlp_grpc.go b/pkg/util/push/otlp_grpc.go new file mode 100644 index 00000000000..8d79da734e6 --- /dev/null +++ b/pkg/util/push/otlp_grpc.go @@ -0,0 +1,188 @@ +package push + +import ( + "context" + "errors" + "net/http" + + "github.com/go-kit/log" + "github.com/go-kit/log/level" + "github.com/gogo/status" + "github.com/prometheus/client_golang/prometheus" + "github.com/weaveworks/common/httpgrpc" + "github.com/weaveworks/common/middleware" + "github.com/weaveworks/common/user" + "go.opentelemetry.io/collector/pdata/pmetric/pmetricotlp" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/metadata" + "google.golang.org/grpc/peer" + + "github.com/cortexproject/cortex/pkg/distributor" + "github.com/cortexproject/cortex/pkg/util" + util_api "github.com/cortexproject/cortex/pkg/util/api" + util_log "github.com/cortexproject/cortex/pkg/util/log" + "github.com/cortexproject/cortex/pkg/util/users" + "github.com/cortexproject/cortex/pkg/util/validation" +) + +// maxPartialSuccessMessageLen limits the size of the partial success message. The +// conversion error joins one error per dropped metric, so it can be very long. +const maxPartialSuccessMessageLen = 1024 + +// OTLPGRPCServer receives OTLP metrics over gRPC, using the standard +// opentelemetry.proto.collector.metrics.v1.MetricsService/Export method. +type OTLPGRPCServer struct { + pmetricotlp.UnimplementedGRPCServer + + overrides *validation.Overrides + cfg distributor.OTLPConfig + sourceIPs *middleware.SourceIPExtractor + push Func + requestTotal *prometheus.CounterVec +} + +// NewOTLPGRPCServer makes a new OTLP gRPC server. Register it with pmetricotlp.RegisterGRPCServer. +func NewOTLPGRPCServer(overrides *validation.Overrides, cfg distributor.OTLPConfig, sourceIPs *middleware.SourceIPExtractor, push Func, requestTotal *prometheus.CounterVec) *OTLPGRPCServer { + return &OTLPGRPCServer{ + overrides: overrides, + cfg: cfg, + sourceIPs: sourceIPs, + push: push, + requestTotal: requestTotal, + } +} + +// Export implements pmetricotlp.GRPCServer. +func (s *OTLPGRPCServer) Export(ctx context.Context, req pmetricotlp.ExportRequest) (pmetricotlp.ExportResponse, error) { + resp := pmetricotlp.NewExportResponse() + + logger := util_log.WithContext(ctx, util_log.Logger) + if s.sourceIPs != nil { + source := s.sourceIPs.Get(httpRequestFromGRPCContext(ctx)) + if source != "" { + ctx = util.AddSourceIPsToOutgoingContext(ctx, source) + logger = util_log.WithSourceIPs(source, logger) + } + } + + userID, err := users.TenantID(ctx) + if err != nil { + if errors.Is(err, user.ErrNoOrgID) { + return resp, status.Error(codes.Unauthenticated, err.Error()) + } + return resp, status.Error(codes.InvalidArgument, err.Error()) + } + + if s.requestTotal != nil { + s.requestTotal.WithLabelValues(labelValueOTLPGRPC).Inc() + } + + prwReq, convErr := convertOTLPToWriteRequest(ctx, req.Metrics(), s.cfg, s.overrides, userID, logger) + // The conversion stops early when the context is done. In that case the error does + // not describe dropped metrics, so it must not become a partial success. + if ctxErr := ctx.Err(); ctxErr != nil { + return resp, toOTLPGRPCError(ctxErr, logger) + } + if convErr != nil && len(prwReq.Timeseries) == 0 { + return resp, status.Error(codes.InvalidArgument, convErr.Error()) + } + + if _, err := s.push(ctx, prwReq); err != nil { + if grpcErr := toOTLPGRPCError(err, logger); grpcErr != nil { + return resp, grpcErr + } + } + + if convErr != nil { + // RejectedDataPoints stays 0. The converter does not count the dropped data + // points, and the converted series do not map 1:1 to OTLP data points, so no + // correct count is available. + resp.PartialSuccess().SetErrorMessage(truncate(convErr.Error(), maxPartialSuccessMessageLen)) + } + + return resp, nil +} + +// toOTLPGRPCError converts an error from the distributor push into a gRPC error with a +// status code that OTLP clients use to decide if they must retry. It returns nil when the +// push must be reported as successful, for example when the HA tracker deduplicated it. +// +// The returned error keeps the original httpgrpc.HTTPResponse as its only detail, so the +// gRPC server instrumentation still records the HTTP status code (for example "429"). +func toOTLPGRPCError(err error, logger log.Logger) error { + switch { + case errors.Is(err, context.Canceled): + err = httpgrpc.Errorf(util_api.StatusClientClosedRequest, "%s", err.Error()) + case errors.Is(err, context.DeadlineExceeded): + err = httpgrpc.Errorf(http.StatusGatewayTimeout, "%s", err.Error()) + } + + httpResp, ok := httpgrpc.HTTPResponseFromError(err) + if !ok { + httpResp = &httpgrpc.HTTPResponse{Code: http.StatusInternalServerError, Body: []byte(err.Error())} + } + + httpCode := int(httpResp.GetCode()) + if httpCode/100 == 2 { + return nil + } + + if httpCode/100 == 5 { + level.Error(logger).Log("msg", "push error", "err", err) + } else if httpCode != http.StatusTooManyRequests && httpCode != util_api.StatusClientClosedRequest { + level.Warn(logger).Log("msg", "push refused", "err", err) + } + + st := status.New(otlpGRPCCode(httpCode), string(httpResp.Body)) + if withDetail, detailErr := st.WithDetails(httpResp); detailErr == nil { + st = withDetail + } + return st.Err() +} + +// otlpGRPCCode maps an HTTP status code to the gRPC code that gives the retry behavior +// the OTLP specification requires. A 429 maps to Unavailable, because ResourceExhausted +// without RetryInfo is a permanent error for OTLP clients. +func otlpGRPCCode(httpCode int) codes.Code { + switch { + case httpCode == http.StatusUnauthorized: + return codes.Unauthenticated + case httpCode == http.StatusForbidden: + return codes.PermissionDenied + case httpCode == http.StatusTooManyRequests: + return codes.Unavailable + case httpCode == util_api.StatusClientClosedRequest: + return codes.Canceled + case httpCode == http.StatusGatewayTimeout: + return codes.DeadlineExceeded + case httpCode/100 == 4: + return codes.InvalidArgument + default: + return codes.Unavailable + } +} + +// httpRequestFromGRPCContext makes an HTTP request that holds the incoming gRPC +// metadata as headers and the peer address as RemoteAddr. It lets the gRPC path use +// the same middleware.SourceIPExtractor as the HTTP path. +func httpRequestFromGRPCContext(ctx context.Context) *http.Request { + r := &http.Request{Header: http.Header{}} + if md, ok := metadata.FromIncomingContext(ctx); ok { + for k, values := range md { + for _, v := range values { + r.Header.Add(k, v) + } + } + } + if p, ok := peer.FromContext(ctx); ok && p.Addr != nil { + r.RemoteAddr = p.Addr.String() + } + return r +} + +func truncate(s string, maxLen int) string { + if len(s) <= maxLen { + return s + } + return s[:maxLen] +} diff --git a/pkg/util/push/otlp_grpc_test.go b/pkg/util/push/otlp_grpc_test.go new file mode 100644 index 00000000000..a72c452247e --- /dev/null +++ b/pkg/util/push/otlp_grpc_test.go @@ -0,0 +1,306 @@ +package push + +import ( + "context" + "errors" + "fmt" + "net" + "net/http" + "testing" + "time" + + "github.com/prometheus/client_golang/prometheus" + "github.com/prometheus/client_golang/prometheus/testutil" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "github.com/weaveworks/common/httpgrpc" + "github.com/weaveworks/common/middleware" + "go.opentelemetry.io/collector/pdata/pmetric" + "go.opentelemetry.io/collector/pdata/pmetric/pmetricotlp" + "google.golang.org/grpc" + "google.golang.org/grpc/codes" + "google.golang.org/grpc/credentials/insecure" + "google.golang.org/grpc/encoding/gzip" + "google.golang.org/grpc/metadata" + grpcstatus "google.golang.org/grpc/status" + "google.golang.org/grpc/test/bufconn" + + "github.com/cortexproject/cortex/pkg/cortexpb" + "github.com/cortexproject/cortex/pkg/distributor" + "github.com/cortexproject/cortex/pkg/querier" + "github.com/cortexproject/cortex/pkg/util" + util_api "github.com/cortexproject/cortex/pkg/util/api" + "github.com/cortexproject/cortex/pkg/util/users" + "github.com/cortexproject/cortex/pkg/util/validation" +) + +const testOTLPGRPCTenant = "user-1" + +// startOTLPGRPCServer starts an in-memory gRPC server with the OTLP receiver and the same +// tenant interceptor that Cortex uses when auth is enabled. It returns a real OTLP client. +func startOTLPGRPCServer(t *testing.T, srv *OTLPGRPCServer) pmetricotlp.GRPCClient { + t.Helper() + + listen := bufconn.Listen(1024 * 1024) + server := grpc.NewServer(grpc.UnaryInterceptor(middleware.ServerUserHeaderInterceptor)) + pmetricotlp.RegisterGRPCServer(server, srv) + go func() { _ = server.Serve(listen) }() + t.Cleanup(server.Stop) + + conn, err := grpc.NewClient("passthrough://bufnet", + grpc.WithContextDialer(func(context.Context, string) (net.Conn, error) { return listen.Dial() }), + grpc.WithTransportCredentials(insecure.NewCredentials()), + grpc.WithDefaultCallOptions(grpc.UseCompressor(gzip.Name)), + ) + require.NoError(t, err) + t.Cleanup(func() { _ = conn.Close() }) + + return pmetricotlp.NewGRPCClient(conn) +} + +func newTestOTLPGRPCServer(cfg distributor.OTLPConfig, sourceIPs *middleware.SourceIPExtractor, push Func, requestTotal *prometheus.CounterVec) *OTLPGRPCServer { + overrides := validation.NewOverrides(querier.DefaultLimitsConfig(), nil) + return NewOTLPGRPCServer(overrides, cfg, sourceIPs, push, requestTotal) +} + +func tenantContext() context.Context { + return metadata.AppendToOutgoingContext(context.Background(), "x-scope-orgid", testOTLPGRPCTenant) +} + +func newRequestTotal() *prometheus.CounterVec { + return prometheus.NewCounterVec(prometheus.CounterOpts{Name: "test_requests_total"}, []string{"type"}) +} + +// exportRequestWithSums makes a request with one sum per given temporality. +func exportRequestWithSums(temporalities ...pmetric.AggregationTemporality) pmetricotlp.ExportRequest { + md := pmetric.NewMetrics() + metrics := md.ResourceMetrics().AppendEmpty().ScopeMetrics().AppendEmpty().Metrics() + for i, temporality := range temporalities { + createOtelSum(fmt.Sprintf("test_sum_%d", i), "", temporality, time.Now()).CopyTo(metrics.AppendEmpty()) + } + return pmetricotlp.NewExportRequestFromMetrics(md) +} + +func TestOTLPGRPCServer_Success(t *testing.T) { + requestTotal := newRequestTotal() + srv := newTestOTLPGRPCServer(distributor.OTLPConfig{}, nil, verifyOTLPWriteRequestHandler(t, cortexpb.API), requestTotal) + client := startOTLPGRPCServer(t, srv) + + resp, err := client.Export(tenantContext(), generateOTLPWriteRequest()) + require.NoError(t, err) + assert.Equal(t, int64(0), resp.PartialSuccess().RejectedDataPoints()) + assert.Empty(t, resp.PartialSuccess().ErrorMessage()) + + assert.Equal(t, 1.0, testutil.ToFloat64(requestTotal.WithLabelValues(labelValueOTLPGRPC))) + assert.Equal(t, 0.0, testutil.ToFloat64(requestTotal.WithLabelValues(labelValueOTLP))) +} + +func TestOTLPGRPCServer_PushReceivesTenant(t *testing.T) { + var gotTenant string + push := func(ctx context.Context, _ *cortexpb.WriteRequest) (*cortexpb.WriteResponse, error) { + gotTenant, _ = users.TenantID(ctx) + return &cortexpb.WriteResponse{}, nil + } + client := startOTLPGRPCServer(t, newTestOTLPGRPCServer(distributor.OTLPConfig{}, nil, push, nil)) + + _, err := client.Export(tenantContext(), generateOTLPWriteRequest()) + require.NoError(t, err) + assert.Equal(t, testOTLPGRPCTenant, gotTenant) +} + +func TestOTLPGRPCServer_MissingTenant(t *testing.T) { + pushCalled := false + push := func(context.Context, *cortexpb.WriteRequest) (*cortexpb.WriteResponse, error) { + pushCalled = true + return &cortexpb.WriteResponse{}, nil + } + + t.Run("rejected by the auth interceptor", func(t *testing.T) { + client := startOTLPGRPCServer(t, newTestOTLPGRPCServer(distributor.OTLPConfig{}, nil, push, nil)) + _, err := client.Export(context.Background(), generateOTLPWriteRequest()) + require.Error(t, err) + assert.Contains(t, err.Error(), "no org id") + }) + + t.Run("rejected by the server when no interceptor sets the tenant", func(t *testing.T) { + srv := newTestOTLPGRPCServer(distributor.OTLPConfig{}, nil, push, nil) + _, err := srv.Export(context.Background(), generateOTLPWriteRequest()) + require.Error(t, err) + assert.Equal(t, codes.Unauthenticated, grpcstatus.Code(err)) + }) + + assert.False(t, pushCalled) +} + +func TestOTLPGRPCServer_DeltaTemporality(t *testing.T) { + tests := map[string]struct { + allowDelta bool + temporalities []pmetric.AggregationTemporality + expectedCode codes.Code + expectPush bool + expectPartialErrMsg bool + }{ + "delta not allowed, only delta: rejected": { + temporalities: []pmetric.AggregationTemporality{pmetric.AggregationTemporalityDelta}, + expectedCode: codes.InvalidArgument, + }, + "delta not allowed, delta and cumulative: partial success": { + temporalities: []pmetric.AggregationTemporality{pmetric.AggregationTemporalityCumulative, pmetric.AggregationTemporalityDelta}, + expectedCode: codes.OK, + expectPush: true, + expectPartialErrMsg: true, + }, + "delta allowed, delta and cumulative: success": { + allowDelta: true, + temporalities: []pmetric.AggregationTemporality{pmetric.AggregationTemporalityCumulative, pmetric.AggregationTemporalityDelta}, + expectedCode: codes.OK, + expectPush: true, + }, + } + + for name, tc := range tests { + t.Run(name, func(t *testing.T) { + pushCalled := false + push := func(_ context.Context, req *cortexpb.WriteRequest) (*cortexpb.WriteResponse, error) { + pushCalled = true + assert.NotEmpty(t, req.Timeseries) + return &cortexpb.WriteResponse{}, nil + } + cfg := distributor.OTLPConfig{AllowDeltaTemporality: tc.allowDelta} + client := startOTLPGRPCServer(t, newTestOTLPGRPCServer(cfg, nil, push, nil)) + + resp, err := client.Export(tenantContext(), exportRequestWithSums(tc.temporalities...)) + assert.Equal(t, tc.expectedCode, grpcstatus.Code(err)) + assert.Equal(t, tc.expectPush, pushCalled) + if err != nil { + return + } + + // The converter does not count dropped data points, so this must always be 0. + assert.Equal(t, int64(0), resp.PartialSuccess().RejectedDataPoints()) + if tc.expectPartialErrMsg { + assert.Contains(t, resp.PartialSuccess().ErrorMessage(), "invalid temporality and type combination") + } else { + assert.Empty(t, resp.PartialSuccess().ErrorMessage()) + } + }) + } +} + +func TestOTLPGRPCServer_PushErrors(t *testing.T) { + tests := map[string]struct { + pushErr error + expectedCode codes.Code + expectedHTTPCode int32 + }{ + "HA dedup (202) is a success": { + pushErr: httpgrpc.Errorf(http.StatusAccepted, "deduplicated"), + expectedCode: codes.OK, + }, + "400 is not retryable": { + pushErr: httpgrpc.Errorf(http.StatusBadRequest, "bad labels"), + expectedCode: codes.InvalidArgument, + expectedHTTPCode: http.StatusBadRequest, + }, + "413 is not retryable": { + pushErr: httpgrpc.Errorf(http.StatusRequestEntityTooLarge, "too large"), + expectedCode: codes.InvalidArgument, + expectedHTTPCode: http.StatusRequestEntityTooLarge, + }, + "429 is retryable": { + pushErr: httpgrpc.Errorf(http.StatusTooManyRequests, "rate limited"), + expectedCode: codes.Unavailable, + expectedHTTPCode: http.StatusTooManyRequests, + }, + "503 is retryable": { + pushErr: httpgrpc.Errorf(http.StatusServiceUnavailable, "too many inflight"), + expectedCode: codes.Unavailable, + expectedHTTPCode: http.StatusServiceUnavailable, + }, + "500 is retryable": { + pushErr: httpgrpc.Errorf(http.StatusInternalServerError, "ingester failed"), + expectedCode: codes.Unavailable, + expectedHTTPCode: http.StatusInternalServerError, + }, + "plain error is retryable": { + pushErr: errors.New("too many unhealthy instances in the ring"), + expectedCode: codes.Unavailable, + expectedHTTPCode: http.StatusInternalServerError, + }, + "context canceled": { + pushErr: context.Canceled, + expectedCode: codes.Canceled, + expectedHTTPCode: util_api.StatusClientClosedRequest, + }, + "wrapped context canceled": { + pushErr: fmt.Errorf("push failed: %w", context.Canceled), + expectedCode: codes.Canceled, + expectedHTTPCode: util_api.StatusClientClosedRequest, + }, + "deadline exceeded": { + pushErr: context.DeadlineExceeded, + expectedCode: codes.DeadlineExceeded, + expectedHTTPCode: http.StatusGatewayTimeout, + }, + } + + for name, tc := range tests { + t.Run(name, func(t *testing.T) { + push := func(context.Context, *cortexpb.WriteRequest) (*cortexpb.WriteResponse, error) { + return nil, tc.pushErr + } + client := startOTLPGRPCServer(t, newTestOTLPGRPCServer(distributor.OTLPConfig{}, nil, push, nil)) + + _, err := client.Export(tenantContext(), generateOTLPWriteRequest()) + require.Equal(t, tc.expectedCode, grpcstatus.Code(err)) + if tc.expectedCode == codes.OK { + return + } + + // The HTTP code must stay in the error detail, because the gRPC server + // instrumentation uses it for the status_code label. + httpResp, ok := httpgrpc.HTTPResponseFromError(err) + require.True(t, ok) + assert.Equal(t, tc.expectedHTTPCode, httpResp.Code) + }) + } +} + +func TestOTLPGRPCServer_SourceIPs(t *testing.T) { + sourceIPs, err := middleware.NewSourceIPs("", "") + require.NoError(t, err) + + var gotSource string + push := func(ctx context.Context, _ *cortexpb.WriteRequest) (*cortexpb.WriteResponse, error) { + gotSource = util.GetSourceIPsFromOutgoingCtx(ctx) + return &cortexpb.WriteResponse{}, nil + } + client := startOTLPGRPCServer(t, newTestOTLPGRPCServer(distributor.OTLPConfig{}, sourceIPs, push, nil)) + + ctx := metadata.AppendToOutgoingContext(tenantContext(), "x-forwarded-for", "1.2.3.4") + _, err = client.Export(ctx, generateOTLPWriteRequest()) + require.NoError(t, err) + assert.Contains(t, gotSource, "1.2.3.4") +} + +func TestOTLPGRPCCode(t *testing.T) { + tests := map[int]codes.Code{ + http.StatusBadRequest: codes.InvalidArgument, + http.StatusUnauthorized: codes.Unauthenticated, + http.StatusForbidden: codes.PermissionDenied, + http.StatusRequestEntityTooLarge: codes.InvalidArgument, + http.StatusTooManyRequests: codes.Unavailable, + util_api.StatusClientClosedRequest: codes.Canceled, + http.StatusInternalServerError: codes.Unavailable, + http.StatusServiceUnavailable: codes.Unavailable, + http.StatusGatewayTimeout: codes.DeadlineExceeded, + } + for httpCode, expected := range tests { + assert.Equal(t, expected, otlpGRPCCode(httpCode), "HTTP %d", httpCode) + } +} + +func TestTruncate(t *testing.T) { + assert.Equal(t, "abc", truncate("abc", 5)) + assert.Equal(t, "ab", truncate("abc", 2)) +} diff --git a/pkg/util/push/push.go b/pkg/util/push/push.go index 50860635ecf..8a2a177d7f3 100644 --- a/pkg/util/push/push.go +++ b/pkg/util/push/push.go @@ -38,10 +38,11 @@ const ( rw20WrittenHistogramsHeader = "X-Prometheus-Remote-Write-Histograms-Written" rw20WrittenExemplarsHeader = "X-Prometheus-Remote-Write-Exemplars-Written" - labelValuePRW1 = "prw1" - labelValuePRW2 = "prw2" - labelValueOTLP = "otlp" - labelValueUnknown = "unknown" + labelValuePRW1 = "prw1" + labelValuePRW2 = "prw2" + labelValueOTLP = "otlp" + labelValueOTLPGRPC = "otlp_grpc" + labelValueUnknown = "unknown" ) // Func defines the type of the push. It is similar to http.HandlerFunc. diff --git a/schemas/cortex-config-schema.json b/schemas/cortex-config-schema.json index 70c01e2a1c5..b4b1c83673d 100644 --- a/schemas/cortex-config-schema.json +++ b/schemas/cortex-config-schema.json @@ -4420,6 +4420,12 @@ "description": "Deprecated: Use `-distributor.enable-type-and-unit-labels` flag instead.", "type": "boolean", "x-cli-flag": "distributor.otlp.enable-type-and-unit-labels" + }, + "grpc_enabled": { + "default": false, + "description": "EXPERIMENTAL: If true, the distributor accepts OTLP metrics over gRPC (opentelemetry.proto.collector.metrics.v1.MetricsService/Export) on the gRPC server port. The tenant is read from the X-Scope-OrgID gRPC metadata. The maximum request size is set by -server.grpc-max-recv-msg-size-bytes, not by -distributor.otlp-max-recv-msg-size.", + "type": "boolean", + "x-cli-flag": "distributor.otlp.grpc-enabled" } }, "type": "object"