Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
29 changes: 29 additions & 0 deletions docs/api/_index.md
Original file line number Diff line number Diff line change
Expand Up @@ -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` |
Expand Down Expand Up @@ -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 |

@SungJin1212 SungJin1212 Oct 2, 2026 •

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Current ServerUserHeaderInterceptor behavior emits UNKNOWN as ServerUserHeaderInterceptor before Export runs in the missing tenant case, so the client gets UNKNOWN, the same as for every other gRPC method.

I think we should change the interceptor in fakeauth.SetupAuthMiddleware to return codes.Unauthenticated for user.ErrNoOrgID, since this is the first gRPC API we expose to external clients.

| 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)._

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The authentication link points to a section that describes the X-Scope-OrgID HTTP header, not gRPC metadata.


### Distributor ring status

```
Expand Down
8 changes: 8 additions & 0 deletions docs/configuration/config-file-reference.md
Original file line number Diff line number Diff line change
Expand Up @@ -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: <boolean> | 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: <boolean> | default = false]
```

### `etcd_config`
Expand Down
1 change: 1 addition & 0 deletions docs/configuration/v1-guarantees.md
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
29 changes: 29 additions & 0 deletions docs/guides/open-telemetry-collector.md
Original file line number Diff line number Diff line change
Expand Up @@ -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: <cortex-distributor>:9095
compression: gzip
tls:
insecure: true
headers:
X-Scope-OrgId: <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.

Expand Down
6 changes: 6 additions & 0 deletions integration/e2ecortex/client.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
72 changes: 72 additions & 0 deletions integration/otlp_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand Down Expand Up @@ -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)
}
7 changes: 7 additions & 0 deletions pkg/api/api.go
Original file line number Diff line number Diff line change
Expand Up @@ -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"
Expand All @@ -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"
)
Expand Down Expand Up @@ -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")
Expand Down
29 changes: 29 additions & 0 deletions pkg/api/api_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 (
Expand Down Expand Up @@ -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)
})
}
}
27 changes: 27 additions & 0 deletions pkg/cortexpb/codec.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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)
Expand Down Expand Up @@ -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)
Expand Down
34 changes: 34 additions & 0 deletions pkg/cortexpb/codec_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
2 changes: 2 additions & 0 deletions pkg/distributor/distributor.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down
Loading