From fd393505a61559140e1f04298436c511b18d96bd Mon Sep 17 00:00:00 2001 From: Friedrich Gonzalez <1517449+friedrichg@users.noreply.github.com> Date: Thu, 1 Oct 2026 10:26:10 -0700 Subject: [PATCH] Distributor: add experimental OTLP metrics ingestion over gRPC Add the OTLP metrics gRPC service (opentelemetry.proto.collector.metrics.v1.MetricsService/Export) to the distributor gRPC server port, so the OpenTelemetry Collector otlp exporter can push to Cortex. Enable it with -distributor.otlp.grpc-enabled. The HTTP and gRPC receivers now share convertOTLPToWriteRequest. The gRPC receiver maps distributor errors to the gRPC codes that OTLP clients use to decide if they retry, and it returns a partial success when some metrics cannot be converted. The Cortex "proto" gRPC codec replaced the pdata codec, because it is registered after it. The Cortex codec now encodes and decodes pdata messages, so OTLP gRPC requests work with either init order. Signed-off-by: Friedrich Gonzalez <1517449+friedrichg@users.noreply.github.com> --- CHANGELOG.md | 1 + docs/api/_index.md | 29 ++ docs/configuration/config-file-reference.md | 8 + docs/configuration/v1-guarantees.md | 1 + docs/guides/open-telemetry-collector.md | 29 ++ integration/e2ecortex/client.go | 6 + integration/otlp_test.go | 72 +++++ pkg/api/api.go | 7 + pkg/api/api_test.go | 29 ++ pkg/cortexpb/codec.go | 27 ++ pkg/cortexpb/codec_test.go | 34 +++ pkg/distributor/distributor.go | 2 + pkg/util/push/otlp.go | 62 ++-- pkg/util/push/otlp_grpc.go | 188 ++++++++++++ pkg/util/push/otlp_grpc_test.go | 306 ++++++++++++++++++++ pkg/util/push/push.go | 9 +- schemas/cortex-config-schema.json | 6 + 17 files changed, 786 insertions(+), 30 deletions(-) create mode 100644 pkg/util/push/otlp_grpc.go create mode 100644 pkg/util/push/otlp_grpc_test.go 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"