This is an automated email from the ASF dual-hosted git repository.
mfordjody pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/dubbo-kubernetes.git
The following commit(s) were added to refs/heads/master by this push:
new 97e861c4 feat: support HTTPRoute request retries (#981)
97e861c4 is described below
commit 97e861c4022acad6c06bf8625323c28a97434ed2
Author: mfordjody <[email protected]>
AuthorDate: Tue Jul 28 12:39:42 2026 +0800
feat: support HTTPRoute request retries (#981)
* feat: support HTTPRoute request retries
* fix: unblock security and e2e checks
* fix: tolerate restricted review gate token
---
.github/workflows/review-gate.yml | 5 +-
dubbod/discovery/cmd/app/grpc_outbound.go | 240 ++++++++++++++++++---
dubbod/discovery/cmd/app/grpc_outbound_test.go | 207 +++++++++++++++++-
dubbod/discovery/pkg/networking/grpcgen/rds.go | 43 ++++
.../discovery/pkg/networking/grpcgen/rds_test.go | 50 +++++
go.mod | 6 +-
go.sum | 12 +-
samples/app/httproute.yaml | 13 +-
tests/e2e/run.sh | 5 +-
9 files changed, 530 insertions(+), 51 deletions(-)
diff --git a/.github/workflows/review-gate.yml
b/.github/workflows/review-gate.yml
index e18c2eec..9a0b5741 100644
--- a/.github/workflows/review-gate.yml
+++ b/.github/workflows/review-gate.yml
@@ -91,7 +91,10 @@ jobs:
name: 'lgtm',
});
} catch (error) {
- if (error.status !== 404) throw error;
+ // Fork-triggered pull_request_target tokens may not be allowed
+ // to mutate labels. The commit status below is authoritative.
+ if (![403, 404].includes(error.status)) throw error;
+ core.info(`lgtm label removal skipped: ${error.message}`);
}
}
diff --git a/dubbod/discovery/cmd/app/grpc_outbound.go
b/dubbod/discovery/cmd/app/grpc_outbound.go
index 61e1e905..c789f9a7 100644
--- a/dubbod/discovery/cmd/app/grpc_outbound.go
+++ b/dubbod/discovery/cmd/app/grpc_outbound.go
@@ -51,6 +51,7 @@ import (
"google.golang.org/grpc/credentials/insecure"
"google.golang.org/protobuf/proto"
"google.golang.org/protobuf/types/known/anypb"
+ "google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/structpb"
)
@@ -101,9 +102,19 @@ type xdsRouteSnapshot struct {
Host string `json:"host"`
Port int `json:"port"`
Timeout string `json:"timeout,omitempty"`
+ Retry *xdsRetryPolicy `json:"retry,omitempty"`
Destinations []xdsDestination `json:"destinations"`
}
+type xdsRetryPolicy struct {
+ Attempts uint32 `json:"attempts"`
+ RetryOn string `json:"retryOn"`
+ StatusCodes []uint32 `json:"statusCodes,omitempty"`
+ PerTryTimeout string `json:"perTryTimeout,omitempty"`
+ Backoff string `json:"backoff,omitempty"`
+ MaxBackoff string `json:"maxBackoff,omitempty"`
+}
+
type sampleADSClient struct {
conn *grpc.ClientConn
stream
discovery.AggregatedDiscoveryService_StreamAggregatedResourcesClient
@@ -118,6 +129,7 @@ type sampleADSClient struct {
subs map[string][]string
route map[string]uint32
routeTimeout time.Duration
+ routeRetry *xdsRetryPolicy
endpoints map[string][]xdsEndpoint
clusterTLS map[string]*tlsv1.UpstreamTlsContext
updates chan struct{}
@@ -427,13 +439,14 @@ func (c *sampleADSClient) handleResponse(resp
*discovery.DiscoveryResponse) erro
return c.subscribe(v1.RouteType, routeNames)
}
case v1.RouteType:
- weights, clusters, timeout, err :=
routeWeightsFromRoutes(resp.Resources, c.path, c.requestHeaders)
+ weights, clusters, timeout, retry, err :=
routeWeightsFromRoutes(resp.Resources, c.path, c.requestHeaders)
if err != nil {
return err
}
c.mu.Lock()
c.route = weights
c.routeTimeout = timeout
+ c.routeRetry = retry
c.mu.Unlock()
c.notify()
if len(clusters) > 0 {
@@ -512,6 +525,7 @@ func (c *sampleADSClient) readySnapshot(expected
map[string]uint32) (xdsRouteSna
if c.routeTimeout > 0 {
snapshot.Timeout = c.routeTimeout.String()
}
+ snapshot.Retry = cloneRetryPolicy(c.routeRetry)
for clusterName, weight := range c.route {
if weight == 0 {
continue
@@ -570,12 +584,12 @@ func routeNamesFromListeners(resources []*anypb.Any)
([]string, error) {
return sortedUnique(out), nil
}
-func routeWeightsFromRoutes(resources []*anypb.Any, requestPath string,
requestHeaders http.Header) (map[string]uint32, []string, time.Duration, error)
{
+func routeWeightsFromRoutes(resources []*anypb.Any, requestPath string,
requestHeaders http.Header) (map[string]uint32, []string, time.Duration,
*xdsRetryPolicy, error) {
weights := map[string]uint32{}
for _, resource := range resources {
rc := &routev1.RouteConfiguration{}
if err := proto.Unmarshal(resource.Value, rc); err != nil {
- return nil, nil, 0, err
+ return nil, nil, 0, nil, err
}
for _, vh := range rc.GetVirtualHosts() {
for _, rt := range vh.GetRoutes() {
@@ -584,11 +598,11 @@ func routeWeightsFromRoutes(resources []*anypb.Any,
requestPath string, requestH
}
action := rt.GetRoute()
addRouteActionWeights(weights, action)
- return weights,
sortedWeightClusterNames(weights), routeActionTimeout(action), nil
+ return weights,
sortedWeightClusterNames(weights), routeActionTimeout(action),
routeActionRetryPolicy(action), nil
}
}
}
- return weights, sortedWeightClusterNames(weights), 0, nil
+ return weights, sortedWeightClusterNames(weights), 0, nil, nil
}
func routeActionTimeout(action *routev1.RouteAction) time.Duration {
@@ -598,6 +612,53 @@ func routeActionTimeout(action *routev1.RouteAction)
time.Duration {
return action.GetTimeout().AsDuration()
}
+func routeActionRetryPolicy(action *routev1.RouteAction) *xdsRetryPolicy {
+ if action == nil || action.GetRetryPolicy() == nil {
+ return nil
+ }
+ policy := action.GetRetryPolicy()
+ retry := &xdsRetryPolicy{
+ Attempts: 1,
+ RetryOn: policy.GetRetryOn(),
+ StatusCodes: append([]uint32(nil),
policy.GetRetriableStatusCodes()...),
+ }
+ if policy.GetNumRetries() != nil {
+ retry.Attempts = policy.GetNumRetries().GetValue()
+ }
+ if timeout := positiveProtoDuration(policy.GetPerTryTimeout()); timeout
> 0 {
+ retry.PerTryTimeout = timeout.String()
+ }
+ if backoff := policy.GetRetryBackOff(); backoff != nil {
+ if base := positiveProtoDuration(backoff.GetBaseInterval());
base > 0 {
+ retry.Backoff = base.String()
+ }
+ if maximum := positiveProtoDuration(backoff.GetMaxInterval());
maximum > 0 {
+ retry.MaxBackoff = maximum.String()
+ }
+ }
+ return retry
+}
+
+func positiveProtoDuration(value *durationpb.Duration) time.Duration {
+ if value == nil {
+ return 0
+ }
+ duration := value.AsDuration()
+ if duration <= 0 {
+ return 0
+ }
+ return duration
+}
+
+func cloneRetryPolicy(policy *xdsRetryPolicy) *xdsRetryPolicy {
+ if policy == nil {
+ return nil
+ }
+ cloned := *policy
+ cloned.StatusCodes = append([]uint32(nil), policy.StatusCodes...)
+ return &cloned
+}
+
func addRouteActionWeights(weights map[string]uint32, action
*routev1.RouteAction) {
if action == nil {
return
@@ -787,24 +848,63 @@ func runSampleRequestsWithOutput(ctx context.Context,
adsClient *sampleADSClient
}
}
}
- destination, endpoint, err := picker.Next()
+ requestCtx := ctx
+ cancelRequest := func() {}
+ if timeout, ok := routeTimeoutDuration(snapshot.Timeout); ok {
+ requestCtx, cancelRequest = context.WithTimeout(ctx,
timeout)
+ }
+ line, err := runSampleRequest(requestCtx, adsClient, clients,
picker, snapshot)
+ cancelRequest()
if err != nil {
return nil, err
}
+ output = append(output, line+"\n")
+ if writer != nil {
+ fmt.Fprintln(writer, line)
+ }
+ if requestInterval > 0 && i+1 < count {
+ timer := time.NewTimer(requestInterval)
+ select {
+ case <-ctx.Done():
+ if !timer.Stop() {
+ <-timer.C
+ }
+ return output, ctx.Err()
+ case <-timer.C:
+ }
+ }
+ }
+ return output, nil
+}
+
+func runSampleRequest(ctx context.Context, adsClient *sampleADSClient, clients
*sampleRequestClients, picker *smoothWeightedPicker, snapshot xdsRouteSnapshot)
(string, error) {
+ retry := snapshot.Retry
+ maxRetries := uint32(0)
+ if retry != nil {
+ maxRetries = retry.Attempts
+ }
+
+ var lastErr error
+ for attempt := uint32(0); attempt <= maxRetries; attempt++ {
+ destination, endpoint, err := picker.Next()
+ if err != nil {
+ return "", err
+ }
httpClient, scheme, err :=
clients.clientForDestination(destination)
if err != nil {
- return nil, err
+ return "", err
}
- requestCtx := ctx
- cancel := func() {}
- if timeout, ok := routeTimeoutDuration(snapshot.Timeout); ok {
- requestCtx, cancel = context.WithTimeout(ctx, timeout)
+
+ attemptCtx := ctx
+ cancelAttempt := func() {}
+ if timeout, ok := retryDuration(retry, func(policy
*xdsRetryPolicy) string { return policy.PerTryTimeout }); ok {
+ attemptCtx, cancelAttempt = context.WithTimeout(ctx,
timeout)
}
- req, err := http.NewRequestWithContext(requestCtx,
http.MethodGet,
+ req, err := http.NewRequestWithContext(attemptCtx,
http.MethodGet,
fmt.Sprintf("%s://%s%s", scheme,
net.JoinHostPort(endpoint.Address, strconv.Itoa(int(endpoint.Port))),
adsClient.path), nil)
if err != nil {
- cancel()
- return nil, err
+ cancelAttempt()
+ return "", err
}
req.Host = snapshot.Host
for name, values := range adsClient.requestHeaders {
@@ -812,35 +912,105 @@ func runSampleRequestsWithOutput(ctx context.Context,
adsClient *sampleADSClient
req.Header.Add(name, value)
}
}
- resp, err := httpClient.Do(req)
- if err != nil {
- cancel()
- return nil, err
+
+ resp, requestErr := httpClient.Do(req)
+ if requestErr != nil {
+ cancelAttempt()
+ lastErr = requestErr
+ if attempt < maxRetries &&
retryAllowsTransportFailure(retry) {
+ if err := waitRetryBackoff(ctx, retry,
attempt); err != nil {
+ return "", err
+ }
+ continue
+ }
+ return "", requestErr
}
+
body, readErr := io.ReadAll(resp.Body)
_ = resp.Body.Close()
- cancel()
+ cancelAttempt()
if readErr != nil {
- return nil, readErr
- }
- line := strings.TrimSpace(string(body))
- output = append(output, line+"\n")
- if writer != nil {
- fmt.Fprintln(writer, line)
- }
- if requestInterval > 0 && i+1 < count {
- timer := time.NewTimer(requestInterval)
- select {
- case <-ctx.Done():
- if !timer.Stop() {
- <-timer.C
+ lastErr = readErr
+ if attempt < maxRetries && retryOnIncludes(retry,
"reset") {
+ if err := waitRetryBackoff(ctx, retry,
attempt); err != nil {
+ return "", err
}
- return output, ctx.Err()
- case <-timer.C:
+ continue
+ }
+ return "", readErr
+ }
+ if attempt < maxRetries && retryStatusConfigured(retry,
uint32(resp.StatusCode)) {
+ lastErr = fmt.Errorf("upstream returned retryable
status %d", resp.StatusCode)
+ if err := waitRetryBackoff(ctx, retry, attempt); err !=
nil {
+ return "", err
}
+ continue
}
+ return strings.TrimSpace(string(body)), nil
}
- return output, nil
+ return "", lastErr
+}
+
+func retryAllowsTransportFailure(policy *xdsRetryPolicy) bool {
+ return retryOnIncludes(policy, "connect-failure") ||
retryOnIncludes(policy, "reset")
+}
+
+func retryOnIncludes(policy *xdsRetryPolicy, condition string) bool {
+ if policy == nil {
+ return false
+ }
+ for _, value := range strings.Split(policy.RetryOn, ",") {
+ if strings.TrimSpace(value) == condition {
+ return true
+ }
+ }
+ return false
+}
+
+func retryStatusConfigured(policy *xdsRetryPolicy, status uint32) bool {
+ if policy == nil || !retryOnIncludes(policy, "retriable-status-codes") {
+ return false
+ }
+ for _, configured := range policy.StatusCodes {
+ if configured == status {
+ return true
+ }
+ }
+ return false
+}
+
+func waitRetryBackoff(ctx context.Context, policy *xdsRetryPolicy, retryIndex
uint32) error {
+ base, ok := retryDuration(policy, func(value *xdsRetryPolicy) string {
return value.Backoff })
+ if !ok {
+ return nil
+ }
+ maximum, hasMaximum := retryDuration(policy, func(value
*xdsRetryPolicy) string { return value.MaxBackoff })
+ delay := base
+ for i := uint32(0); i < retryIndex && (!hasMaximum || delay < maximum);
i++ {
+ if delay > time.Duration(1<<62) {
+ delay = time.Duration(1 << 62)
+ break
+ }
+ delay *= 2
+ }
+ if hasMaximum && delay > maximum {
+ delay = maximum
+ }
+ timer := time.NewTimer(delay)
+ defer timer.Stop()
+ select {
+ case <-ctx.Done():
+ return ctx.Err()
+ case <-timer.C:
+ return nil
+ }
+}
+
+func retryDuration(policy *xdsRetryPolicy, field func(*xdsRetryPolicy) string)
(time.Duration, bool) {
+ if policy == nil {
+ return 0, false
+ }
+ return routeTimeoutDuration(field(policy))
}
func routeTimeoutDuration(value string) (time.Duration, bool) {
diff --git a/dubbod/discovery/cmd/app/grpc_outbound_test.go
b/dubbod/discovery/cmd/app/grpc_outbound_test.go
index ee82cd69..880493c6 100644
--- a/dubbod/discovery/cmd/app/grpc_outbound_test.go
+++ b/dubbod/discovery/cmd/app/grpc_outbound_test.go
@@ -24,8 +24,10 @@ import (
"net/http/httptest"
neturl "net/url"
"os"
+ "slices"
"strconv"
"strings"
+ "sync/atomic"
"testing"
"time"
@@ -296,6 +298,164 @@ func TestRunSampleRequestsUsesRouteTimeout(t *testing.T) {
}
}
+func TestRunSampleRequestsRetriesConfiguredStatus(t *testing.T) {
+ var requests atomic.Int32
+ server := httptest.NewServer(http.HandlerFunc(func(w
http.ResponseWriter, _ *http.Request) {
+ if requests.Add(1) == 1 {
+ w.WriteHeader(http.StatusServiceUnavailable)
+ _, _ = w.Write([]byte("retry"))
+ return
+ }
+ _, _ = w.Write([]byte("ok"))
+ }))
+ defer server.Close()
+
+ endpoint := endpointForServer(t, server)
+ snapshot := xdsRouteSnapshot{
+ Host: "reviews.moviereview.svc.cluster.local",
+ Port: 9080,
+ Retry: &xdsRetryPolicy{
+ Attempts: 2,
+ RetryOn:
"connect-failure,reset,retriable-status-codes",
+ StatusCodes: []uint32{503},
+ Backoff: "1ms",
+ MaxBackoff: "10ms",
+ },
+ Destinations: []xdsDestination{{
+ Cluster:
"outbound|9080|v1|reviews.moviereview.svc.cluster.local",
+ Host: "reviews.moviereview.svc.cluster.local",
+ Weight: 100,
+ Endpoints: []xdsEndpoint{endpoint},
+ }},
+ }
+ client := &sampleADSClient{
+ host: snapshot.Host,
+ port: snapshot.Port,
+ path: "/",
+ route:
map[string]uint32{snapshot.Destinations[0].Cluster: 100},
+ routeRetry: cloneRetryPolicy(snapshot.Retry),
+ endpoints:
map[string][]xdsEndpoint{snapshot.Destinations[0].Cluster: {endpoint}},
+ clusterTLS: map[string]*tlsv1.UpstreamTlsContext{},
+ requestHeaders: http.Header{},
+ }
+
+ output, err := runSampleRequestsWithOutput(context.Background(),
client, snapshot, 1, 0, time.Second, nil)
+ if err != nil {
+ t.Fatalf("runSampleRequestsWithOutput() error = %v", err)
+ }
+ if requests.Load() != 2 {
+ t.Fatalf("requests = %d, want 2", requests.Load())
+ }
+ if len(output) != 1 || strings.TrimSpace(output[0]) != "ok" {
+ t.Fatalf("output = %v, want ok", output)
+ }
+}
+
+func TestRunSampleRequestsRetriesConnectionFailureOnNextEndpoint(t *testing.T)
{
+ goodServer := httptest.NewServer(http.HandlerFunc(func(w
http.ResponseWriter, _ *http.Request) {
+ _, _ = w.Write([]byte("recovered"))
+ }))
+ defer goodServer.Close()
+
+ unusedListener, err := net.Listen("tcp", "127.0.0.1:0")
+ if err != nil {
+ t.Fatalf("net.Listen() error = %v", err)
+ }
+ badAddress := unusedListener.Addr().(*net.TCPAddr)
+ _ = unusedListener.Close()
+
+ goodEndpoint := endpointForServer(t, goodServer)
+ badEndpoint := xdsEndpoint{Address: badAddress.IP.String(), Port:
uint32(badAddress.Port)}
+ snapshot := xdsRouteSnapshot{
+ Host: "reviews.moviereview.svc.cluster.local",
+ Port: 9080,
+ Retry: &xdsRetryPolicy{
+ Attempts: 1,
+ RetryOn: "connect-failure,reset",
+ },
+ Destinations: []xdsDestination{{
+ Cluster:
"outbound|9080|v1|reviews.moviereview.svc.cluster.local",
+ Host: "reviews.moviereview.svc.cluster.local",
+ Weight: 100,
+ Endpoints: []xdsEndpoint{
+ badEndpoint,
+ goodEndpoint,
+ },
+ }},
+ }
+ client := &sampleADSClient{
+ host: snapshot.Host,
+ port: snapshot.Port,
+ path: "/",
+ route:
map[string]uint32{snapshot.Destinations[0].Cluster: 100},
+ routeRetry: cloneRetryPolicy(snapshot.Retry),
+ endpoints:
map[string][]xdsEndpoint{snapshot.Destinations[0].Cluster: {badEndpoint,
goodEndpoint}},
+ clusterTLS: map[string]*tlsv1.UpstreamTlsContext{},
+ requestHeaders: http.Header{},
+ }
+
+ output, err := runSampleRequestsWithOutput(context.Background(),
client, snapshot, 1, 0, time.Second, nil)
+ if err != nil {
+ t.Fatalf("runSampleRequestsWithOutput() error = %v", err)
+ }
+ if len(output) != 1 || strings.TrimSpace(output[0]) != "recovered" {
+ t.Fatalf("output = %v, want recovered", output)
+ }
+}
+
+func TestRunSampleRequestsRetriesPerTryTimeoutWithinTotalBudget(t *testing.T) {
+ var requests atomic.Int32
+ server := httptest.NewServer(http.HandlerFunc(func(w
http.ResponseWriter, _ *http.Request) {
+ if requests.Add(1) == 1 {
+ time.Sleep(50 * time.Millisecond)
+ _, _ = w.Write([]byte("late"))
+ return
+ }
+ _, _ = w.Write([]byte("recovered"))
+ }))
+ defer server.Close()
+
+ endpoint := endpointForServer(t, server)
+ snapshot := xdsRouteSnapshot{
+ Host: "reviews.moviereview.svc.cluster.local",
+ Port: 9080,
+ Timeout: "500ms",
+ Retry: &xdsRetryPolicy{
+ Attempts: 1,
+ RetryOn: "connect-failure,reset",
+ PerTryTimeout: "5ms",
+ },
+ Destinations: []xdsDestination{{
+ Cluster:
"outbound|9080|v1|reviews.moviereview.svc.cluster.local",
+ Host: "reviews.moviereview.svc.cluster.local",
+ Weight: 100,
+ Endpoints: []xdsEndpoint{endpoint},
+ }},
+ }
+ client := &sampleADSClient{
+ host: snapshot.Host,
+ port: snapshot.Port,
+ path: "/",
+ route:
map[string]uint32{snapshot.Destinations[0].Cluster: 100},
+ routeTimeout: 500 * time.Millisecond,
+ routeRetry: cloneRetryPolicy(snapshot.Retry),
+ endpoints:
map[string][]xdsEndpoint{snapshot.Destinations[0].Cluster: {endpoint}},
+ clusterTLS: map[string]*tlsv1.UpstreamTlsContext{},
+ requestHeaders: http.Header{},
+ }
+
+ output, err := runSampleRequestsWithOutput(context.Background(),
client, snapshot, 1, 0, time.Second, nil)
+ if err != nil {
+ t.Fatalf("runSampleRequestsWithOutput() error = %v", err)
+ }
+ if requests.Load() != 2 {
+ t.Fatalf("requests = %d, want 2", requests.Load())
+ }
+ if len(output) != 1 || strings.TrimSpace(output[0]) != "recovered" {
+ t.Fatalf("output = %v, want recovered", output)
+ }
+}
+
func TestRouteWeightsFromRoutesFiltersHeaderMatchedRoute(t *testing.T) {
resources := []*anypb.Any{mustAnyRouteConfig(t,
&routev1.RouteConfiguration{
VirtualHosts: []*routev1.VirtualHost{{
@@ -325,7 +485,7 @@ func TestRouteWeightsFromRoutesFiltersHeaderMatchedRoute(t
*testing.T) {
}},
})}
- weights, _, _, err := routeWeightsFromRoutes(resources, "/",
http.Header{"End-User": []string{"terminal-user"}})
+ weights, _, _, _, err := routeWeightsFromRoutes(resources, "/",
http.Header{"End-User": []string{"terminal-user"}})
if err != nil {
t.Fatalf("routeWeightsFromRoutes() error = %v", err)
}
@@ -333,7 +493,7 @@ func TestRouteWeightsFromRoutesFiltersHeaderMatchedRoute(t
*testing.T) {
t.Fatalf("terminal-user weights = %v, want v1=100 only",
weights)
}
- weights, _, _, err = routeWeightsFromRoutes(resources, "/", nil)
+ weights, _, _, _, err = routeWeightsFromRoutes(resources, "/", nil)
if err != nil {
t.Fatalf("routeWeightsFromRoutes() fallback error = %v", err)
}
@@ -344,6 +504,47 @@ func TestRouteWeightsFromRoutesFiltersHeaderMatchedRoute(t
*testing.T) {
}
}
+func TestRouteWeightsFromRoutesReadsRetryPolicy(t *testing.T) {
+ action := weightedRouteAction(map[string]uint32{
+ "outbound|9080|v1|reviews.moviereview.svc.cluster.local": 100,
+ })
+ action.Route.RetryPolicy = &routev1.RetryPolicy{
+ RetryOn:
"connect-failure,reset,retriable-status-codes",
+ NumRetries: wrapperspb.UInt32(3),
+ PerTryTimeout: durationpb.New(250 * time.Millisecond),
+ RetriableStatusCodes: []uint32{500, 503},
+ RetryBackOff: &routev1.RetryPolicy_RetryBackOff{
+ BaseInterval: durationpb.New(100 * time.Millisecond),
+ MaxInterval: durationpb.New(time.Second),
+ },
+ }
+ resources := []*anypb.Any{mustAnyRouteConfig(t,
&routev1.RouteConfiguration{
+ VirtualHosts: []*routev1.VirtualHost{{
+ Routes: []*routev1.Route{{
+ Match: &routev1.RouteMatch{PathSpecifier:
&routev1.RouteMatch_Prefix{Prefix: "/"}},
+ Action: action,
+ }},
+ }},
+ })}
+
+ _, _, _, retry, err := routeWeightsFromRoutes(resources, "/", nil)
+ if err != nil {
+ t.Fatalf("routeWeightsFromRoutes() error = %v", err)
+ }
+ if retry == nil {
+ t.Fatal("retry = nil")
+ }
+ if retry.Attempts != 3 || retry.RetryOn !=
"connect-failure,reset,retriable-status-codes" {
+ t.Fatalf("retry = %#v", retry)
+ }
+ if !slices.Equal(retry.StatusCodes, []uint32{500, 503}) {
+ t.Fatalf("status codes = %v", retry.StatusCodes)
+ }
+ if retry.PerTryTimeout != "250ms" || retry.Backoff != "100ms" ||
retry.MaxBackoff != "1s" {
+ t.Fatalf("retry durations = %#v", retry)
+ }
+}
+
func TestRouteWeightsFromRoutesReadsRequestTimeout(t *testing.T) {
resources := []*anypb.Any{mustAnyRouteConfig(t,
&routev1.RouteConfiguration{
VirtualHosts: []*routev1.VirtualHost{{
@@ -364,7 +565,7 @@ func TestRouteWeightsFromRoutesReadsRequestTimeout(t
*testing.T) {
}},
})}
- weights, _, timeout, err := routeWeightsFromRoutes(resources,
"/reviews", nil)
+ weights, _, timeout, _, err := routeWeightsFromRoutes(resources,
"/reviews", nil)
if err != nil {
t.Fatalf("routeWeightsFromRoutes() error = %v", err)
}
diff --git a/dubbod/discovery/pkg/networking/grpcgen/rds.go
b/dubbod/discovery/pkg/networking/grpcgen/rds.go
index 42d6f5bf..b01afeaf 100644
--- a/dubbod/discovery/pkg/networking/grpcgen/rds.go
+++ b/dubbod/discovery/pkg/networking/grpcgen/rds.go
@@ -417,6 +417,7 @@ func buildRoutesFromGatewayHTTPRoute(httpRoutes
[]config.Config, hostName host.N
routeAction.Timeout = timeout
}
}
+ routeAction.RetryPolicy =
gatewayAPIRetryPolicy(rule.Retry, rule.Timeouts)
builtRoute := &route.Route{
Match: routeMatch,
@@ -434,6 +435,48 @@ func buildRoutesFromGatewayHTTPRoute(httpRoutes
[]config.Config, hostName host.N
return allRoutes
}
+func gatewayAPIRetryPolicy(retry *sigsk8siogatewayapiapisv1.HTTPRouteRetry,
timeouts *sigsk8siogatewayapiapisv1.HTTPRouteTimeouts) *route.RetryPolicy {
+ if retry == nil {
+ return nil
+ }
+
+ attempts := uint32(1)
+ if retry.Attempts != nil {
+ if *retry.Attempts <= 0 {
+ return nil
+ }
+ attempts = uint32(*retry.Attempts)
+ }
+
+ retryOn := []string{"connect-failure", "reset"}
+ statusCodes := make([]uint32, 0, len(retry.Codes))
+ for _, code := range retry.Codes {
+ if code < 400 || code > 599 {
+ continue
+ }
+ statusCodes = append(statusCodes, uint32(code))
+ }
+ if len(statusCodes) > 0 {
+ retryOn = append(retryOn, "retriable-status-codes")
+ }
+
+ policy := &route.RetryPolicy{
+ RetryOn: strings.Join(retryOn, ","),
+ NumRetries: wrapperspb.UInt32(attempts),
+ RetriableStatusCodes: statusCodes,
+ }
+ if timeouts != nil {
+ policy.PerTryTimeout =
gatewayAPIDurationToProto(timeouts.BackendRequest)
+ }
+ if base := gatewayAPIDurationToProto(retry.Backoff); base != nil &&
base.AsDuration() > 0 {
+ policy.RetryBackOff = &route.RetryPolicy_RetryBackOff{
+ BaseInterval: base,
+ MaxInterval: durationpb.New(10 * base.AsDuration()),
+ }
+ }
+ return policy
+}
+
func gatewayAPIDurationToProto(duration *sigsk8siogatewayapiapisv1.Duration)
*durationpb.Duration {
if duration == nil {
return nil
diff --git a/dubbod/discovery/pkg/networking/grpcgen/rds_test.go
b/dubbod/discovery/pkg/networking/grpcgen/rds_test.go
index 83d43946..526fca2f 100644
--- a/dubbod/discovery/pkg/networking/grpcgen/rds_test.go
+++ b/dubbod/discovery/pkg/networking/grpcgen/rds_test.go
@@ -135,6 +135,56 @@ func TestBuildHTTPRouteSetsGatewayAPIRequestTimeout(t
*testing.T) {
}
}
+func TestBuildHTTPRouteSetsGatewayAPIRetryPolicy(t *testing.T) {
+ cfg := newServiceAttachedHTTPRouteConfig("reviews-retry",
"moviereview", "reviews", 9080)
+ spec := cfg.Spec.(*gatewayv1.HTTPRouteSpec)
+ spec.Rules[0].Timeouts = &gatewayv1.HTTPRouteTimeouts{
+ Request: ptrTo(gatewayv1.Duration("2s")),
+ BackendRequest: ptrTo(gatewayv1.Duration("250ms")),
+ }
+ spec.Rules[0].Retry = &gatewayv1.HTTPRouteRetry{
+ Codes: []gatewayv1.HTTPRouteRetryStatusCode{500, 503},
+ Attempts: ptrTo(3),
+ Backoff: ptrTo(gatewayv1.Duration("100ms")),
+ }
+ push := newRDSTestPushContext(t, []config.Config{cfg}, []*model.Service{
+ newRDSTestService("reviews", "moviereview",
"reviews.moviereview.svc.cluster.local", 9080),
+ newRDSTestService("reviews-v1", "moviereview",
"reviews-v1.moviereview.svc.cluster.local", 9080),
+ newRDSTestService("reviews-v2", "moviereview",
"reviews-v2.moviereview.svc.cluster.local", 9080),
+ })
+
+ rc := buildHTTPRoute(
+ &model.Proxy{ID: "moviepage.moviereview", Type:
model.Proxyless},
+ push,
+ "outbound|9080||reviews.moviereview.svc.cluster.local",
+ )
+ if rc == nil {
+ t.Fatal("buildHTTPRoute() returned nil")
+ }
+ retry := rc.VirtualHosts[0].Routes[0].GetRoute().GetRetryPolicy()
+ if retry == nil {
+ t.Fatal("retry policy = nil")
+ }
+ if got := retry.GetRetryOn(); got !=
"connect-failure,reset,retriable-status-codes" {
+ t.Fatalf("retry_on = %q", got)
+ }
+ if got := retry.GetNumRetries().GetValue(); got != 3 {
+ t.Fatalf("num_retries = %d, want 3", got)
+ }
+ if got := retry.GetRetriableStatusCodes(); !reflect.DeepEqual(got,
[]uint32{500, 503}) {
+ t.Fatalf("retriable_status_codes = %v", got)
+ }
+ if got := retry.GetPerTryTimeout().AsDuration(); got !=
250*time.Millisecond {
+ t.Fatalf("per_try_timeout = %v, want 250ms", got)
+ }
+ if got := retry.GetRetryBackOff().GetBaseInterval().AsDuration(); got
!= 100*time.Millisecond {
+ t.Fatalf("retry backoff base = %v, want 100ms", got)
+ }
+ if got := retry.GetRetryBackOff().GetMaxInterval().AsDuration(); got !=
time.Second {
+ t.Fatalf("retry backoff max = %v, want 1s", got)
+ }
+}
+
func newRDSTestPushContext(t *testing.T, configs []config.Config, services
[]*model.Service) *model.PushContext {
t.Helper()
diff --git a/go.mod b/go.mod
index d4fc22e0..18e16210 100644
--- a/go.mod
+++ b/go.mod
@@ -54,7 +54,7 @@ require (
github.com/heroku/color v0.0.6
github.com/kdubbo/api v0.0.0-20260713105721-558a18fe4f33
github.com/kdubbo/client-go v0.0.0-20260713105923-9a87feab8d96
- github.com/kdubbo/xds-api v0.0.0-20260722141036-e1ddf72d7580
+ github.com/kdubbo/xds-api v0.0.0-20260728021336-34287e74f6f2
github.com/moby/moby/client v0.4.1
github.com/moby/term v0.5.2
github.com/ory/viper v1.7.5
@@ -72,7 +72,7 @@ require (
golang.org/x/term v0.45.0
golang.org/x/time v0.15.0
gomodules.xyz/jsonpatch/v2 v2.5.0
- google.golang.org/grpc v1.80.0
+ google.golang.org/grpc v1.82.1
google.golang.org/protobuf v1.36.11
gopkg.in/yaml.v3 v3.0.1
helm.sh/helm/v3 v3.18.6
@@ -265,7 +265,7 @@ require (
golang.org/x/oauth2 v0.36.0 // indirect
golang.org/x/sync v0.22.0 // indirect
golang.org/x/text v0.40.0 // indirect
- google.golang.org/genproto/googleapis/api
v0.0.0-20260406210006-6f92a3bedf2d // indirect
+ google.golang.org/genproto/googleapis/api
v0.0.0-20260414002931-afd174a4e478 // indirect
google.golang.org/genproto/googleapis/rpc
v0.0.0-20260427160629-7cedc36a6bc4 // indirect
gopkg.in/evanphx/json-patch.v4 v4.13.0 // indirect
gopkg.in/inf.v0 v0.9.1 // indirect
diff --git a/go.sum b/go.sum
index eafa9299..8c160efd 100644
--- a/go.sum
+++ b/go.sum
@@ -372,8 +372,8 @@ github.com/kdubbo/api v0.0.0-20260713105721-558a18fe4f33
h1://TxoBFL/igcCqokYHOc
github.com/kdubbo/api v0.0.0-20260713105721-558a18fe4f33/go.mod
h1:8BtJiIovg7QCPsCxXcw3gDf922VcvYq5ihOSvj49Rq8=
github.com/kdubbo/client-go v0.0.0-20260713105923-9a87feab8d96
h1:co2FdpetF4dIozVHMatAd2yYPtKBqdpCgSbpdnXEkKw=
github.com/kdubbo/client-go v0.0.0-20260713105923-9a87feab8d96/go.mod
h1:NUl5DTlsovJaoYDj0SyHi263H+1Z11EXP7qnk6TPKjo=
-github.com/kdubbo/xds-api v0.0.0-20260722141036-e1ddf72d7580
h1:pzUsk+g0CqS+gyBv11LEz1nlJUyfnWNVpGj7uiwG/4o=
-github.com/kdubbo/xds-api v0.0.0-20260722141036-e1ddf72d7580/go.mod
h1:o2HDUgL1ntaDbWomZ4cD2tt8jBamuG2qRtjXOa1zZ0Q=
+github.com/kdubbo/xds-api v0.0.0-20260728021336-34287e74f6f2
h1:8tFRMU/dvoQJLaPKM/neaXjRhV/5FfxEjLy0uFJ0KG8=
+github.com/kdubbo/xds-api v0.0.0-20260728021336-34287e74f6f2/go.mod
h1:o2HDUgL1ntaDbWomZ4cD2tt8jBamuG2qRtjXOa1zZ0Q=
github.com/kevinburke/ssh_config v1.2.0
h1:x584FjTGwHzMwvHx18PXxbBVzfnxogHaAReU4gf13a4=
github.com/kevinburke/ssh_config v1.2.0/go.mod
h1:CT57kijsi8u/K/BOFA39wgDQJ9CxiF4nAY/ojJ6r6mM=
github.com/kisielk/errcheck v1.5.0/go.mod
h1:pFxgyoBC7bSaBwPgfKdkLd5X25qrDl4LWUI2bnpBCr8=
@@ -780,8 +780,8 @@ google.golang.org/appengine v1.4.0/go.mod
h1:xpcJRLb0r/rnEns0DIKYYv+WjYCduHsrkT7
google.golang.org/genproto v0.0.0-20180817151627-c66870c02cf8/go.mod
h1:JiN7NxoALGmiZfu7CAH4rXhgtRTLTxftemlI0sWmxmc=
google.golang.org/genproto v0.0.0-20190819201941-24fa4b261c55/go.mod
h1:DMBHOl98Agz4BDEuKkezgsaosCRResVns1a3J2ZsMNc=
google.golang.org/genproto v0.0.0-20200423170343-7949de9c1215/go.mod
h1:55QSHmfGQM9UVYDPBsyGGes0y52j32PQ3BqQfXhyH3c=
-google.golang.org/genproto/googleapis/api v0.0.0-20260406210006-6f92a3bedf2d
h1:/aDRtSZJjyLQzm75d+a1wOJaqyKBMvIAfeQmoa3ORiI=
-google.golang.org/genproto/googleapis/api
v0.0.0-20260406210006-6f92a3bedf2d/go.mod
h1:etfGUgejTiadZAUaEP14NP97xi1RGeawqkjDARA/UOs=
+google.golang.org/genproto/googleapis/api v0.0.0-20260414002931-afd174a4e478
h1:yQugLulqltosq0B/f8l4w9VryjV+N/5gcW0jQ3N8Qec=
+google.golang.org/genproto/googleapis/api
v0.0.0-20260414002931-afd174a4e478/go.mod
h1:C6ADNqOxbgdUUeRTU+LCHDPB9ttAMCTff6auwCVa4uc=
google.golang.org/genproto/googleapis/rpc v0.0.0-20260427160629-7cedc36a6bc4
h1:tEkOQcXgF6dH1G+MVKZrfpYvozGrzb91k6ha7jireSM=
google.golang.org/genproto/googleapis/rpc
v0.0.0-20260427160629-7cedc36a6bc4/go.mod
h1:4Hqkh8ycfw05ld/3BWL7rJOSfebL2Q+DVDeRgYgxUU8=
google.golang.org/grpc v1.19.0/go.mod
h1:mqu4LbDTu4XGKhr4mRzUsmM4RtVoemTSY81AxZiDr8c=
@@ -789,8 +789,8 @@ google.golang.org/grpc v1.23.0/go.mod
h1:Y5yQAOtifL1yxbo5wqy6BxZv8vAUGQwXBOALyac
google.golang.org/grpc v1.25.1/go.mod
h1:c3i+UQWmh7LiEpx4sFZnkU36qjEYZ0imhYfXVyQciAY=
google.golang.org/grpc v1.27.0/go.mod
h1:qbnxyOmOxrQa7FizSgH+ReBfzJrCY1pSN7KXBS8abTk=
google.golang.org/grpc v1.29.1/go.mod
h1:itym6AZVZYACWQqET3MqgPpjcuV5QH3BxFS3IjizoKk=
-google.golang.org/grpc v1.80.0 h1:Xr6m2WmWZLETvUNvIUmeD5OAagMw3FiKmMlTdViWsHM=
-google.golang.org/grpc v1.80.0/go.mod
h1:ho/dLnxwi3EDJA4Zghp7k2Ec1+c2jqup0bFkw07bwF4=
+google.golang.org/grpc v1.82.1 h1:NnAxzGRA0677vCa4BUkOAnO5+FfQqVl9iUXeD0IqcGE=
+google.golang.org/grpc v1.82.1/go.mod
h1:yzTZ1TB1Z3SG+LIYaI+WiE8D5+PZ3ArnrSp8zF3+/ZA=
google.golang.org/protobuf v1.36.11
h1:fV6ZwhNocDyBLK0dj+fg8ektcVegBBuEolpbTQyBNVE=
google.golang.org/protobuf v1.36.11/go.mod
h1:HTf+CrKn2C3g5S8VImy6tdcUvCska2kB7j23XfzDpco=
gopkg.in/alecthomas/kingpin.v2 v2.2.6/go.mod
h1:FMv+mEhP44yOT+4EoQTLFTRgOQ1FBLkstjWtayDeSgw=
diff --git a/samples/app/httproute.yaml b/samples/app/httproute.yaml
index 91328734..21417a0a 100644
--- a/samples/app/httproute.yaml
+++ b/samples/app/httproute.yaml
@@ -25,7 +25,18 @@ spec:
name: nginx
port: 80
rules:
- - backendRefs:
+ - retry:
+ attempts: 3
+ codes:
+ - 500
+ - 502
+ - 503
+ - 504
+ backoff: 100ms
+ timeouts:
+ request: 5s
+ backendRequest: 1s
+ backendRefs:
- name: nginx-v1
port: 80
weight: 50
diff --git a/tests/e2e/run.sh b/tests/e2e/run.sh
index 672b9571..bae668af 100755
--- a/tests/e2e/run.sh
+++ b/tests/e2e/run.sh
@@ -126,9 +126,10 @@ if [[ "${IMAGE}" != "${UPGRADE_FROM_IMAGE}" ]]; then
fi
log "installing Gateway API CRDs"
-# Pin to the sigs.k8s.io/gateway-api version in go.mod.
+# Pin to the sigs.k8s.io/gateway-api version in go.mod. HTTPRoute retry is an
+# Extended feature carried by the experimental CRD bundle.
GATEWAY_API_VERSION="${GATEWAY_API_VERSION:-v1.4.1}"
-"${KUBECTL[@]}" apply -f
"https://github.com/kubernetes-sigs/gateway-api/releases/download/${GATEWAY_API_VERSION}/standard-install.yaml"
+"${KUBECTL[@]}" apply --server-side -f
"https://github.com/kubernetes-sigs/gateway-api/releases/download/${GATEWAY_API_VERSION}/experimental-install.yaml"
log "installing base chart (CRDs)"
helm upgrade --install dubbo-base "${ROOT}/manifests/charts/base" \