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

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
3 changes: 3 additions & 0 deletions clients/pkg/promtail/scrapeconfig/scrapeconfig.go
Original file line number Diff line number Diff line change
Expand Up @@ -456,6 +456,9 @@ type PushTargetConfig struct {

// If promtail should maintain the incoming log timestamp or replace it with the current time.
KeepTimestamp bool `yaml:"use_incoming_timestamp"`

// MaxSendMsgSize is the maximum size of the sent message.
MaxSendMsgSize int `yaml:"max_send_msg_size"`
}

// DefaultScrapeConfig is the default Config.
Expand Down
7 changes: 6 additions & 1 deletion clients/pkg/promtail/targets/lokipush/pushtarget.go
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,11 @@ func (t *PushTarget) run() error {
return err
}

// Default to 100MB if not set to align with the default value in the distributor.
if t.config.MaxSendMsgSize == 0 {
t.config.MaxSendMsgSize = 100 << 20
}

t.server = srv
t.server.HTTP.Path("/loki/api/v1/push").Methods("POST").Handler(http.HandlerFunc(t.handleLoki))
t.server.HTTP.Path("/promtail/api/v1/raw").Methods("POST").Handler(http.HandlerFunc(t.handlePlaintext))
Expand All @@ -111,7 +116,7 @@ func (t *PushTarget) run() error {
func (t *PushTarget) handleLoki(w http.ResponseWriter, r *http.Request) {
logger := util_log.WithContext(r.Context(), util_log.Logger)
userID, _ := tenant.TenantID(r.Context())
req, err := push.ParseRequest(logger, userID, r, push.EmptyLimits{}, push.ParseLokiRequest, nil, nil, false)
req, err := push.ParseRequest(logger, userID, t.config.MaxSendMsgSize, r, push.EmptyLimits{}, push.ParseLokiRequest, nil, nil, false)
if err != nil {
level.Warn(t.logger).Log("msg", "failed to parse incoming push request", "err", err.Error())
http.Error(w, err.Error(), http.StatusBadRequest)
Expand Down
6 changes: 6 additions & 0 deletions docs/sources/setup/upgrade/_index.md
Original file line number Diff line number Diff line change
Expand Up @@ -35,6 +35,12 @@ The output is incredibly verbose as it shows the entire internal config struct u

## Main / Unreleased

#### Distributor Max Receive Limits for uncompressed bytes

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`.

Supported clients should check the configuration options for max send message size if applicable.

## 3.4.0

### Loki 3.4.0
Expand Down
4 changes: 4 additions & 0 deletions docs/sources/shared/configuration.md
Original file line number Diff line number Diff line change
Expand Up @@ -2516,6 +2516,10 @@ ring:
# CLI flag: -distributor.push-worker-count
[push_worker_count: <int> | default = 256]

# The maximum size of a received message.
# CLI flag: -distributor.max-recv-msg-size
[max_recv_msg_size: <int> | default = 104857600]

rate_store:
# The max number of concurrent requests to make to ingester stream apis
# CLI flag: -distributor.rate-store.max-request-parallelism
Expand Down
4 changes: 4 additions & 0 deletions pkg/distributor/distributor.go
Original file line number Diff line number Diff line change
Expand Up @@ -89,6 +89,9 @@ type Config struct {
DistributorRing RingConfig `yaml:"ring,omitempty"`
PushWorkerCount int `yaml:"push_worker_count"`

// Request parser
MaxRecvMsgSize int `yaml:"max_recv_msg_size"`

// For testing.
factory ring_client.PoolFactory `yaml:"-"`

Expand Down Expand Up @@ -118,6 +121,7 @@ func (cfg *Config) RegisterFlags(fs *flag.FlagSet) {
cfg.RateStore.RegisterFlagsWithPrefix("distributor.rate-store", fs)
cfg.WriteFailuresLogging.RegisterFlagsWithPrefix("distributor.write-failures-logging", fs)
cfg.TenantTopic.RegisterFlags(fs)
fs.IntVar(&cfg.MaxRecvMsgSize, "distributor.max-recv-msg-size", 100<<20, "The maximum size of a received message.")
fs.IntVar(&cfg.PushWorkerCount, "distributor.push-worker-count", 256, "Number of workers to push batches to ingesters.")
fs.BoolVar(&cfg.KafkaEnabled, "distributor.kafka-writes-enabled", false, "Enable writes to Kafka during Push requests.")
fs.BoolVar(&cfg.IngesterEnabled, "distributor.ingester-writes-enabled", true, "Enable writes to Ingesters during Push requests. Defaults to true.")
Expand Down
39 changes: 30 additions & 9 deletions pkg/distributor/http.go
Original file line number Diff line number Diff line change
Expand Up @@ -45,9 +45,29 @@ func (d *Distributor) pushHandler(w http.ResponseWriter, r *http.Request, pushRe
streamResolver := newRequestScopedStreamResolver(tenantID, d.validator.Limits, logger)

logPushRequestStreams := d.tenantConfigs.LogPushRequestStreams(tenantID)
req, err := push.ParseRequest(logger, tenantID, r, d.validator.Limits, pushRequestParser, d.usageTracker, streamResolver, logPushRequestStreams)
req, err := push.ParseRequest(logger, tenantID, d.cfg.MaxRecvMsgSize, r, d.validator.Limits, pushRequestParser, d.usageTracker, streamResolver, logPushRequestStreams)
if err != nil {
if !errors.Is(err, push.ErrAllLogsFiltered) {
switch {
case errors.Is(err, push.ErrRequestBodyTooLarge):
if d.tenantConfigs.LogPushRequest(tenantID) {
level.Debug(logger).Log(
"msg", "push request failed",
"code", http.StatusRequestEntityTooLarge,
"err", err,
)
}
d.writeFailuresManager.Log(tenantID, fmt.Errorf("couldn't decompress push request: %w", err))

// We count the compressed request body size here
// because the request body could not be decompressed
// and thus we don't know the uncompressed size.
// In addition we don't add the metric label values for
// `retention_hours` and `policy` because we don't know the labels.
validation.DiscardedBytes.WithLabelValues(validation.RequestBodyTooLarge, tenantID).Add(float64(r.ContentLength))
errorWriter(w, err.Error(), http.StatusRequestEntityTooLarge, logger)
return

case !errors.Is(err, push.ErrAllLogsFiltered):
if d.tenantConfigs.LogPushRequest(tenantID) {
level.Debug(logger).Log(
"msg", "push request failed",
Expand All @@ -59,15 +79,16 @@ func (d *Distributor) pushHandler(w http.ResponseWriter, r *http.Request, pushRe

errorWriter(w, err.Error(), http.StatusBadRequest, logger)
return
}

if d.tenantConfigs.LogPushRequest(tenantID) {
level.Debug(logger).Log(
"msg", "successful push request filtered all lines",
)
default:
if d.tenantConfigs.LogPushRequest(tenantID) {
level.Debug(logger).Log(
"msg", "successful push request filtered all lines",
)
}
w.WriteHeader(http.StatusNoContent)
return
}
w.WriteHeader(http.StatusNoContent)
return
}

if logPushRequestStreams {
Expand Down
1 change: 1 addition & 0 deletions pkg/distributor/http_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -126,6 +126,7 @@ func (p *fakeParser) parseRequest(
_ string,
_ *http.Request,
_ push.Limits,
_ int,
_ push.UsageTracker,
_ push.StreamResolver,
_ bool,
Expand Down
16 changes: 13 additions & 3 deletions pkg/loghttp/push/otlp.go
Original file line number Diff line number Diff line change
Expand Up @@ -34,11 +34,13 @@ const (

OTLPSeverityNumber = "severity_number"
OTLPSeverityText = "severity_text"

messageSizeLargerErrFmt = "%w than max (%d vs %d)"
)

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

func extractLogs(r *http.Request, pushStats *Stats) (plog.Logs, error) {
func extractLogs(r *http.Request, maxRecvMsgSize int, pushStats *Stats) (plog.Logs, error) {
pushStats.ContentEncoding = r.Header.Get(contentEnc)
// bodySize should always reflect the compressed size of the request body
bodySize := loki_util.NewSizeReader(r.Body)
var body io.Reader = bodySize
if maxRecvMsgSize > 0 {
// Read from LimitReader with limit max+1. So if the underlying
// reader is over limit, the result will be bigger than max.
body = io.LimitReader(bodySize, int64(maxRecvMsgSize)+1)
}
if pushStats.ContentEncoding == gzipContentEncoding {
r, err := gzip.NewReader(bodySize)
if err != nil {
Expand All @@ -64,6 +71,9 @@ func extractLogs(r *http.Request, pushStats *Stats) (plog.Logs, error) {
}
buf, err := io.ReadAll(body)
if err != nil {
if size := bodySize.Size(); size > int64(maxRecvMsgSize) && maxRecvMsgSize > 0 {
return plog.NewLogs(), fmt.Errorf(messageSizeLargerErrFmt, loki_util.ErrMessageSizeTooLarge, size, maxRecvMsgSize)
}
return plog.NewLogs(), err
}

Expand Down
19 changes: 12 additions & 7 deletions pkg/loghttp/push/push.go
Original file line number Diff line number Diff line change
Expand Up @@ -5,7 +5,6 @@ import (
"compress/gzip"
"fmt"
"io"
"math"
"mime"
"net/http"
"strconv"
Expand Down Expand Up @@ -69,7 +68,10 @@ const (
AggregatedMetricLabel = "__aggregated_metric__"
)

var ErrAllLogsFiltered = errors.New("all logs lines filtered during parsing")
var (
ErrAllLogsFiltered = errors.New("all logs lines filtered during parsing")
ErrRequestBodyTooLarge = errors.New("request body too large")
)

type TenantsRetention interface {
RetentionPeriodFor(userID string, lbs labels.Labels) time.Duration
Expand Down Expand Up @@ -103,7 +105,7 @@ type StreamResolver interface {
}

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

func ParseRequest(logger log.Logger, userID string, r *http.Request, limits Limits, pushRequestParser RequestParser, tracker UsageTracker, streamResolver StreamResolver, logPushRequestStreams bool) (*logproto.PushRequest, error) {
req, pushStats, err := pushRequestParser(userID, r, limits, tracker, streamResolver, logPushRequestStreams, logger)
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) {
req, pushStats, err := pushRequestParser(userID, r, limits, maxRecvMsgSize, tracker, streamResolver, logPushRequestStreams, logger)
if err != nil && !errors.Is(err, ErrAllLogsFiltered) {
if errors.Is(err, loki_util.ErrMessageSizeTooLarge) {
return nil, fmt.Errorf("%w: %s", ErrRequestBodyTooLarge, err.Error())
}
return nil, err
}

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

func ParseLokiRequest(userID string, r *http.Request, limits Limits, tracker UsageTracker, streamResolver StreamResolver, logPushRequestStreams bool, logger log.Logger) (*logproto.PushRequest, *Stats, error) {
func ParseLokiRequest(userID string, r *http.Request, limits Limits, maxRecvMsgSize int, tracker UsageTracker, streamResolver StreamResolver, logPushRequestStreams bool, logger log.Logger) (*logproto.PushRequest, *Stats, error) {
// Body
var body io.Reader
// bodySize should always reflect the compressed size of the request body
Expand Down Expand Up @@ -263,7 +268,7 @@ func ParseLokiRequest(userID string, r *http.Request, limits Limits, tracker Usa
default:
// When no content-type header is set or when it is set to
// `application/x-protobuf`: expect snappy compression.
if err := util.ParseProtoReader(r.Context(), body, int(r.ContentLength), math.MaxInt32, &req, util.RawSnappy); err != nil {
if err := util.ParseProtoReader(r.Context(), body, int(r.ContentLength), maxRecvMsgSize, &req, util.RawSnappy); err != nil {
return nil, nil, err
}
}
Expand Down
9 changes: 5 additions & 4 deletions pkg/loghttp/push/push_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -300,6 +300,7 @@ func TestParseRequest(t *testing.T) {
data, err := ParseRequest(
util_log.Logger,
"fake",
100<<20,
request,
test.fakeLimits,
ParseLokiRequest,
Expand Down Expand Up @@ -437,7 +438,7 @@ func Test_ServiceDetection(t *testing.T) {

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

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

limits := &fakeLimits{enabled: true}
streamResolver := newMockStreamResolver("fake", limits)
data, err := ParseRequest(util_log.Logger, "fake", request, limits, ParseOTLPRequest, tracker, streamResolver, false)
data, err := ParseRequest(util_log.Logger, "fake", 100<<20, request, limits, ParseOTLPRequest, tracker, streamResolver, false)
require.NoError(t, err)
require.Equal(t, labels.FromStrings("k8s_job_name", "bar", LabelServiceName, "bar").String(), data.Streams[0].Labels)
})
Expand All @@ -464,7 +465,7 @@ func Test_ServiceDetection(t *testing.T) {
indexAttributes: []string{"special"},
}
streamResolver := newMockStreamResolver("fake", limits)
data, err := ParseRequest(util_log.Logger, "fake", request, limits, ParseOTLPRequest, tracker, streamResolver, false)
data, err := ParseRequest(util_log.Logger, "fake", 100<<20, request, limits, ParseOTLPRequest, tracker, streamResolver, false)
require.NoError(t, err)
require.Equal(t, labels.FromStrings("special", "sauce", LabelServiceName, "sauce").String(), data.Streams[0].Labels)
})
Expand All @@ -479,7 +480,7 @@ func Test_ServiceDetection(t *testing.T) {
indexAttributes: []string{},
}
streamResolver := newMockStreamResolver("fake", limits)
data, err := ParseRequest(util_log.Logger, "fake", request, limits, ParseOTLPRequest, tracker, streamResolver, false)
data, err := ParseRequest(util_log.Logger, "fake", 100<<20, request, limits, ParseOTLPRequest, tracker, streamResolver, false)
require.NoError(t, err)
require.Equal(t, labels.FromStrings(LabelServiceName, ServiceUnknown).String(), data.Streams[0].Labels)
})
Expand Down
13 changes: 8 additions & 5 deletions pkg/util/http.go
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,7 @@ import (
"bytes"
"context"
"encoding/json"
"errors"
"flag"
"fmt"
"html/template"
Expand All @@ -21,7 +22,9 @@ import (
"gopkg.in/yaml.v2"
)

const messageSizeLargerErrFmt = "received message larger than max (%d vs %d)"
const messageSizeLargerErrFmt = "%w than max (%d vs %d)"

var ErrMessageSizeTooLarge = errors.New("message size too large")

const (
HTTPRateLimited = "rate_limited"
Expand Down Expand Up @@ -193,11 +196,11 @@ func ParseProtoReader(ctx context.Context, reader io.Reader, expectedSize, maxSi
func decompressRequest(reader io.Reader, expectedSize, maxSize int, compression CompressionType, sp opentracing.Span) (body []byte, err error) {
defer func() {
if err != nil && len(body) > maxSize {
err = fmt.Errorf(messageSizeLargerErrFmt, len(body), maxSize)
err = fmt.Errorf(messageSizeLargerErrFmt, ErrMessageSizeTooLarge, len(body), maxSize)
}
}()
if expectedSize > maxSize {
return nil, fmt.Errorf(messageSizeLargerErrFmt, expectedSize, maxSize)
return nil, fmt.Errorf(messageSizeLargerErrFmt, ErrMessageSizeTooLarge, expectedSize, maxSize)
}
buffer, ok := tryBufferFromReader(reader)
if ok {
Expand Down Expand Up @@ -237,7 +240,7 @@ func decompressFromReader(reader io.Reader, expectedSize, maxSize int, compressi
func decompressFromBuffer(buffer *bytes.Buffer, maxSize int, compression CompressionType, sp opentracing.Span) ([]byte, error) {
bufBytes := buffer.Bytes()
if len(bufBytes) > maxSize {
return nil, fmt.Errorf(messageSizeLargerErrFmt, len(bufBytes), maxSize)
return nil, fmt.Errorf(messageSizeLargerErrFmt, ErrMessageSizeTooLarge, len(bufBytes), maxSize)
}
switch compression {
case NoCompression:
Expand All @@ -252,7 +255,7 @@ func decompressFromBuffer(buffer *bytes.Buffer, maxSize int, compression Compres
return nil, err
}
if size > maxSize {
return nil, fmt.Errorf(messageSizeLargerErrFmt, size, maxSize)
return nil, fmt.Errorf(messageSizeLargerErrFmt, ErrMessageSizeTooLarge, size, maxSize)
}
body, err := snappy.Decode(nil, bufBytes)
if err != nil {
Expand Down
4 changes: 4 additions & 0 deletions pkg/validation/validate.go
Original file line number Diff line number Diff line change
Expand Up @@ -14,6 +14,10 @@ const (
ReasonLabel = "reason"
MissingStreamsErrorMsg = "error at least one valid stream is required for ingestion"

// RequestBodyTooLarge is a reason when decompressing the request body is too large.
RequestBodyTooLarge = "request_body_too_large"
RequestBodyTooLargeErrorMsg = "request body too large: %d bytes, limit: %d bytes"

// InvalidLabels is a reason for discarding log lines which have labels that cannot be parsed.
InvalidLabels = "invalid_labels"
MissingLabels = "missing_labels"
Expand Down