Skip to content
Draft
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
45 changes: 42 additions & 3 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -960,6 +960,45 @@ By default the ratelimit gRPC server binds to `0.0.0.0:8081`. To change this set
socket then set `GRPC_UDS`, e.g. `GRPC_UDS=/<dir>/ratelimit.sock` and leave
`GRPC_HOST` and `GRPC_PORT` unmodified.

`GRPC_MAX_CONCURRENT_STREAMS` limits active streams **per gRPC connection**.
It defaults to `0`, which leaves grpc-go's effectively unlimited default in
place. Valid caps are `1` through `4294967294`; grpc-go treats `4294967295` as
unlimited, so that value is rejected. A valid cap is advertised in HTTP/2
SETTINGS. A client respecting that setting waits for stream capacity or its own
deadline; it does not receive an immediate `ResourceExhausted` response from
this limit. A peer that sends streams past the advertised cap can instead receive
HTTP/2 `REFUSED_STREAM`.
Size the cap using the number of upstream connections and the desired
per-instance memory budget. This setting does not limit HTTP `/json` or the
number of connections; use `MAX_CONCURRENT_REQUESTS` for a per-instance bound
on admitted calls.

**Request admission and deadlines**

The following settings can protect the service from accumulating synchronous request handlers when its cache is slow:

