Skip to content

Commit 1d99f4d

Browse files
authored
feat(distributor): Add MaxRecvMsgSize config for uncompressed message size limits (#16915)
1 parent ab8ab01 commit 1d99f4d

12 files changed

Lines changed: 96 additions & 29 deletions

File tree

‎clients/pkg/promtail/scrapeconfig/scrapeconfig.go‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -456,6 +456,9 @@ type PushTargetConfig struct {
456456

457457
// If promtail should maintain the incoming log timestamp or replace it with the current time.
458458
KeepTimestamp bool `yaml:"use_incoming_timestamp"`
459+
460+
// MaxSendMsgSize is the maximum size of the sent message.
461+
MaxSendMsgSize int `yaml:"max_send_msg_size"`
459462
}
460463

461464
// DefaultScrapeConfig is the default Config.

‎clients/pkg/promtail/targets/lokipush/pushtarget.go‎

Lines changed: 6 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -93,6 +93,11 @@ func (t *PushTarget) run() error {
9393
return err
9494
}
9595

96+
// Default to 100MB if not set to align with the default value in the distributor.
97+
if t.config.MaxSendMsgSize == 0 {
98+
t.config.MaxSendMsgSize = 100 << 20
99+
}
100+
96101
t.server = srv
97102
t.server.HTTP.Path("/loki/api/v1/push").Methods("POST").Handler(http.HandlerFunc(t.handleLoki))
98103
t.server.HTTP.Path("/promtail/api/v1/raw").Methods("POST").Handler(http.HandlerFunc(t.handlePlaintext))
@@ -111,7 +116,7 @@ func (t *PushTarget) run() error {
111116
func (t *PushTarget) handleLoki(w http.ResponseWriter, r *http.Request) {
112117
logger := util_log.WithContext(r.Context(), util_log.Logger)
113118
userID, _ := tenant.TenantID(r.Context())
114-
req, err := push.ParseRequest(logger, userID, r, push.EmptyLimits{}, push.ParseLokiRequest, nil, nil, false)
119+
req, err := push.ParseRequest(logger, userID, t.config.MaxSendMsgSize, r, push.EmptyLimits{}, push.ParseLokiRequest, nil, nil, false)
115120
if err != nil {
116121
level.Warn(t.logger).Log("msg", "failed to parse incoming push request", "err", err.Error())
117122
http.Error(w, err.Error(), http.StatusBadRequest)

‎docs/sources/setup/upgrade/_index.md‎

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -35,6 +35,12 @@ The output is incredibly verbose as it shows the entire internal config struct u
3535

3636
## Main / Unreleased
3737

38+
#### Distributor Max Receive Limits for uncompressed bytes
39+
40+
The next Loki release introduces a new configuration option (i.e. `-distibutor.max-recv-msg-size`) for the distributors to control the max receive size of uncompressed stream data. The new options's default value is set to `100MB`.
41+
42+
Supported clients should check the configuration options for max send message size if applicable.
43+
3844
## 3.4.0
3945

4046
### Loki 3.4.0

‎docs/sources/shared/configuration.md‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2516,6 +2516,10 @@ ring:
25162516
# CLI flag: -distributor.push-worker-count
25172517
[push_worker_count: <int> | default = 256]
25182518
2519+
# The maximum size of a received message.
2520+
# CLI flag: -distributor.max-recv-msg-size
2521+
[max_recv_msg_size: <int> | default = 104857600]
2522+
25192523
rate_store:
25202524
# The max number of concurrent requests to make to ingester stream apis
25212525
# CLI flag: -distributor.rate-store.max-request-parallelism

‎pkg/distributor/distributor.go‎

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -89,6 +89,9 @@ type Config struct {
8989
DistributorRing RingConfig `yaml:"ring,omitempty"`
9090
PushWorkerCount int `yaml:"push_worker_count"`
9191

92+
// Request parser
93+
MaxRecvMsgSize int `yaml:"max_recv_msg_size"`
94+
9295
// For testing.
9396
factory ring_client.PoolFactory `yaml:"-"`
9497

@@ -118,6 +121,7 @@ func (cfg *Config) RegisterFlags(fs *flag.FlagSet) {
118121
cfg.RateStore.RegisterFlagsWithPrefix("distributor.rate-store", fs)
119122
cfg.WriteFailuresLogging.RegisterFlagsWithPrefix("distributor.write-failures-logging", fs)
120123
cfg.TenantTopic.RegisterFlags(fs)
124+
fs.IntVar(&cfg.MaxRecvMsgSize, "distributor.max-recv-msg-size", 100<<20, "The maximum size of a received message.")
121125
fs.IntVar(&cfg.PushWorkerCount, "distributor.push-worker-count", 256, "Number of workers to push batches to ingesters.")
122126
fs.BoolVar(&cfg.KafkaEnabled, "distributor.kafka-writes-enabled", false, "Enable writes to Kafka during Push requests.")
123127
fs.BoolVar(&cfg.IngesterEnabled, "distributor.ingester-writes-enabled", true, "Enable writes to Ingesters during Push requests. Defaults to true.")

‎pkg/distributor/http.go‎

Lines changed: 30 additions & 9 deletions
Original file line numberDiff line numberDiff line change
@@ -45,9 +45,29 @@ func (d *Distributor) pushHandler(w http.ResponseWriter, r *http.Request, pushRe
4545
streamResolver := newRequestScopedStreamResolver(tenantID, d.validator.Limits, logger)
4646

4747
logPushRequestStreams := d.tenantConfigs.LogPushRequestStreams(tenantID)
48-
req, err := push.ParseRequest(logger, tenantID, r, d.validator.Limits, pushRequestParser, d.usageTracker, streamResolver, logPushRequestStreams)
48+
req, err := push.ParseRequest(logger, tenantID, d.cfg.MaxRecvMsgSize, r, d.validator.Limits, pushRequestParser, d.usageTracker, streamResolver, logPushRequestStreams)
4949
if err != nil {
50-
if !errors.Is(err, push.ErrAllLogsFiltered) {
50+
switch {
51+
case errors.Is(err, push.ErrRequestBodyTooLarge):
52+
if d.tenantConfigs.LogPushRequest(tenantID) {
53+
level.Debug(logger).Log(
54+
"msg", "push request failed",
55+
"code", http.StatusRequestEntityTooLarge,
56+
"err", err,
57+
)
58+
}
59+
d.writeFailuresManager.Log(tenantID, fmt.Errorf("couldn't decompress push request: %w", err))
60+
61+
// We count the compressed request body size here
62+
// because the request body could not be decompressed
63+
// and thus we don't know the uncompressed size.
64+
// In addition we don't add the metric label values for
65+
// `retention_hours` and `policy` because we don't know the labels.
66+
validation.DiscardedBytes.WithLabelValues(validation.RequestBodyTooLarge, tenantID).Add(float64(r.ContentLength))
67+
errorWriter(w, err.Error(), http.StatusRequestEntityTooLarge, logger)
68+
return
69+
70+
case !errors.Is(err, push.ErrAllLogsFiltered):
5171
if d.tenantConfigs.LogPushRequest(tenantID) {
5272
level.Debug(logger).Log(
5373
"msg", "push request failed",
@@ -59,15 +79,16 @@ func (d *Distributor) pushHandler(w http.ResponseWriter, r *http.Request, pushRe
5979

6080
errorWriter(w, err.Error(), http.StatusBadRequest, logger)
6181
return
62-
}
6382

64-
if d.tenantConfigs.LogPushRequest(tenantID) {
65-
level.Debug(logger).Log(
66-
"msg", "successful push request filtered all lines",
67-
)
83+
default:
84+
if d.tenantConfigs.LogPushRequest(tenantID) {
85+
level.Debug(logger).Log(
86+
"msg", "successful push request filtered all lines",
87+
)
88+
}
89+
w.WriteHeader(http.StatusNoContent)
90+
return
6891
}
69-
w.WriteHeader(http.StatusNoContent)
70-
return
7192
}
7293

7394
if logPushRequestStreams {

‎pkg/distributor/http_test.go‎

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -126,6 +126,7 @@ func (p *fakeParser) parseRequest(
126126
_ string,
127127
_ *http.Request,
128128
_ push.Limits,
129+
_ int,
129130
_ push.UsageTracker,
130131
_ push.StreamResolver,
131132
_ bool,

‎pkg/loghttp/push/otlp.go‎

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -34,11 +34,13 @@ const (
3434

3535
OTLPSeverityNumber = "severity_number"
3636
OTLPSeverityText = "severity_text"
37+
38+
messageSizeLargerErrFmt = "%w than max (%d vs %d)"
3739
)
3840

39-
func ParseOTLPRequest(userID string, r *http.Request, limits Limits, tracker UsageTracker, streamResolver StreamResolver, logPushRequestStreams bool, logger log.Logger) (*logproto.PushRequest, *Stats, error) {
41+
func ParseOTLPRequest(userID string, r *http.Request, limits Limits, maxRecvMsgSize int, tracker UsageTracker, streamResolver StreamResolver, logPushRequestStreams bool, logger log.Logger) (*logproto.PushRequest, *Stats, error) {
4042
stats := NewPushStats()
41-
otlpLogs, err := extractLogs(r, stats)
43+
otlpLogs, err := extractLogs(r, maxRecvMsgSize, stats)
4244
if err != nil {
4345
return nil, nil, err
4446
}
@@ -47,11 +49,16 @@ func ParseOTLPRequest(userID string, r *http.Request, limits Limits, tracker Usa
4749
return req, stats, nil
4850
}
4951

50-
func extractLogs(r *http.Request, pushStats *Stats) (plog.Logs, error) {
52+
func extractLogs(r *http.Request, maxRecvMsgSize int, pushStats *Stats) (plog.Logs, error) {
5153
pushStats.ContentEncoding = r.Header.Get(contentEnc)
5254
// bodySize should always reflect the compressed size of the request body
5355
bodySize := loki_util.NewSizeReader(r.Body)
5456
var body io.Reader = bodySize
57+
if maxRecvMsgSize > 0 {
58+
// Read from LimitReader with limit max+1. So if the underlying
59+
// reader is over limit, the result will be bigger than max.
60+
body = io.LimitReader(bodySize, int64(maxRecvMsgSize)+1)
61+
}
5562
if pushStats.ContentEncoding == gzipContentEncoding {
5663
r, err := gzip.NewReader(bodySize)
5764
if err != nil {
@@ -64,6 +71,9 @@ func extractLogs(r *http.Request, pushStats *Stats) (plog.Logs, error) {
6471
}
6572
buf, err := io.ReadAll(body)
6673
if err != nil {
74+
if size := bodySize.Size(); size > int64(maxRecvMsgSize) && maxRecvMsgSize > 0 {
75+
return plog.NewLogs(), fmt.Errorf(messageSizeLargerErrFmt, loki_util.ErrMessageSizeTooLarge, size, maxRecvMsgSize)
76+
}
6777
return plog.NewLogs(), err
6878
}
6979

‎pkg/loghttp/push/push.go‎

Lines changed: 12 additions & 7 deletions
Original file line numberDiff line numberDiff line change
@@ -5,7 +5,6 @@ import (
55
"compress/gzip"
66
"fmt"
77
"io"
8-
"math"
98
"mime"
109
"net/http"
1110
"strconv"
@@ -69,7 +68,10 @@ const (
6968
AggregatedMetricLabel = "__aggregated_metric__"
7069
)
7170

72-
var ErrAllLogsFiltered = errors.New("all logs lines filtered during parsing")
71+
var (
72+
ErrAllLogsFiltered = errors.New("all logs lines filtered during parsing")
73+
ErrRequestBodyTooLarge = errors.New("request body too large")
74+
)
7375

7476
type TenantsRetention interface {
7577
RetentionPeriodFor(userID string, lbs labels.Labels) time.Duration
@@ -103,7 +105,7 @@ type StreamResolver interface {
103105
}
104106

105107
type (
106-
RequestParser func(userID string, r *http.Request, limits Limits, tracker UsageTracker, streamResolver StreamResolver, logPushRequestStreams bool, logger log.Logger) (*logproto.PushRequest, *Stats, error)
108+
RequestParser func(userID string, r *http.Request, limits Limits, maxRecvMsgSize int, tracker UsageTracker, streamResolver StreamResolver, logPushRequestStreams bool, logger log.Logger) (*logproto.PushRequest, *Stats, error)
107109
RequestParserWrapper func(inner RequestParser) RequestParser
108110
ErrorWriter func(w http.ResponseWriter, errorStr string, code int, logger log.Logger)
109111
)
@@ -137,9 +139,12 @@ type Stats struct {
137139
IsAggregatedMetric bool
138140
}
139141

140-
func ParseRequest(logger log.Logger, userID string, r *http.Request, limits Limits, pushRequestParser RequestParser, tracker UsageTracker, streamResolver StreamResolver, logPushRequestStreams bool) (*logproto.PushRequest, error) {
141-
req, pushStats, err := pushRequestParser(userID, r, limits, tracker, streamResolver, logPushRequestStreams, logger)
142+
func ParseRequest(logger log.Logger, userID string, maxRecvMsgSize int, r *http.Request, limits Limits, pushRequestParser RequestParser, tracker UsageTracker, streamResolver StreamResolver, logPushRequestStreams bool) (*logproto.PushRequest, error) {
143+
req, pushStats, err := pushRequestParser(userID, r, limits, maxRecvMsgSize, tracker, streamResolver, logPushRequestStreams, logger)
142144
if err != nil && !errors.Is(err, ErrAllLogsFiltered) {
145+
if errors.Is(err, loki_util.ErrMessageSizeTooLarge) {
146+
return nil, fmt.Errorf("%w: %s", ErrRequestBodyTooLarge, err.Error())
147+
}
143148
return nil, err
144149
}
145150

@@ -203,7 +208,7 @@ func ParseRequest(logger log.Logger, userID string, r *http.Request, limits Limi
203208
return req, err
204209
}
205210

206-
func ParseLokiRequest(userID string, r *http.Request, limits Limits, tracker UsageTracker, streamResolver StreamResolver, logPushRequestStreams bool, logger log.Logger) (*logproto.PushRequest, *Stats, error) {
211+
func ParseLokiRequest(userID string, r *http.Request, limits Limits, maxRecvMsgSize int, tracker UsageTracker, streamResolver StreamResolver, logPushRequestStreams bool, logger log.Logger) (*logproto.PushRequest, *Stats, error) {
207212
// Body
208213
var body io.Reader
209214
// bodySize should always reflect the compressed size of the request body
@@ -263,7 +268,7 @@ func ParseLokiRequest(userID string, r *http.Request, limits Limits, tracker Usa
263268
default:
264269
// When no content-type header is set or when it is set to
265270
// `application/x-protobuf`: expect snappy compression.
266-
if err := util.ParseProtoReader(r.Context(), body, int(r.ContentLength), math.MaxInt32, &req, util.RawSnappy); err != nil {
271+
if err := util.ParseProtoReader(r.Context(), body, int(r.ContentLength), maxRecvMsgSize, &req, util.RawSnappy); err != nil {
267272
return nil, nil, err
268273
}
269274
}

‎pkg/loghttp/push/push_test.go‎

Lines changed: 5 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -300,6 +300,7 @@ func TestParseRequest(t *testing.T) {
300300
data, err := ParseRequest(
301301
util_log.Logger,
302302
"fake",
303+
100<<20,
303304
request,
304305
test.fakeLimits,
305306
ParseLokiRequest,
@@ -437,7 +438,7 @@ func Test_ServiceDetection(t *testing.T) {
437438

438439
limits := &fakeLimits{enabled: true, labels: []string{"foo"}}
439440
streamResolver := newMockStreamResolver("fake", limits)
440-
data, err := ParseRequest(util_log.Logger, "fake", request, limits, ParseLokiRequest, tracker, streamResolver, false)
441+
data, err := ParseRequest(util_log.Logger, "fake", 100<<20, request, limits, ParseLokiRequest, tracker, streamResolver, false)
441442

442443
require.NoError(t, err)
443444
require.Equal(t, labels.FromStrings("foo", "bar", LabelServiceName, "bar").String(), data.Streams[0].Labels)
@@ -449,7 +450,7 @@ func Test_ServiceDetection(t *testing.T) {
449450

450451
limits := &fakeLimits{enabled: true}
451452
streamResolver := newMockStreamResolver("fake", limits)
452-
data, err := ParseRequest(util_log.Logger, "fake", request, limits, ParseOTLPRequest, tracker, streamResolver, false)
453+
data, err := ParseRequest(util_log.Logger, "fake", 100<<20, request, limits, ParseOTLPRequest, tracker, streamResolver, false)
453454
require.NoError(t, err)
454455
require.Equal(t, labels.FromStrings("k8s_job_name", "bar", LabelServiceName, "bar").String(), data.Streams[0].Labels)
455456
})
@@ -464,7 +465,7 @@ func Test_ServiceDetection(t *testing.T) {
464465
indexAttributes: []string{"special"},
465466
}
466467
streamResolver := newMockStreamResolver("fake", limits)
467-
data, err := ParseRequest(util_log.Logger, "fake", request, limits, ParseOTLPRequest, tracker, streamResolver, false)
468+
data, err := ParseRequest(util_log.Logger, "fake", 100<<20, request, limits, ParseOTLPRequest, tracker, streamResolver, false)
468469
require.NoError(t, err)
469470
require.Equal(t, labels.FromStrings("special", "sauce", LabelServiceName, "sauce").String(), data.Streams[0].Labels)
470471
})
@@ -479,7 +480,7 @@ func Test_ServiceDetection(t *testing.T) {
479480
indexAttributes: []string{},
480481
}
481482
streamResolver := newMockStreamResolver("fake", limits)
482-
data, err := ParseRequest(util_log.Logger, "fake", request, limits, ParseOTLPRequest, tracker, streamResolver, false)
483+
data, err := ParseRequest(util_log.Logger, "fake", 100<<20, request, limits, ParseOTLPRequest, tracker, streamResolver, false)
483484
require.NoError(t, err)
484485
require.Equal(t, labels.FromStrings(LabelServiceName, ServiceUnknown).String(), data.Streams[0].Labels)
485486
})

0 commit comments

Comments
 (0)