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" \

Reply via email to