1. `MAX_CONCURRENT_REQUESTS`: maximum admitted `ShouldRateLimit` calls per service instance, shared by gRPC and HTTP. Default: `0` (disabled). Excess calls fail immediately; there is no admission queue.
1. `REQUEST_TIMEOUT`: deadline budget passed to admitted calls and their cache operations. Default: `0` (use only the caller's deadline). An earlier caller deadline is preserved. This setting alone does not limit concurrency.

Both settings must be non-negative and take effect at process startup. Reloading descriptor configuration does not reset the admission limit or release occupied slots. Admission occurs after the request has been decoded, before configuration lookup, tracing attributes, or cache work. It does not limit connection counts, request sizes, descriptor counts, or work a backend starts asynchronously. Set `GRPC_MAX_CONCURRENT_STREAMS` separately to limit how many gRPC streams each connection can have open before calls reach admission.

A slot stays occupied until the synchronous cache call and handler processing return, including after caller cancellation. `REQUEST_TIMEOUT` supplies a cancellation signal; it is **not a guarantee of prompt backend cleanup**. In particular, Radix can continue draining a cancelled response after returning one call and can hold later calls behind that response even after their deadlines expire. Such later handlers retain their slots. Commands already sent to Redis may still execute. Recovery must be checked by observing completed calls and successful new requests after the backend recovers.

Overload returns gRPC `ResourceExhausted`, not an `OVER_LIMIT` quota decision. Cancellation and expiry return `Canceled` and `DeadlineExceeded` after the synchronous work returns. Global and descriptor shadow modes do not override these service errors. Envoy handles them according to its separate `failure_mode_deny` setting; check that policy before enabling admission limits, and avoid immediate retry loops.

HTTP `/json` now preserves its request context, including when both limits are disabled. Service error codes `ResourceExhausted`, `DeadlineExceeded`, and `Canceled` map to HTTP `503`, `504`, and `408` respectively. A successful quota rejection continues to use `429`.

With either request limit enabled, the service records the following metrics without per-request or per-descriptor labels. The names below are exported by the default Prometheus mapper; custom mappers need equivalent entries.

| Metric | Meaning |
| ---------------------------------------------------------------- | ------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------------ |
| `ratelimit_service_request_admission_admitted_total` | Calls admitted for processing. |
| `ratelimit_service_request_admission_rejected_total` | Calls rejected because all admission slots were occupied; excludes already-cancelled callers and quota decisions. |
| `ratelimit_service_request_admission_in_flight` | Admitted handlers whose synchronous processing and completion-metric recording have not finished. This is not a count of all outstanding Redis operations or internally retained responses. |
| `ratelimit_service_request_admission_completed_duration_seconds` | Histogram of admitted synchronous processing time up to recording the completion metric, including time retained after caller cancellation. It excludes the recording call's own wait and has no sample while cache work is still blocked. |

The corresponding StatsD prefix is `ratelimit.service.request_admission`, with suffixes `admitted`, `rejected`, `in_flight`, and `completed_duration`. The duration is emitted in milliseconds and converted to seconds by the default Prometheus mapper, independently of `PROMETHEUS_RESPONSE_TIME_AS_MILLISECONDS`.

# Request Fields

For information on the fields of a Ratelimit gRPC request please read the information
Expand Down Expand Up @@ -1343,10 +1382,10 @@ The deployment type can be specified with the `REDIS_TYPE` / `REDIS_PERSECOND_TY

### Connection Timeout

Controls the maximum duration for Redis connection establishment, read operations, and write operations.
Controls the timeout for Redis connection establishment, not command I/O. `REQUEST_TIMEOUT` supplies a request deadline subject to the cancellation and cleanup limitations described above.

1. `REDIS_TIMEOUT`: sets the timeout for Redis connection and I/O operations. Default: `10s`
1. `REDIS_PERSECOND_TIMEOUT`: sets the timeout for per-second Redis connection and I/O operations. Default: `10s`
1. `REDIS_TIMEOUT`: timeout for Redis connection establishment. Default: `10s`
1. `REDIS_PERSECOND_TIMEOUT`: timeout for per-second Redis connection establishment. Default: `10s`

### Pool On-Empty Behavior

Expand Down
117 changes: 117 additions & 0 deletions src/server/grpc_streams_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,117 @@
package server

import (
"context"
"io"
"net"
"testing"
"time"

"github.com/stretchr/testify/require"
"golang.org/x/net/http2"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/grpc/health"
healthpb "google.golang.org/grpc/health/grpc_health_v1"
"google.golang.org/grpc/status"
"google.golang.org/grpc/test/bufconn"

"github.com/envoyproxy/ratelimit/src/settings"
)

func TestGrpcMaxConcurrentStreamsAdvertised(t *testing.T) {
for _, tc := range []struct {
name string
limit uint32
}{
{name: "default unlimited"},
{name: "configured limit", limit: 16},
} {
t.Run(tc.name, func(t *testing.T) {
s := settings.Settings{GrpcMaxConcurrentStreams: tc.limit}
grpcServer := grpc.NewServer(grpcServerOptions(s)...)
listener := bufconn.Listen(1024 * 1024)
serveDone := make(chan error, 1)
go func() { serveDone <- grpcServer.Serve(listener) }()
t.Cleanup(func() {
grpcServer.Stop()
require.NoError(t, <-serveDone)
})

conn, err := listener.Dial()
require.NoError(t, err)
defer conn.Close()
require.NoError(t, conn.SetDeadline(time.Now().Add(5*time.Second)))
_, err = io.WriteString(conn, http2.ClientPreface)
require.NoError(t, err)
framer := http2.NewFramer(conn, conn)
require.NoError(t, framer.WriteSettings())
frame, err := framer.ReadFrame()
require.NoError(t, err)
serverSettings, ok := frame.(*http2.SettingsFrame)
require.True(t, ok, "first server frame must be SETTINGS")

var advertised uint32
found := false
require.NoError(t, serverSettings.ForeachSetting(func(setting http2.Setting) error {
if setting.ID == http2.SettingMaxConcurrentStreams {
advertised = setting.Val
found = true
}
return nil
}))
if tc.limit == 0 {
require.False(t, found, "zero must preserve grpc-go's default")
} else {
require.True(t, found, "server must advertise the configured cap")
require.Equal(t, tc.limit, advertised)
}
})
}
}

func TestGrpcStreamLimitWaitsUntilCapacityOrCallerDeadline(t *testing.T) {
s := settings.Settings{GrpcMaxConcurrentStreams: 1}
grpcServer := grpc.NewServer(grpcServerOptions(s)...)
healthpb.RegisterHealthServer(grpcServer, health.NewServer())
listener := bufconn.Listen(1024 * 1024)
serveDone := make(chan error, 1)
go func() { serveDone <- grpcServer.Serve(listener) }()
t.Cleanup(func() {
grpcServer.Stop()
require.NoError(t, <-serveDone)
})

clientConn, err := grpc.NewClient("passthrough:///bufnet",
grpc.WithTransportCredentials(insecure.NewCredentials()),
grpc.WithContextDialer(func(context.Context, string) (net.Conn, error) { return listener.Dial() }),
)
require.NoError(t, err)
defer clientConn.Close()
client := healthpb.NewHealthClient(clientConn)
request := &healthpb.HealthCheckRequest{Service: "ratelimit"}

firstCtx, cancelFirst := context.WithTimeout(context.Background(), 5*time.Second)
defer cancelFirst()
first, err := client.Watch(firstCtx, request)
require.NoError(t, err)
_, err = first.Recv()
require.NoError(t, err) // First Watch now holds the only active stream.

secondCtx, cancelSecond := context.WithTimeout(context.Background(), 500*time.Millisecond)
defer cancelSecond()
second, err := client.Watch(secondCtx, request)
if err == nil {
_, err = second.Recv()
}
require.Equal(t, codes.DeadlineExceeded, status.Code(err))

cancelFirst()
thirdCtx, cancelThird := context.WithTimeout(context.Background(), 3*time.Second)
defer cancelThird()
third, err := client.Watch(thirdCtx, request)
require.NoError(t, err)
_, err = third.Recv()
require.NoError(t, err)
}
47 changes: 33 additions & 14 deletions src/server/server_impl.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,8 +26,10 @@ import (
gostats "github.com/lyft/gostats"
logger "github.com/sirupsen/logrus"
"google.golang.org/grpc"
"google.golang.org/grpc/codes"
"google.golang.org/grpc/health"
healthpb "google.golang.org/grpc/health/grpc_health_v1"
"google.golang.org/grpc/status"

"github.com/envoyproxy/ratelimit/src/limiter"
"github.com/envoyproxy/ratelimit/src/settings"
Expand Down Expand Up @@ -85,7 +87,7 @@ func NewJsonHandler(svc pb.RateLimitServiceServer) func(http.ResponseWriter, *ht
return func(writer http.ResponseWriter, request *http.Request) {
var req pb.RateLimitRequest

ctx := context.Background()
ctx := request.Context()

body, err := io.ReadAll(request.Body)
if err != nil {
Expand All @@ -103,7 +105,16 @@ func NewJsonHandler(svc pb.RateLimitServiceServer) func(http.ResponseWriter, *ht
resp, err := svc.ShouldRateLimit(ctx, &req)
if err != nil {
logger.Warnf("error: %s", err.Error())
writeHttpStatus(writer, http.StatusBadRequest)
httpStatus := http.StatusBadRequest
switch status.Code(err) {
case codes.ResourceExhausted:
httpStatus = http.StatusServiceUnavailable
case codes.DeadlineExceeded:
httpStatus = http.StatusGatewayTimeout
case codes.Canceled:
httpStatus = http.StatusRequestTimeout
}
writeHttpStatus(writer, httpStatus)
return
}

Expand Down Expand Up @@ -246,18 +257,7 @@ func newServer(s settings.Settings, name string, statsManager stats.Manager, loc
ret.store.AddStatGenerator(limiter.NewLocalCacheStats(localCache, ret.scope.Scope("localcache")))
}

keepaliveOpt := grpc.KeepaliveParams(keepalive.ServerParameters{
MaxConnectionAge: s.GrpcMaxConnectionAge,
MaxConnectionAgeGrace: s.GrpcMaxConnectionAgeGrace,
})
grpcOptions := []grpc.ServerOption{
keepaliveOpt,
grpc.ChainUnaryInterceptor(
s.GrpcUnaryInterceptor, // chain otel interceptor after the input interceptor
otelgrpc.UnaryServerInterceptor(),
),
grpc.StreamInterceptor(otelgrpc.StreamServerInterceptor()),
}
grpcOptions := grpcServerOptions(s)
if s.GrpcServerUseTLS {
grpcServerTlsConfig := s.GrpcServerTlsConfig
ret.grpcCertProvider = provider.NewCertProvider(s, ret.store, s.GrpcServerTlsCert, s.GrpcServerTlsKey)
Expand Down Expand Up @@ -349,6 +349,25 @@ func newServer(s settings.Settings, name string, statsManager stats.Manager, loc
return ret
}

func grpcServerOptions(s settings.Settings) []grpc.ServerOption {
keepaliveOpt := grpc.KeepaliveParams(keepalive.ServerParameters{
MaxConnectionAge: s.GrpcMaxConnectionAge,
MaxConnectionAgeGrace: s.GrpcMaxConnectionAgeGrace,
})
grpcOptions := []grpc.ServerOption{
keepaliveOpt,
grpc.ChainUnaryInterceptor(
s.GrpcUnaryInterceptor, // chain otel interceptor after the input interceptor
otelgrpc.UnaryServerInterceptor(),
),
grpc.StreamInterceptor(otelgrpc.StreamServerInterceptor()),
}
if s.GrpcMaxConcurrentStreams > 0 {
grpcOptions = append(grpcOptions, grpc.MaxConcurrentStreams(s.GrpcMaxConcurrentStreams))
}
return grpcOptions
}

func (server *server) Stop() {
server.grpcServer.GracefulStop()
server.listenerMu.Lock()
Expand Down
26 changes: 26 additions & 0 deletions src/service/ratelimit.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,6 +10,7 @@ import (
"go.opentelemetry.io/otel"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/trace"
"google.golang.org/grpc/status"
"google.golang.org/protobuf/types/known/structpb"

"github.com/envoyproxy/ratelimit/src/settings"
Expand Down Expand Up @@ -59,6 +60,7 @@ type service struct {
globalQuotaMode bool
responseDynamicMetadataEnabled bool
useCalendarMonthRateLimit bool
requestLimits *requestLimits
}

func (this *service) SetConfig(updateEvent provider.ConfigUpdateEvent, healthyWithAtLeastOneConfigLoad bool) {
Expand Down Expand Up @@ -443,6 +445,26 @@ func (this *service) ShouldRateLimit(
ctx context.Context,
request *pb.RateLimitRequest,
) (finalResponse *pb.RateLimitResponse, finalError error) {
if this.requestLimits != nil {
var release func()
ctx, release, finalError = this.requestLimits.acquire(ctx)
if finalError != nil {
return nil, finalError
}
defer release()
// Do not start cache work if cancellation raced with admission. Check
// again on return because a backend may finish after its deadline.
if err := ctx.Err(); err != nil {
return nil, status.FromContextError(err).Err()
}
defer func() {
if err := ctx.Err(); err != nil {
finalResponse = nil
finalError = status.FromContextError(err).Err()
}
}()
}

logger.Debugf("ShouldRateLimit: %+v", request)
// Generate trace
_, span := tracer.Start(
Expand Down Expand Up @@ -493,6 +515,7 @@ func (this *service) GetCurrentConfig() (config.RateLimitConfig, bool, bool) {

func NewService(cache limiter.RateLimitCache, configProvider provider.RateLimitConfigProvider, statsManager stats.Manager,
health *server.HealthChecker, clock utils.TimeSource, shadowMode, forceStart bool, healthyWithAtLeastOneConfigLoad bool,
options ...ServiceOption,
) RateLimitServiceServer {
newService := &service{
configLock: sync.RWMutex{},
Expand All @@ -505,6 +528,9 @@ func NewService(cache limiter.RateLimitCache, configProvider provider.RateLimitC
globalQuotaMode: false,
customHeaderClock: clock,
}
for _, option := range options {
option(newService)
}

if !forceStart {
logger.Info("Waiting for initial ratelimit config update event")
Expand Down
Loading