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 588f80c2 Supports fault injection
588f80c2 is described below
commit 588f80c26c0fdf4a56f0eef101403e06b18d197a
Author: mfordjody <[email protected]>
AuthorDate: Mon Aug 3 12:33:20 2026 +0800
Supports fault injection
Merged after CI, Maintainer Approval Gate, dependency publication, and
/lgtm workflow validation.
---
.github/workflows/commands.yml | 3 +-
architecture/tests/PERFORMANCE.md | 115 -----------------
dubbod/discovery/cmd/app/grpc_outbound.go | 86 ++++++++++++-
dubbod/discovery/cmd/app/grpc_outbound_test.go | 97 ++++++++++++++-
.../pkg/bootstrap/proxyless_grpc_controller.go | 49 +++++++-
.../bootstrap/proxyless_grpc_controller_test.go | 49 ++++++++
dubbod/discovery/pkg/bootstrap/server.go | 5 +-
.../pkg/config/kube/crdclient/types.gen.go | 51 ++++++++
dubbod/discovery/pkg/model/push_context.go | 90 ++++++++++++++
dubbod/discovery/pkg/networking/grpcgen/rds.go | 46 ++++++-
.../discovery/pkg/networking/grpcgen/rds_test.go | 59 +++++++++
go.mod | 6 +-
go.sum | 12 +-
manifests/charts/base/files/crd-all.gen.yaml | 136 +++++++++++++++++++++
pkg/config/schema/collections/collections.gen.go | 19 +++
pkg/config/schema/gvk/resources.gen.go | 7 ++
pkg/config/schema/gvr/resources.gen.go | 3 +
pkg/config/schema/kind/resources.gen.go | 5 +
pkg/config/schema/kubeclient/resources.gen.go | 13 ++
pkg/config/schema/kubetypes/resources.gen.go | 4 +
pkg/config/schema/metadata.yaml | 10 ++
pkg/config/validation/helpers.go | 11 ++
pkg/config/validation/validators.go | 50 ++++++++
pkg/config/validation/validators_test.go | 87 +++++++++++++
samples/httpbin/fault-injection.yaml | 15 +++
25 files changed, 882 insertions(+), 146 deletions(-)
diff --git a/.github/workflows/commands.yml b/.github/workflows/commands.yml
index 8c45f579..19f092ba 100644
--- a/.github/workflows/commands.yml
+++ b/.github/workflows/commands.yml
@@ -56,8 +56,7 @@ jobs:
timeout-minutes: 10
if: |
github.event_name == 'issue_comment' &&
- github.repository == 'apache/dubbo-kubernetes' &&
- startsWith(github.event.comment.body, '/')
+ github.repository == 'apache/dubbo-kubernetes'
env:
BOT_APP_ID: ${{ secrets.BOT_APP_ID }}
steps:
diff --git a/architecture/tests/PERFORMANCE.md
b/architecture/tests/PERFORMANCE.md
deleted file mode 100644
index d20b631b..00000000
--- a/architecture/tests/PERFORMANCE.md
+++ /dev/null
@@ -1,115 +0,0 @@
-<!--
-Licensed to the Apache Software Foundation (ASF) under one or more
-contributor license agreements. See the NOTICE file distributed with
-this work for additional information regarding copyright ownership.
-The ASF licenses this file to You under the Apache License, Version 2.0
-(the "License"); you may not use this file except in compliance with
-the License. You may obtain a copy of the License at
-
- http://www.apache.org/licenses/LICENSE-2.0
-
-Unless required by applicable law or agreed to in writing, software
-distributed under the License is distributed on an "AS IS" BASIS,
-WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
-See the License for the specific language governing permissions and
-limitations under the License.
--->
-
-# Performance Engineering
-
-The performance suite keeps two layers separate:
-
-1. Control-plane Go benchmarks measure indexing, client selection, and xDS
- response generation without Kubernetes or network noise.
-2. Cluster load tests measure end-to-end request latency and resource usage.
- Those results must identify the cluster, traffic generator, workload, and
- observability setup; they must not be compared directly with Go benchmark
- timings.
-
-## Control-plane scale matrix
-
-`BenchmarkProxylessPushScale` exercises the real targeted endpoint-update path:
-
-```text
-clientsForPush -> pushConnection -> EDS generation -> DiscoveryResponse send
-```
-
-Every synthetic proxyless connection watches every generated service. A
-single service endpoint update therefore selects every connection and emits
-one EDS response per connection. The matrix changes one dimension at a time
-from the `100 services / 10 endpoints / 100 connections` baseline:
-
-| Dimension | Cases |
-|---|---|
-| Services | 10, 100, 1,000 |
-| Endpoints per service | 1, 10, 100 |
-| Connections | 10, 100, 1,000 |
-
-`BenchmarkProxylessFullPushScale` uses the same dimensions but forces a full
-push that emits every watched service. Its baseline uses 10 connections, and
-its largest cases cap the generated work at 100,000 endpoint copies per
-operation:
-
-| Dimension | Cases |
-|---|---|
-| Services | 10, 100, 1,000 |
-| Endpoints per service | 1, 10, 100 |
-| Connections | 1, 10, 100 |
-
-Both benchmarks report `ns/op`, `B/op`, `allocs/op`, the three scale values,
-and the number of generated endpoint copies per operation.
-`ns/op` is the serial cost of completing one synthetic push wave on the
-benchmark CPU. Production pushes may execute concurrently, so this number is
-for regression and scaling comparisons rather than a direct production SLO.
-
-The existing `BenchmarkClientsForPushProxylessTargetedScale` separately
-measures connection selection through 100,000 connected clients, while
-`BenchmarkKRTFetch` measures indexed and label-scan lookups over 10,000
-workloads.
-
-## Running benchmarks
-
-Run the representative targeted and full-push scale cases used by pull-request
-CI:
-
-```bash
-make benchmark-smoke
-```
-
-Run the complete matrix with repeated samples:
-
-```bash
-make benchmark BENCHMARK_COUNT=5 BENCHMARK_TIME=1s BENCHMARK_CPU=1
-```
-
-The Make variables can narrow the benchmark expression or adjust measurement
-time and CPU count:
-
-```bash
-make benchmark \
- BENCHMARK_PATTERN=BenchmarkProxylessPushScale \
- BENCHMARK_COUNT=10 \
- BENCHMARK_TIME=2s \
- BENCHMARK_CPU=1
-```
-
-Use the same Go version, CPU model, power mode, `BENCHMARK_CPU`, and background
-load for baseline and candidate runs. Keep at least five samples and compare
-them with `benchstat`; do not draw a capacity conclusion from one run.
-
-## Continuous evidence
-
-The main CI workflow runs `make benchmark-smoke` to catch fixture failures and
-gross scale regressions on every change. The scheduled `Performance
-Benchmarks` workflow runs five samples of the complete matrix each week and
-retains the raw output plus runner details for 30 days. Shared runner variance
-makes a universal absolute latency threshold unreliable; compare like-for-like
-runs and investigate changes in both time and allocations.
-
-## Cluster load tests
-
-Use `samples/httpbin/httpbin.yaml` as the minimal data-plane target and the
-Prometheus sample under `samples/addons/` for control-plane metrics. Record at
-least request rate, p50/p95/p99 latency, error rate, CPU, memory, connected xDS
-clients, services, and endpoints. Change one scale dimension at a time and
-include warm-up, steady-state, and configuration-churn phases in the result.
diff --git a/dubbod/discovery/cmd/app/grpc_outbound.go
b/dubbod/discovery/cmd/app/grpc_outbound.go
index c789f9a7..57e191c9 100644
--- a/dubbod/discovery/cmd/app/grpc_outbound.go
+++ b/dubbod/discovery/cmd/app/grpc_outbound.go
@@ -21,6 +21,7 @@ import (
"encoding/json"
"fmt"
"io"
+ "math/rand/v2"
"net"
"net/http"
"os"
@@ -53,6 +54,7 @@ import (
"google.golang.org/protobuf/types/known/anypb"
"google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/structpb"
+ "google.golang.org/protobuf/types/known/wrapperspb"
)
type grpcOutboundOptions struct {
@@ -103,6 +105,7 @@ type xdsRouteSnapshot struct {
Port int `json:"port"`
Timeout string `json:"timeout,omitempty"`
Retry *xdsRetryPolicy `json:"retry,omitempty"`
+ Fault *xdsFaultPolicy `json:"fault,omitempty"`
Destinations []xdsDestination `json:"destinations"`
}
@@ -115,6 +118,13 @@ type xdsRetryPolicy struct {
MaxBackoff string `json:"maxBackoff,omitempty"`
}
+type xdsFaultPolicy struct {
+ Delay string `json:"delay,omitempty"`
+ DelayPercentage uint32 `json:"delayPercentage,omitempty"`
+ AbortStatus uint32 `json:"abortStatus,omitempty"`
+ AbortPercentage uint32 `json:"abortPercentage,omitempty"`
+}
+
type sampleADSClient struct {
conn *grpc.ClientConn
stream
discovery.AggregatedDiscoveryService_StreamAggregatedResourcesClient
@@ -130,6 +140,7 @@ type sampleADSClient struct {
route map[string]uint32
routeTimeout time.Duration
routeRetry *xdsRetryPolicy
+ routeFault *xdsFaultPolicy
endpoints map[string][]xdsEndpoint
clusterTLS map[string]*tlsv1.UpstreamTlsContext
updates chan struct{}
@@ -439,7 +450,7 @@ func (c *sampleADSClient) handleResponse(resp
*discovery.DiscoveryResponse) erro
return c.subscribe(v1.RouteType, routeNames)
}
case v1.RouteType:
- weights, clusters, timeout, retry, err :=
routeWeightsFromRoutes(resp.Resources, c.path, c.requestHeaders)
+ weights, clusters, timeout, retry, fault, err :=
routeWeightsFromRoutes(resp.Resources, c.path, c.requestHeaders)
if err != nil {
return err
}
@@ -447,6 +458,7 @@ func (c *sampleADSClient) handleResponse(resp
*discovery.DiscoveryResponse) erro
c.route = weights
c.routeTimeout = timeout
c.routeRetry = retry
+ c.routeFault = fault
c.mu.Unlock()
c.notify()
if len(clusters) > 0 {
@@ -526,6 +538,7 @@ func (c *sampleADSClient) readySnapshot(expected
map[string]uint32) (xdsRouteSna
snapshot.Timeout = c.routeTimeout.String()
}
snapshot.Retry = cloneRetryPolicy(c.routeRetry)
+ snapshot.Fault = cloneFaultPolicy(c.routeFault)
for clusterName, weight := range c.route {
if weight == 0 {
continue
@@ -584,12 +597,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,
*xdsRetryPolicy, error) {
+func routeWeightsFromRoutes(resources []*anypb.Any, requestPath string,
requestHeaders http.Header) (map[string]uint32, []string, time.Duration,
*xdsRetryPolicy, *xdsFaultPolicy, 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, nil, err
+ return nil, nil, 0, nil, nil, err
}
for _, vh := range rc.GetVirtualHosts() {
for _, rt := range vh.GetRoutes() {
@@ -598,11 +611,11 @@ func routeWeightsFromRoutes(resources []*anypb.Any,
requestPath string, requestH
}
action := rt.GetRoute()
addRouteActionWeights(weights, action)
- return weights,
sortedWeightClusterNames(weights), routeActionTimeout(action),
routeActionRetryPolicy(action), nil
+ return weights,
sortedWeightClusterNames(weights), routeActionTimeout(action),
routeActionRetryPolicy(action), routeActionFaultPolicy(action), nil
}
}
}
- return weights, sortedWeightClusterNames(weights), 0, nil, nil
+ return weights, sortedWeightClusterNames(weights), 0, nil, nil, nil
}
func routeActionTimeout(action *routev1.RouteAction) time.Duration {
@@ -659,6 +672,41 @@ func cloneRetryPolicy(policy *xdsRetryPolicy)
*xdsRetryPolicy {
return &cloned
}
+func routeActionFaultPolicy(action *routev1.RouteAction) *xdsFaultPolicy {
+ if action == nil || action.GetFaultPolicy() == nil {
+ return nil
+ }
+ fault := action.GetFaultPolicy()
+ out := &xdsFaultPolicy{}
+ if delay := fault.GetDelay(); delay != nil &&
positiveProtoDuration(delay.GetFixedDelay()) > 0 {
+ out.Delay = delay.GetFixedDelay().AsDuration().String()
+ out.DelayPercentage = percentageValue(delay.GetPercentage())
+ }
+ if abort := fault.GetAbort(); abort != nil && abort.GetHttpStatus() !=
0 {
+ out.AbortStatus = abort.GetHttpStatus()
+ out.AbortPercentage = percentageValue(abort.GetPercentage())
+ }
+ if out.Delay == "" && out.AbortStatus == 0 {
+ return nil
+ }
+ return out
+}
+
+func percentageValue(value *wrapperspb.UInt32Value) uint32 {
+ if value == nil {
+ return 100
+ }
+ return value.GetValue()
+}
+
+func cloneFaultPolicy(policy *xdsFaultPolicy) *xdsFaultPolicy {
+ if policy == nil {
+ return nil
+ }
+ cloned := *policy
+ return &cloned
+}
+
func addRouteActionWeights(weights map[string]uint32, action
*routev1.RouteAction) {
if action == nil {
return
@@ -878,6 +926,23 @@ func runSampleRequestsWithOutput(ctx context.Context,
adsClient *sampleADSClient
}
func runSampleRequest(ctx context.Context, adsClient *sampleADSClient, clients
*sampleRequestClients, picker *smoothWeightedPicker, snapshot xdsRouteSnapshot)
(string, error) {
+ if fault := snapshot.Fault; fault != nil {
+ if delay, ok := routeTimeoutDuration(fault.Delay); ok &&
faultPercentageMatches(fault.DelayPercentage) {
+ timer := time.NewTimer(delay)
+ select {
+ case <-ctx.Done():
+ if !timer.Stop() {
+ <-timer.C
+ }
+ return "", ctx.Err()
+ case <-timer.C:
+ }
+ }
+ if fault.AbortStatus != 0 &&
faultPercentageMatches(fault.AbortPercentage) {
+ return "", fmt.Errorf("fault injected HTTP status %d",
fault.AbortStatus)
+ }
+ }
+
retry := snapshot.Retry
maxRetries := uint32(0)
if retry != nil {
@@ -951,6 +1016,17 @@ func runSampleRequest(ctx context.Context, adsClient
*sampleADSClient, clients *
return "", lastErr
}
+func faultPercentageMatches(percentage uint32) bool {
+ switch {
+ case percentage == 0:
+ return false
+ case percentage >= 100:
+ return true
+ default:
+ return rand.Uint32N(100) < percentage
+ }
+}
+
func retryAllowsTransportFailure(policy *xdsRetryPolicy) bool {
return retryOnIncludes(policy, "connect-failure") ||
retryOnIncludes(policy, "reset")
}
diff --git a/dubbod/discovery/cmd/app/grpc_outbound_test.go
b/dubbod/discovery/cmd/app/grpc_outbound_test.go
index 880493c6..2337413c 100644
--- a/dubbod/discovery/cmd/app/grpc_outbound_test.go
+++ b/dubbod/discovery/cmd/app/grpc_outbound_test.go
@@ -485,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)
}
@@ -493,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)
}
@@ -527,7 +527,7 @@ func TestRouteWeightsFromRoutesReadsRetryPolicy(t
*testing.T) {
}},
})}
- _, _, _, retry, err := routeWeightsFromRoutes(resources, "/", nil)
+ _, _, _, retry, _, err := routeWeightsFromRoutes(resources, "/", nil)
if err != nil {
t.Fatalf("routeWeightsFromRoutes() error = %v", err)
}
@@ -545,6 +545,95 @@ func TestRouteWeightsFromRoutesReadsRetryPolicy(t
*testing.T) {
}
}
+func TestRouteWeightsFromRoutesReadsFaultPolicy(t *testing.T) {
+ action := weightedRouteAction(map[string]uint32{
+ "outbound|9080|v1|reviews.moviereview.svc.cluster.local": 100,
+ })
+ action.Route.FaultPolicy = &routev1.FaultPolicy{
+ Delay: &routev1.FaultDelay{
+ FixedDelay: durationpb.New(250 * time.Millisecond),
+ Percentage: wrapperspb.UInt32(20),
+ },
+ Abort: &routev1.FaultAbort{
+ HttpStatus: http.StatusServiceUnavailable,
+ Percentage: wrapperspb.UInt32(10),
+ },
+ }
+ resources := []*anypb.Any{mustAnyRouteConfig(t,
&routev1.RouteConfiguration{
+ VirtualHosts: []*routev1.VirtualHost{{
+ Routes: []*routev1.Route{{
+ Match: &routev1.RouteMatch{PathSpecifier:
&routev1.RouteMatch_Prefix{Prefix: "/"}},
+ Action: action,
+ }},
+ }},
+ })}
+
+ _, _, _, _, fault, err := routeWeightsFromRoutes(resources, "/", nil)
+ if err != nil {
+ t.Fatalf("routeWeightsFromRoutes() error = %v", err)
+ }
+ if fault == nil {
+ t.Fatal("fault = nil")
+ }
+ if fault.Delay != "250ms" || fault.DelayPercentage != 20 ||
+ fault.AbortStatus != http.StatusServiceUnavailable ||
fault.AbortPercentage != 10 {
+ t.Fatalf("fault = %#v", fault)
+ }
+}
+
+func TestRunSampleRequestAppliesFaultBeforeUpstreamAndRetry(t *testing.T) {
+ var requests atomic.Int32
+ server := httptest.NewServer(http.HandlerFunc(func(w
http.ResponseWriter, _ *http.Request) {
+ requests.Add(1)
+ _, _ = w.Write([]byte("unexpected"))
+ }))
+ defer server.Close()
+
+ endpoint := endpointForServer(t, server)
+ snapshot := xdsRouteSnapshot{
+ Host: "reviews.moviereview.svc.cluster.local",
+ Port: 9080,
+ Fault: &xdsFaultPolicy{
+ Delay: "20ms",
+ DelayPercentage: 100,
+ AbortStatus: http.StatusServiceUnavailable,
+ AbortPercentage: 100,
+ },
+ Retry: &xdsRetryPolicy{
+ Attempts: 2,
+ RetryOn: "retriable-status-codes",
+ StatusCodes: []uint32{http.StatusServiceUnavailable},
+ },
+ Destinations: []xdsDestination{{
+ Cluster:
"outbound|9080||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: "/",
+ requestHeaders: http.Header{},
+ }
+ picker, err := newSmoothWeightedPicker(snapshot)
+ if err != nil {
+ t.Fatal(err)
+ }
+ started := time.Now()
+ _, err = runSampleRequest(context.Background(), client,
newSampleRequestClients(client, time.Second), picker, snapshot)
+ if err == nil || err.Error() != "fault injected HTTP status 503" {
+ t.Fatalf("runSampleRequest() error = %v", err)
+ }
+ if elapsed := time.Since(started); elapsed < 20*time.Millisecond {
+ t.Fatalf("fault delay elapsed = %v, want at least 20ms",
elapsed)
+ }
+ if requests.Load() != 0 {
+ t.Fatalf("upstream requests = %d, want 0", requests.Load())
+ }
+}
+
func TestRouteWeightsFromRoutesReadsRequestTimeout(t *testing.T) {
resources := []*anypb.Any{mustAnyRouteConfig(t,
&routev1.RouteConfiguration{
VirtualHosts: []*routev1.VirtualHost{{
@@ -565,7 +654,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/bootstrap/proxyless_grpc_controller.go
b/dubbod/discovery/pkg/bootstrap/proxyless_grpc_controller.go
index ad98f61b..1f310477 100644
--- a/dubbod/discovery/pkg/bootstrap/proxyless_grpc_controller.go
+++ b/dubbod/discovery/pkg/bootstrap/proxyless_grpc_controller.go
@@ -240,7 +240,7 @@ func proxylessGRPCRuntimeConfigNeedsUpdate(req
*discoverymodel.PushRequest) bool
}
for cfg := range req.ConfigsUpdated {
switch cfg.Kind {
- case kind.HTTPRoute, kind.BackendTLSPolicy,
kind.CircuitBreakerPolicy, kind.PeerAuthentication, kind.RequestAuthentication,
kind.AuthorizationPolicy, kind.Service, kind.EndpointSlice, kind.Endpoints,
kind.Pod, kind.Namespace:
+ case kind.HTTPRoute, kind.BackendTLSPolicy,
kind.CircuitBreakerPolicy, kind.FaultInjectionPolicy, kind.PeerAuthentication,
kind.RequestAuthentication, kind.AuthorizationPolicy, kind.Service,
kind.EndpointSlice, kind.Endpoints, kind.Pod, kind.Namespace:
return true
}
}
@@ -409,9 +409,25 @@ type proxylessGRPCServiceRuntimeConfig struct {
}
type proxylessGRPCPortRuntimeConfig struct {
- Name string `json:"name,omitempty"`
- Port int `json:"port"`
- MTLSMode string `json:"mtlsMode,omitempty"`
+ Name string `json:"name,omitempty"`
+ Port int `json:"port"`
+ MTLSMode string `json:"mtlsMode,omitempty"`
+ Fault *proxylessGRPCFaultRuntimeConfig `json:"fault,omitempty"`
+}
+
+type proxylessGRPCFaultRuntimeConfig struct {
+ Delay *proxylessGRPCFaultDelayRuntimeConfig `json:"delay,omitempty"`
+ Abort *proxylessGRPCFaultAbortRuntimeConfig `json:"abort,omitempty"`
+}
+
+type proxylessGRPCFaultDelayRuntimeConfig struct {
+ FixedDelay string `json:"fixedDelay"`
+ Percentage uint32 `json:"percentage"`
+}
+
+type proxylessGRPCFaultAbortRuntimeConfig struct {
+ HTTPStatus uint32 `json:"httpStatus"`
+ Percentage uint32 `json:"percentage"`
}
type proxylessGRPCRouteRuntimeConfig struct {
@@ -740,6 +756,7 @@ func buildRuntimeServiceConfig(push
*discoverymodel.PushContext, endpointIndex *
Name: port.Name,
Port: port.Port,
MTLSMode: runtimeInboundMTLSMode(push,
svc.Attributes.Namespace, port.Port),
+ Fault: runtimeFaultInjection(push,
svc.Attributes.Namespace, svc.Attributes.Name, port.Name),
})
cfg.Endpoints = append(cfg.Endpoints,
runtimeEndpointsForService(endpointIndex, svc, port.Port, nil)...)
}
@@ -747,6 +764,30 @@ func buildRuntimeServiceConfig(push
*discoverymodel.PushContext, endpointIndex *
return cfg
}
+func runtimeFaultInjection(push *discoverymodel.PushContext, namespace, name,
portName string) *proxylessGRPCFaultRuntimeConfig {
+ settings, found := push.FaultInjectionForService(namespace, name,
portName)
+ if !found {
+ return nil
+ }
+ fault := &proxylessGRPCFaultRuntimeConfig{}
+ if settings.Delay > 0 {
+ fault.Delay = &proxylessGRPCFaultDelayRuntimeConfig{
+ FixedDelay: settings.Delay.String(),
+ Percentage: settings.DelayPercentage,
+ }
+ }
+ if settings.AbortStatus != 0 {
+ fault.Abort = &proxylessGRPCFaultAbortRuntimeConfig{
+ HTTPStatus: settings.AbortStatus,
+ Percentage: settings.AbortPercentage,
+ }
+ }
+ if fault.Delay == nil && fault.Abort == nil {
+ return nil
+ }
+ return fault
+}
+
func buildRuntimeRouteConfig(push *discoverymodel.PushContext, endpointIndex
*discoverymodel.EndpointIndex, svc *discoverymodel.Service, port int)
proxylessGRPCRouteRuntimeConfig {
cfg := proxylessGRPCRouteRuntimeConfig{
Host: string(svc.Hostname),
diff --git a/dubbod/discovery/pkg/bootstrap/proxyless_grpc_controller_test.go
b/dubbod/discovery/pkg/bootstrap/proxyless_grpc_controller_test.go
index 1898e1b4..f2bcb876 100644
--- a/dubbod/discovery/pkg/bootstrap/proxyless_grpc_controller_test.go
+++ b/dubbod/discovery/pkg/bootstrap/proxyless_grpc_controller_test.go
@@ -36,7 +36,10 @@ import (
"github.com/apache/dubbo-kubernetes/pkg/kube/inject"
"github.com/apache/dubbo-kubernetes/pkg/kube/krt"
"github.com/apache/dubbo-kubernetes/pkg/util/sets"
+ networking "github.com/kdubbo/api/networking/v1alpha3"
security "github.com/kdubbo/api/security/v1alpha3"
+ "google.golang.org/protobuf/types/known/durationpb"
+ "google.golang.org/protobuf/types/known/wrapperspb"
corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/types"
@@ -306,6 +309,45 @@ func
TestBuildRuntimeTrafficConfigCapturesPermissivePeerAuthentication(t *testin
}
}
+func TestBuildRuntimeTrafficConfigCapturesFaultInjection(t *testing.T) {
+ hostname := host.Name("provider.grpc-app.svc.cluster.local")
+ svc := newProxylessRuntimeTestService("provider", "grpc-app",
string(hostname), 17070)
+ push := newProxylessRuntimeTestPushContext(t, []config.Config{{
+ Meta: config.Meta{
+ GroupVersionKind: gvk.FaultInjectionPolicy,
+ Name: "provider-fault",
+ Namespace: "grpc-app",
+ },
+ Spec: &networking.FaultInjectionPolicy{
+ TargetRefs: []*networking.PolicyTargetReference{{
+ Kind: "Service",
+ Name: "provider",
+ SectionName: "grpc",
+ }},
+ Delay: &networking.FaultDelay{
+ FixedDelay: durationpb.New(250 *
time.Millisecond),
+ Percentage: wrapperspb.UInt32(20),
+ },
+ Abort: &networking.FaultAbort{
+ HttpStatus: 503,
+ Percentage: wrapperspb.UInt32(10),
+ },
+ },
+ }}, []*discoverymodel.Service{svc})
+
+ serviceConfig := buildRuntimeServiceConfig(push, nil, svc)
+ if len(serviceConfig.Ports) != 1 || serviceConfig.Ports[0].Fault == nil
{
+ t.Fatalf("ports = %+v, want one port with fault",
serviceConfig.Ports)
+ }
+ fault := serviceConfig.Ports[0].Fault
+ if fault.Delay == nil || fault.Delay.FixedDelay != "250ms" ||
fault.Delay.Percentage != 20 {
+ t.Fatalf("delay = %+v", fault.Delay)
+ }
+ if fault.Abort == nil || fault.Abort.HTTPStatus != 503 ||
fault.Abort.Percentage != 10 {
+ t.Fatalf("abort = %+v", fault.Abort)
+ }
+}
+
func TestProxylessGRPCRuntimeConfigNeedsUpdate(t *testing.T) {
tests := []struct {
name string
@@ -324,6 +366,13 @@ func TestProxylessGRPCRuntimeConfigNeedsUpdate(t
*testing.T) {
},
want: true,
},
+ {
+ name: "fault injection policy",
+ req: &discoverymodel.PushRequest{
+ ConfigsUpdated:
sets.New(discoverymodel.ConfigKey{Kind: kind.FaultInjectionPolicy, Name:
"provider-fault", Namespace: "grpc-app"}),
+ },
+ want: true,
+ },
{
name: "peerauthentication",
req: &discoverymodel.PushRequest{
diff --git a/dubbod/discovery/pkg/bootstrap/server.go
b/dubbod/discovery/pkg/bootstrap/server.go
index f5ac8bd0..f329fbff 100644
--- a/dubbod/discovery/pkg/bootstrap/server.go
+++ b/dubbod/discovery/pkg/bootstrap/server.go
@@ -515,6 +515,8 @@ func (s *Server) initRegistryEventHandlers() {
configKind = kind.BackendTLSPolicy
case "CircuitBreakerPolicy":
configKind = kind.CircuitBreakerPolicy
+ case "FaultInjectionPolicy":
+ configKind = kind.FaultInjectionPolicy
default:
log.Debugf("unknown schema identifier %s for %v,
skipping", schemaID, cfg.GroupVersionKind)
return
@@ -539,7 +541,8 @@ func (s *Server) initRegistryEventHandlers() {
configKind == kind.HTTPRoute ||
configKind == kind.BackendTLSPolicy ||
configKind == kind.ReferenceGrant ||
- configKind == kind.CircuitBreakerPolicy
+ configKind == kind.CircuitBreakerPolicy ||
+ configKind == kind.FaultInjectionPolicy
// Trigger ConfigUpdate to push changes to all connected proxies
s.XDSServer.ConfigUpdate(&model.PushRequest{
diff --git a/dubbod/discovery/pkg/config/kube/crdclient/types.gen.go
b/dubbod/discovery/pkg/config/kube/crdclient/types.gen.go
index 818efced..44ba1f33 100755
--- a/dubbod/discovery/pkg/config/kube/crdclient/types.gen.go
+++ b/dubbod/discovery/pkg/config/kube/crdclient/types.gen.go
@@ -50,6 +50,11 @@ func create(c kube.Client, cfg config.Config, objMeta
metav1.ObjectMeta) (metav1
ObjectMeta: objMeta,
Spec:
*(cfg.Spec.(*githubcomkdubboapinetworkingv1alpha3.CircuitBreakerPolicy)),
}, metav1.CreateOptions{})
+ case gvk.FaultInjectionPolicy:
+ return
c.Dubbo().NetworkingV1alpha3().FaultInjectionPolicies(cfg.Namespace).Create(context.TODO(),
&apigithubcomapachedubbokubernetesapinetworkingv1alpha3.FaultInjectionPolicy{
+ ObjectMeta: objMeta,
+ Spec:
*(cfg.Spec.(*githubcomkdubboapinetworkingv1alpha3.FaultInjectionPolicy)),
+ }, metav1.CreateOptions{})
case gvk.GatewayClass:
return
c.GatewayAPI().GatewayV1().GatewayClasses().Create(context.TODO(),
&sigsk8siogatewayapiapisv1.GatewayClass{
ObjectMeta: objMeta,
@@ -117,6 +122,11 @@ func update(c kube.Client, cfg config.Config, objMeta
metav1.ObjectMeta) (metav1
ObjectMeta: objMeta,
Spec:
*(cfg.Spec.(*githubcomkdubboapinetworkingv1alpha3.CircuitBreakerPolicy)),
}, metav1.UpdateOptions{})
+ case gvk.FaultInjectionPolicy:
+ return
c.Dubbo().NetworkingV1alpha3().FaultInjectionPolicies(cfg.Namespace).Update(context.TODO(),
&apigithubcomapachedubbokubernetesapinetworkingv1alpha3.FaultInjectionPolicy{
+ ObjectMeta: objMeta,
+ Spec:
*(cfg.Spec.(*githubcomkdubboapinetworkingv1alpha3.FaultInjectionPolicy)),
+ }, metav1.UpdateOptions{})
case gvk.GatewayClass:
return
c.GatewayAPI().GatewayV1().GatewayClasses().Update(context.TODO(),
&sigsk8siogatewayapiapisv1.GatewayClass{
ObjectMeta: objMeta,
@@ -184,6 +194,11 @@ func updateStatus(c kube.Client, cfg config.Config,
objMeta metav1.ObjectMeta) (
ObjectMeta: objMeta,
Status:
*(cfg.Status.(*githubcomkdubboapimetav1alpha1.DubboStatus)),
}, metav1.UpdateOptions{})
+ case gvk.FaultInjectionPolicy:
+ return
c.Dubbo().NetworkingV1alpha3().FaultInjectionPolicies(cfg.Namespace).UpdateStatus(context.TODO(),
&apigithubcomapachedubbokubernetesapinetworkingv1alpha3.FaultInjectionPolicy{
+ ObjectMeta: objMeta,
+ Status:
*(cfg.Status.(*githubcomkdubboapimetav1alpha1.DubboStatus)),
+ }, metav1.UpdateOptions{})
case gvk.GatewayClass:
return
c.GatewayAPI().GatewayV1().GatewayClasses().UpdateStatus(context.TODO(),
&sigsk8siogatewayapiapisv1.GatewayClass{
ObjectMeta: objMeta,
@@ -279,6 +294,21 @@ func patch(c kube.Client, orig config.Config, origMeta
metav1.ObjectMeta, mod co
}
return
c.Dubbo().NetworkingV1alpha3().CircuitBreakerPolicies(orig.Namespace).
Patch(context.TODO(), orig.Name, typ, patchBytes,
metav1.PatchOptions{FieldManager: "pilot-discovery"})
+ case gvk.FaultInjectionPolicy:
+ oldRes :=
&apigithubcomapachedubbokubernetesapinetworkingv1alpha3.FaultInjectionPolicy{
+ ObjectMeta: origMeta,
+ Spec:
*(orig.Spec.(*githubcomkdubboapinetworkingv1alpha3.FaultInjectionPolicy)),
+ }
+ modRes :=
&apigithubcomapachedubbokubernetesapinetworkingv1alpha3.FaultInjectionPolicy{
+ ObjectMeta: modMeta,
+ Spec:
*(mod.Spec.(*githubcomkdubboapinetworkingv1alpha3.FaultInjectionPolicy)),
+ }
+ patchBytes, err := genPatchBytes(oldRes, modRes, typ)
+ if err != nil {
+ return nil, err
+ }
+ return
c.Dubbo().NetworkingV1alpha3().FaultInjectionPolicies(orig.Namespace).
+ Patch(context.TODO(), orig.Name, typ, patchBytes,
metav1.PatchOptions{FieldManager: "pilot-discovery"})
case gvk.GatewayClass:
oldRes := &sigsk8siogatewayapiapisv1.GatewayClass{
ObjectMeta: origMeta,
@@ -431,6 +461,8 @@ func delete(c kube.Client, typ config.GroupVersionKind,
name, namespace string,
return
c.GatewayAPI().GatewayV1().BackendTLSPolicies(namespace).Delete(context.TODO(),
name, deleteOptions)
case gvk.CircuitBreakerPolicy:
return
c.Dubbo().NetworkingV1alpha3().CircuitBreakerPolicies(namespace).Delete(context.TODO(),
name, deleteOptions)
+ case gvk.FaultInjectionPolicy:
+ return
c.Dubbo().NetworkingV1alpha3().FaultInjectionPolicies(namespace).Delete(context.TODO(),
name, deleteOptions)
case gvk.GatewayClass:
return
c.GatewayAPI().GatewayV1().GatewayClasses().Delete(context.TODO(), name,
deleteOptions)
case gvk.HTTPRoute:
@@ -620,6 +652,25 @@ var translationMap = map[config.GroupVersionKind]func(r
runtime.Object) config.C
Spec: obj,
}
},
+ gvk.FaultInjectionPolicy: func(r runtime.Object) config.Config {
+ obj :=
r.(*apigithubcomapachedubbokubernetesapinetworkingv1alpha3.FaultInjectionPolicy)
+ return config.Config{
+ Meta: config.Meta{
+ GroupVersionKind: gvk.FaultInjectionPolicy,
+ Name: obj.Name,
+ Namespace: obj.Namespace,
+ Labels: obj.Labels,
+ Annotations: obj.Annotations,
+ ResourceVersion: obj.ResourceVersion,
+ CreationTimestamp: obj.CreationTimestamp.Time,
+ OwnerReferences: obj.OwnerReferences,
+ UID: string(obj.UID),
+ Generation: obj.Generation,
+ },
+ Spec: &obj.Spec,
+ Status: &obj.Status,
+ }
+ },
gvk.GatewayClass: func(r runtime.Object) config.Config {
obj := r.(*sigsk8siogatewayapiapisv1.GatewayClass)
return config.Config{
diff --git a/dubbod/discovery/pkg/model/push_context.go
b/dubbod/discovery/pkg/model/push_context.go
index 99d3affe..40ecc2c8 100644
--- a/dubbod/discovery/pkg/model/push_context.go
+++ b/dubbod/discovery/pkg/model/push_context.go
@@ -40,6 +40,7 @@ import (
"github.com/apache/dubbo-kubernetes/pkg/xds"
meshv1alpha1 "github.com/kdubbo/api/mesh/v1alpha1"
"go.uber.org/atomic"
+ "google.golang.org/protobuf/types/known/wrapperspb"
"k8s.io/apimachinery/pkg/types"
)
@@ -75,6 +76,7 @@ type PushContext struct {
virtualServiceIndex virtualServiceIndex
httpRouteIndex httpRouteIndex
backendTLSPolicyIndex backendTLSPolicyIndex
+ faultInjectionIndex faultInjectionPolicyIndex
destinationRuleIndex destinationRuleIndex
serviceAccounts map[serviceAccountKey][]string
AuthenticationPolicies *AuthenticationPolicies
@@ -591,6 +593,7 @@ func (ps *PushContext) createNewContext(env *Environment) {
ps.initKubernetesGateways(env)
ps.initHTTPRoutes(env)
ps.initBackendTLSPolicies(env)
+ ps.initFaultInjectionPolicies(env)
ps.initAuthenticationPolicies(env)
}
@@ -655,6 +658,14 @@ func (ps *PushContext) updateContext(env *Environment,
oldPushContext *PushConte
ps.backendTLSPolicyIndex = oldPushContext.backendTLSPolicyIndex
}
+ faultInjectionPoliciesChanged := pushReq != nil &&
HasConfigsOfKind(pushReq.ConfigsUpdated, kind.FaultInjectionPolicy)
+ if faultInjectionPoliciesChanged {
+ log.Debugf("FaultInjectionPolicies changed, re-initializing
FaultInjectionPolicy index")
+ ps.initFaultInjectionPolicies(env)
+ } else {
+ ps.faultInjectionIndex = oldPushContext.faultInjectionIndex
+ }
+
authnPoliciesChanged := pushReq != nil && (pushReq.Full ||
authPolicyKindsChanged(pushReq.ConfigsUpdated))
if authnPoliciesChanged || oldPushContext == nil ||
oldPushContext.AuthenticationPolicies == nil {
log.Debugf("security authentication policy changed (full=%v,
configsUpdatedContainingSecurityPolicy=%v), rebuilding authentication policies",
@@ -902,6 +913,85 @@ func (ps *PushContext) BackendTLSForService(namespace,
name string) (BackendTLSS
return settings, found
}
+type FaultInjectionSettings struct {
+ Delay time.Duration
+ DelayPercentage uint32
+ AbortStatus uint32
+ AbortPercentage uint32
+}
+
+type faultInjectionPolicyIndex struct {
+ serviceFaults map[string]FaultInjectionSettings
+}
+
+func (ps *PushContext) initFaultInjectionPolicies(env *Environment) {
+ policies := sortConfigByCreationTime(env.List(gvk.FaultInjectionPolicy,
NamespaceAll))
+ serviceFaults := map[string]FaultInjectionSettings{}
+ for _, cfg := range policies {
+ spec, ok := cfg.Spec.(*networking.FaultInjectionPolicy)
+ if !ok || spec == nil {
+ continue
+ }
+ settings := FaultInjectionSettings{}
+ if delay := spec.GetDelay(); delay != nil {
+ if value := delay.GetFixedDelay(); value != nil &&
value.CheckValid() == nil && value.AsDuration() > 0 {
+ settings.Delay = value.AsDuration()
+ settings.DelayPercentage =
faultPercentage(delay.GetPercentage())
+ }
+ }
+ if abort := spec.GetAbort(); abort != nil &&
abort.GetHttpStatus() >= 400 && abort.GetHttpStatus() <= 599 {
+ settings.AbortStatus = abort.GetHttpStatus()
+ settings.AbortPercentage =
faultPercentage(abort.GetPercentage())
+ }
+ if settings.Delay == 0 && settings.AbortStatus == 0 {
+ continue
+ }
+ for _, target := range spec.GetTargetRefs() {
+ if !isFaultInjectionServiceTarget(target) {
+ continue
+ }
+ key := faultInjectionServiceKey(cfg.Namespace,
target.GetName(), target.GetSectionName())
+ if _, found := serviceFaults[key]; !found {
+ serviceFaults[key] = settings
+ }
+ }
+ }
+ ps.faultInjectionIndex.serviceFaults = serviceFaults
+ log.Debugf("indexed FaultInjectionPolicies for %d service targets",
len(serviceFaults))
+}
+
+func (ps *PushContext) FaultInjectionForService(namespace, name, portName
string) (FaultInjectionSettings, bool) {
+ if ps == nil || ps.faultInjectionIndex.serviceFaults == nil {
+ return FaultInjectionSettings{}, false
+ }
+ if portName != "" {
+ if settings, found :=
ps.faultInjectionIndex.serviceFaults[faultInjectionServiceKey(namespace, name,
portName)]; found {
+ return settings, true
+ }
+ }
+ settings, found :=
ps.faultInjectionIndex.serviceFaults[faultInjectionServiceKey(namespace, name,
"")]
+ return settings, found
+}
+
+func faultPercentage(value *wrapperspb.UInt32Value) uint32 {
+ if value == nil {
+ return 100
+ }
+ return value.GetValue()
+}
+
+func isFaultInjectionServiceTarget(target *networking.PolicyTargetReference)
bool {
+ if target == nil || target.GetName() == "" {
+ return false
+ }
+ group := strings.TrimSpace(target.GetGroup())
+ return (group == "" || group == "core") && target.GetKind() == "Service"
+}
+
+func faultInjectionServiceKey(namespace, name, section string) string {
+ return namespace + "/" + name + "#" + section
+}
+
func supportsSystemBackendTLS(spec
*sigsk8siogatewayapiapisv1.BackendTLSPolicySpec) bool {
if spec == nil || spec.Validation.Hostname == "" {
return false
diff --git a/dubbod/discovery/pkg/networking/grpcgen/rds.go
b/dubbod/discovery/pkg/networking/grpcgen/rds.go
index b01afeaf..46871a5d 100644
--- a/dubbod/discovery/pkg/networking/grpcgen/rds.go
+++ b/dubbod/discovery/pkg/networking/grpcgen/rds.go
@@ -96,8 +96,9 @@ func buildHTTPRoute(node *model.Proxy, push
*model.PushContext, routeName string
}
domains = append(domains, "*") // Wildcard for any domain -
LEAST SPECIFIC
+ faultPolicy := serviceFaultPolicy(push, svc, parsedPort)
outboundRoutes := []*route.Route{
- defaultSingleClusterRoute(routeName),
+ defaultSingleClusterRoute(routeName, faultPolicy),
}
if node.IsRouter() {
@@ -136,7 +137,7 @@ func buildHTTPRoute(node *model.Proxy, push
*model.PushContext, routeName string
}
}
- if routes :=
buildRoutesFromGatewayHTTPRoute(httpRoutes, host.Name("*"), parsedPort);
len(routes) > 0 {
+ if routes :=
buildRoutesFromGatewayHTTPRoute(httpRoutes, host.Name("*"), parsedPort,
faultPolicy); len(routes) > 0 {
log.Infof("built %d routes from Gateway
API HTTPRoute", len(routes))
outboundRoutes = routes
} else {
@@ -145,7 +146,7 @@ func buildHTTPRoute(node *model.Proxy, push
*model.PushContext, routeName string
}
} else if httpRoutes :=
filterHTTPRoutesByService(push.HTTPRouteForHost(host.Name(hostStr)), svc,
parsedPort); len(httpRoutes) > 0 {
log.Infof("found %d service-attached HTTPRoute(s) for
host %s", len(httpRoutes), hostStr)
- if routes :=
buildRoutesFromGatewayHTTPRoute(httpRoutes, host.Name(hostStr), parsedPort);
len(routes) > 0 {
+ if routes :=
buildRoutesFromGatewayHTTPRoute(httpRoutes, host.Name(hostStr), parsedPort,
faultPolicy); len(routes) > 0 {
log.Infof("built %d routes from
service-attached HTTPRoute for host %s", len(routes), hostStr)
outboundRoutes = routes
} else {
@@ -257,7 +258,7 @@ func buildHTTPRoute(node *model.Proxy, push
*model.PushContext, routeName string
}
}
- if routes :=
buildRoutesFromGatewayHTTPRoute(httpRoutes, host.Name("*"), parsedPort);
len(routes) > 0 {
+ if routes :=
buildRoutesFromGatewayHTTPRoute(httpRoutes, host.Name("*"), parsedPort, nil);
len(routes) > 0 {
log.Infof("Gateway Pod inbound listener built
%d routes from HTTPRoute", len(routes))
outboundRoutes = routes
} else {
@@ -313,7 +314,7 @@ func buildHTTPRoute(node *model.Proxy, push
*model.PushContext, routeName string
}
}
-func defaultSingleClusterRoute(clusterName string) *route.Route {
+func defaultSingleClusterRoute(clusterName string, faultPolicy
*route.FaultPolicy) *route.Route {
return &route.Route{
Match: &route.RouteMatch{
PathSpecifier: &route.RouteMatch_Prefix{
@@ -325,13 +326,14 @@ func defaultSingleClusterRoute(clusterName string)
*route.Route {
ClusterSpecifier: &route.RouteAction_Cluster{
Cluster: clusterName,
},
+ FaultPolicy: faultPolicy,
},
},
}
}
// buildRoutesFromGatewayHTTPRoute converts Gateway API HTTPRoute resources to
XDS Route configurations
-func buildRoutesFromGatewayHTTPRoute(httpRoutes []config.Config, hostName
host.Name, defaultPort int) []*route.Route {
+func buildRoutesFromGatewayHTTPRoute(httpRoutes []config.Config, hostName
host.Name, defaultPort int, faultPolicy *route.FaultPolicy) []*route.Route {
if len(httpRoutes) == 0 {
return nil
}
@@ -418,6 +420,7 @@ func buildRoutesFromGatewayHTTPRoute(httpRoutes
[]config.Config, hostName host.N
}
}
routeAction.RetryPolicy =
gatewayAPIRetryPolicy(rule.Retry, rule.Timeouts)
+ routeAction.FaultPolicy = faultPolicy
builtRoute := &route.Route{
Match: routeMatch,
@@ -435,6 +438,37 @@ func buildRoutesFromGatewayHTTPRoute(httpRoutes
[]config.Config, hostName host.N
return allRoutes
}
+func serviceFaultPolicy(push *model.PushContext, svc *model.Service, port int)
*route.FaultPolicy {
+ if push == nil || svc == nil {
+ return nil
+ }
+ portName := ""
+ if servicePort, found := svc.Ports.GetByPort(port); found &&
servicePort != nil {
+ portName = servicePort.Name
+ }
+ settings, found :=
push.FaultInjectionForService(svc.Attributes.Namespace, svc.Attributes.Name,
portName)
+ if !found {
+ return nil
+ }
+ policy := &route.FaultPolicy{}
+ if settings.Delay > 0 {
+ policy.Delay = &route.FaultDelay{
+ FixedDelay: durationpb.New(settings.Delay),
+ Percentage: wrapperspb.UInt32(settings.DelayPercentage),
+ }
+ }
+ if settings.AbortStatus != 0 {
+ policy.Abort = &route.FaultAbort{
+ HttpStatus: settings.AbortStatus,
+ Percentage: wrapperspb.UInt32(settings.AbortPercentage),
+ }
+ }
+ if policy.Delay == nil && policy.Abort == nil {
+ return nil
+ }
+ return policy
+}
+
func gatewayAPIRetryPolicy(retry *sigsk8siogatewayapiapisv1.HTTPRouteRetry,
timeouts *sigsk8siogatewayapiapisv1.HTTPRouteTimeouts) *route.RetryPolicy {
if retry == nil {
return nil
diff --git a/dubbod/discovery/pkg/networking/grpcgen/rds_test.go
b/dubbod/discovery/pkg/networking/grpcgen/rds_test.go
index 526fca2f..dfca75a9 100644
--- a/dubbod/discovery/pkg/networking/grpcgen/rds_test.go
+++ b/dubbod/discovery/pkg/networking/grpcgen/rds_test.go
@@ -30,7 +30,10 @@ import (
"github.com/apache/dubbo-kubernetes/pkg/config/schema/collections"
"github.com/apache/dubbo-kubernetes/pkg/config/schema/gvk"
"github.com/apache/dubbo-kubernetes/pkg/kube/krt"
+ networking "github.com/kdubbo/api/networking/v1alpha3"
route "github.com/kdubbo/xds-api/route/v1"
+ "google.golang.org/protobuf/types/known/durationpb"
+ "google.golang.org/protobuf/types/known/wrapperspb"
gatewayv1 "sigs.k8s.io/gateway-api/apis/v1"
)
@@ -185,6 +188,62 @@ func TestBuildHTTPRouteSetsGatewayAPIRetryPolicy(t
*testing.T) {
}
}
+func TestBuildHTTPRouteSetsServiceFaultInjectionPolicy(t *testing.T) {
+ routeConfig := newServiceAttachedHTTPRouteConfig("reviews-fault",
"moviereview", "reviews", 9080)
+ faultConfig := config.Config{
+ Meta: config.Meta{
+ GroupVersionKind: gvk.FaultInjectionPolicy,
+ Name: "reviews-fault",
+ Namespace: "moviereview",
+ },
+ Spec: &networking.FaultInjectionPolicy{
+ TargetRefs: []*networking.PolicyTargetReference{{
+ Kind: "Service",
+ Name: "reviews",
+ SectionName: "http",
+ }},
+ Delay: &networking.FaultDelay{
+ FixedDelay: durationpb.New(250 *
time.Millisecond),
+ Percentage: wrapperspb.UInt32(20),
+ },
+ Abort: &networking.FaultAbort{
+ HttpStatus: 503,
+ Percentage: wrapperspb.UInt32(10),
+ },
+ },
+ }
+ push := newRDSTestPushContext(t, []config.Config{routeConfig,
faultConfig}, []*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")
+ }
+ fault := rc.VirtualHosts[0].Routes[0].GetRoute().GetFaultPolicy()
+ if fault == nil {
+ t.Fatal("fault policy = nil")
+ }
+ if got := fault.GetDelay().GetFixedDelay().AsDuration(); got !=
250*time.Millisecond {
+ t.Fatalf("fixed delay = %v, want 250ms", got)
+ }
+ if got := fault.GetDelay().GetPercentage().GetValue(); got != 20 {
+ t.Fatalf("delay percentage = %d, want 20", got)
+ }
+ if got := fault.GetAbort().GetHttpStatus(); got != 503 {
+ t.Fatalf("abort status = %d, want 503", got)
+ }
+ if got := fault.GetAbort().GetPercentage().GetValue(); got != 10 {
+ t.Fatalf("abort percentage = %d, want 10", 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 18e16210..c10abf4c 100644
--- a/go.mod
+++ b/go.mod
@@ -52,9 +52,9 @@ require (
github.com/hashicorp/go-multierror v1.1.1
github.com/hashicorp/golang-lru/v2 v2.0.7
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-20260728021336-34287e74f6f2
+ github.com/kdubbo/api v0.0.0-20260728161804-a5971782efe0
+ github.com/kdubbo/client-go v0.0.0-20260729004545-0427ad75f167
+ github.com/kdubbo/xds-api v0.0.0-20260728161804-af6dbc11367a
github.com/moby/moby/client v0.4.1
github.com/moby/term v0.5.2
github.com/ory/viper v1.7.5
diff --git a/go.sum b/go.sum
index 8c160efd..26b0dbb9 100644
--- a/go.sum
+++ b/go.sum
@@ -368,12 +368,12 @@ github.com/jtolds/gls v4.20.0+incompatible/go.mod
h1:QJZ7F/aHp+rZTRtaJ1ow/lLfFfV
github.com/julienschmidt/httprouter v1.2.0/go.mod
h1:SYymIcj16QtmaHHD7aYtjjsJG7VTCxuUUipMqKk8s4w=
github.com/kballard/go-shellquote v0.0.0-20180428030007-95032a82bc51
h1:Z9n2FFNUXsshfwJMBgNA0RU6/i7WVaAegv3PtuIHPMs=
github.com/kballard/go-shellquote v0.0.0-20180428030007-95032a82bc51/go.mod
h1:CzGEWj7cYgsdH8dAjBGEr58BoE7ScuLd+fwFZ44+/x8=
-github.com/kdubbo/api v0.0.0-20260713105721-558a18fe4f33
h1://TxoBFL/igcCqokYHOcA5hA30J6LTpFusDYVQx48Nw=
-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-20260728021336-34287e74f6f2
h1:8tFRMU/dvoQJLaPKM/neaXjRhV/5FfxEjLy0uFJ0KG8=
-github.com/kdubbo/xds-api v0.0.0-20260728021336-34287e74f6f2/go.mod
h1:o2HDUgL1ntaDbWomZ4cD2tt8jBamuG2qRtjXOa1zZ0Q=
+github.com/kdubbo/api v0.0.0-20260728161804-a5971782efe0
h1:t4bS0hQfyiu3QdsL+frOw/KbOjmvQERPTaziWHhLc3U=
+github.com/kdubbo/api v0.0.0-20260728161804-a5971782efe0/go.mod
h1:8BtJiIovg7QCPsCxXcw3gDf922VcvYq5ihOSvj49Rq8=
+github.com/kdubbo/client-go v0.0.0-20260729004545-0427ad75f167
h1:j5nY/UzzGztfLEPrXbNvohGO53PVaRhoIK1GrtrikOs=
+github.com/kdubbo/client-go v0.0.0-20260729004545-0427ad75f167/go.mod
h1:/wrQJoD+yhTTBi/5sC8rHLK+62mk4vOnLaYZ1Sf+DOQ=
+github.com/kdubbo/xds-api v0.0.0-20260728161804-af6dbc11367a
h1:WfpZeq43xfNY+6icInQxNSNDNLsbDj3QebyvLD9OmL8=
+github.com/kdubbo/xds-api v0.0.0-20260728161804-af6dbc11367a/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=
diff --git a/manifests/charts/base/files/crd-all.gen.yaml
b/manifests/charts/base/files/crd-all.gen.yaml
index 422dd049..ca9761c1 100644
--- a/manifests/charts/base/files/crd-all.gen.yaml
+++ b/manifests/charts/base/files/crd-all.gen.yaml
@@ -152,6 +152,142 @@ spec:
---
apiVersion: apiextensions.k8s.io/v1
kind: CustomResourceDefinition
+metadata:
+ annotations:
+ "helm.sh/resource-policy": keep
+ labels:
+ app: dubbo
+ chart: dubbo
+ dubbo: networking
+ heritage: Tiller
+ release: dubbo
+ name: faultinjectionpolicies.networking.dubbo.apache.org
+spec:
+ group: networking.dubbo.apache.org
+ names:
+ categories:
+ - dubbo
+ - networking
+ kind: FaultInjectionPolicy
+ listKind: FaultInjectionPolicyList
+ plural: faultinjectionpolicies
+ shortNames:
+ - fip
+ singular: faultinjectionpolicy
+ scope: Namespaced
+ versions:
+ - additionalPrinterColumns:
+ - description: Gateway API targets.
+ jsonPath: .spec.targetRefs[*].name
+ name: Targets
+ type: string
+ - description: CreationTimestamp is a timestamp representing the server
time when
+ this object was created.
+ jsonPath: .metadata.creationTimestamp
+ name: Age
+ type: date
+ name: v1alpha3
+ schema:
+ openAPIV3Schema:
+ properties:
+ spec:
+ description: 'Gateway API policy attachment for controlled service
faults.
+ See more details at: '
+ properties:
+ abort:
+ description: Abort terminates a matching request before it
reaches
+ the application.
+ properties:
+ httpStatus:
+ description: HTTP status returned by an L7 data plane.
+ maximum: 599
+ minimum: 400
+ type: integer
+ percentage:
+ description: Percentage of requests or connections
affected, from
+ 0 through 100.
+ maximum: 100
+ minimum: 0
+ nullable: true
+ type: integer
+ required:
+ - httpStatus
+ type: object
+ delay:
+ description: Delay pauses a matching request or connection
before
+ forwarding.
+ properties:
+ fixedDelay:
+ description: Fixed delay applied when this fault is
selected.
+ type: string
+ x-kubernetes-validations:
+ - message: must be a valid duration greater than 1ms
+ rule: duration(self) >= duration('1ms')
+ percentage:
+ description: Percentage of requests or connections
affected, from
+ 0 through 100.
+ maximum: 100
+ minimum: 0
+ nullable: true
+ type: integer
+ required:
+ - fixedDelay
+ type: object
+ targetRefs:
+ description: Gateway API policy targets.
+ items:
+ properties:
+ group:
+ description: API group of the target.
+ type: string
+ kind:
+ description: Kind of the target.
+ type: string
+ name:
+ description: Name of the target object.
+ type: string
+ sectionName:
+ description: Optional section name for future per-port
attachment.
+ type: string
+ required:
+ - kind
+ - name
+ type: object
+ type: array
+ required:
+ - targetRefs
+ type: object
+ x-kubernetes-validations:
+ - message: at least one of delay or abort must be set
+ rule: has(self.delay) || has(self.abort)
+ status:
+ properties:
+ conditions:
+ items:
+ properties:
+ observedGeneration:
+ anyOf:
+ - type: integer
+ - type: string
+ x-kubernetes-int-or-string: true
+ reason:
+ type: string
+ status:
+ type: string
+ type:
+ type: string
+ type: object
+ type: array
+ type: object
+ x-kubernetes-preserve-unknown-fields: true
+ type: object
+ served: true
+ storage: true
+ subresources:
+ status: {}
+---
+apiVersion: apiextensions.k8s.io/v1
+kind: CustomResourceDefinition
metadata:
annotations:
"helm.sh/resource-policy": keep
diff --git a/pkg/config/schema/collections/collections.gen.go
b/pkg/config/schema/collections/collections.gen.go
index de2c2a0a..6eb97852 100755
--- a/pkg/config/schema/collections/collections.gen.go
+++ b/pkg/config/schema/collections/collections.gen.go
@@ -167,6 +167,21 @@ var (
ValidateProto: validation.EmptyValidate,
}.MustBuild()
+ FaultInjectionPolicy = resource.Builder{
+ Identifier: "FaultInjectionPolicy",
+ Group: "networking.dubbo.apache.org",
+ Kind: "FaultInjectionPolicy",
+ Plural: "faultinjectionpolicies",
+ Version: "v1alpha3",
+ Proto: "dubbo.networking.v1alpha3.FaultInjectionPolicy",
StatusProto: "dubbo.meta.v1alpha1.DubboStatus",
+ ReflectType:
reflect.TypeOf(&githubcomkdubboapinetworkingv1alpha3.FaultInjectionPolicy{}).Elem(),
StatusType:
reflect.TypeOf(&githubcomkdubboapimetav1alpha1.DubboStatus{}).Elem(),
+ ProtoPackage: "github.com/kdubbo/api/networking/v1alpha3",
StatusPackage: "github.com/kdubbo/api/meta/v1alpha1",
+ ClusterScoped: false,
+ Synthetic: false,
+ Builtin: false,
+ ValidateProto: validation.ValidateFaultInjectionPolicy,
+ }.MustBuild()
+
GatewayClass = resource.Builder{
Identifier: "GatewayClass",
Group: "gateway.networking.k8s.io",
@@ -520,6 +535,7 @@ var (
MustAdd(Deployment).
MustAdd(EndpointSlice).
MustAdd(Endpoints).
+ MustAdd(FaultInjectionPolicy).
MustAdd(GatewayClass).
MustAdd(HTTPRoute).
MustAdd(HorizontalPodAutoscaler).
@@ -575,6 +591,7 @@ var (
Dubbo = collection.NewSchemasBuilder().
MustAdd(AuthorizationPolicy).
MustAdd(CircuitBreakerPolicy).
+ MustAdd(FaultInjectionPolicy).
MustAdd(PeerAuthentication).
MustAdd(RequestAuthentication).
MustAdd(ServiceEntry).
@@ -587,6 +604,7 @@ var (
MustAdd(AuthorizationPolicy).
MustAdd(BackendTLSPolicy).
MustAdd(CircuitBreakerPolicy).
+ MustAdd(FaultInjectionPolicy).
MustAdd(GatewayClass).
MustAdd(HTTPRoute).
MustAdd(KubernetesGateway).
@@ -603,6 +621,7 @@ var (
MustAdd(AuthorizationPolicy).
MustAdd(BackendTLSPolicy).
MustAdd(CircuitBreakerPolicy).
+ MustAdd(FaultInjectionPolicy).
MustAdd(GatewayClass).
MustAdd(HTTPRoute).
MustAdd(KubernetesGateway).
diff --git a/pkg/config/schema/gvk/resources.gen.go
b/pkg/config/schema/gvk/resources.gen.go
index 8cd44e99..27a3ccbb 100755
--- a/pkg/config/schema/gvk/resources.gen.go
+++ b/pkg/config/schema/gvk/resources.gen.go
@@ -20,6 +20,7 @@ var (
Deployment = config.GroupVersionKind{Group: "apps",
Version: "v1", Kind: "Deployment"}
EndpointSlice = config.GroupVersionKind{Group:
"discovery.k8s.io", Version: "v1", Kind: "EndpointSlice"}
Endpoints = config.GroupVersionKind{Group: "",
Version: "v1", Kind: "Endpoints"}
+ FaultInjectionPolicy = config.GroupVersionKind{Group:
"networking.dubbo.apache.org", Version: "v1alpha3", Kind:
"FaultInjectionPolicy"}
GatewayClass = config.GroupVersionKind{Group:
"gateway.networking.k8s.io", Version: "v1", Kind: "GatewayClass"}
GatewayClass_v1 = config.GroupVersionKind{Group:
"gateway.networking.k8s.io", Version: "v1", Kind: "GatewayClass"}
HTTPRoute = config.GroupVersionKind{Group:
"gateway.networking.k8s.io", Version: "v1", Kind: "HTTPRoute"}
@@ -71,6 +72,8 @@ func ToGVR(g config.GroupVersionKind)
(schema.GroupVersionResource, bool) {
return gvr.EndpointSlice, true
case Endpoints:
return gvr.Endpoints, true
+ case FaultInjectionPolicy:
+ return gvr.FaultInjectionPolicy, true
case GatewayClass:
return gvr.GatewayClass, true
case GatewayClass_v1:
@@ -148,6 +151,8 @@ func MustToKind(g config.GroupVersionKind) kind.Kind {
return kind.EndpointSlice
case Endpoints:
return kind.Endpoints
+ case FaultInjectionPolicy:
+ return kind.FaultInjectionPolicy
case GatewayClass:
return kind.GatewayClass
case HTTPRoute:
@@ -228,6 +233,8 @@ func FromGVR(g schema.GroupVersionResource)
(config.GroupVersionKind, bool) {
return EndpointSlice, true
case gvr.Endpoints:
return Endpoints, true
+ case gvr.FaultInjectionPolicy:
+ return FaultInjectionPolicy, true
case gvr.GatewayClass:
return GatewayClass, true
case gvr.HTTPRoute:
diff --git a/pkg/config/schema/gvr/resources.gen.go
b/pkg/config/schema/gvr/resources.gen.go
index 62f90f2a..e9188777 100755
--- a/pkg/config/schema/gvr/resources.gen.go
+++ b/pkg/config/schema/gvr/resources.gen.go
@@ -15,6 +15,7 @@ var (
Deployment = schema.GroupVersionResource{Group:
"apps", Version: "v1", Resource: "deployments"}
EndpointSlice = schema.GroupVersionResource{Group:
"discovery.k8s.io", Version: "v1", Resource: "endpointslices"}
Endpoints = schema.GroupVersionResource{Group: "",
Version: "v1", Resource: "endpoints"}
+ FaultInjectionPolicy = schema.GroupVersionResource{Group:
"networking.dubbo.apache.org", Version: "v1alpha3", Resource:
"faultinjectionpolicies"}
GatewayClass = schema.GroupVersionResource{Group:
"gateway.networking.k8s.io", Version: "v1", Resource: "gatewayclasses"}
GatewayClass_v1 = schema.GroupVersionResource{Group:
"gateway.networking.k8s.io", Version: "v1", Resource: "gatewayclasses"}
HTTPRoute = schema.GroupVersionResource{Group:
"gateway.networking.k8s.io", Version: "v1", Resource: "httproutes"}
@@ -65,6 +66,8 @@ func IsClusterScoped(g schema.GroupVersionResource) bool {
return false
case Endpoints:
return false
+ case FaultInjectionPolicy:
+ return false
case GatewayClass:
return true
case GatewayClass_v1:
diff --git a/pkg/config/schema/kind/resources.gen.go
b/pkg/config/schema/kind/resources.gen.go
index 97193ecb..5e4e84b3 100755
--- a/pkg/config/schema/kind/resources.gen.go
+++ b/pkg/config/schema/kind/resources.gen.go
@@ -15,6 +15,7 @@ const (
Deployment
EndpointSlice
Endpoints
+ FaultInjectionPolicy
GatewayClass
HTTPRoute
HorizontalPodAutoscaler
@@ -63,6 +64,8 @@ func (k Kind) String() string {
return "EndpointSlice"
case Endpoints:
return "Endpoints"
+ case FaultInjectionPolicy:
+ return "FaultInjectionPolicy"
case GatewayClass:
return "GatewayClass"
case HTTPRoute:
@@ -136,6 +139,8 @@ func FromString(s string) Kind {
return EndpointSlice
case "Endpoints":
return Endpoints
+ case "FaultInjectionPolicy":
+ return FaultInjectionPolicy
case "GatewayClass":
return GatewayClass
case "HTTPRoute":
diff --git a/pkg/config/schema/kubeclient/resources.gen.go
b/pkg/config/schema/kubeclient/resources.gen.go
index e63b7a7c..c8ca80d3 100755
--- a/pkg/config/schema/kubeclient/resources.gen.go
+++ b/pkg/config/schema/kubeclient/resources.gen.go
@@ -52,6 +52,8 @@ func GetWriteClient[T runtime.Object](c ClientGetter,
namespace string) ktypes.W
return
c.Kube().DiscoveryV1().EndpointSlices(namespace).(ktypes.WriteAPI[T])
case *k8sioapicorev1.Endpoints:
return
c.Kube().CoreV1().Endpoints(namespace).(ktypes.WriteAPI[T])
+ case
*apigithubcomapachedubbokubernetesapinetworkingv1alpha3.FaultInjectionPolicy:
+ return
c.Dubbo().NetworkingV1alpha3().FaultInjectionPolicies(namespace).(ktypes.WriteAPI[T])
case *sigsk8siogatewayapiapisv1.GatewayClass:
return
c.GatewayAPI().GatewayV1().GatewayClasses().(ktypes.WriteAPI[T])
case *sigsk8siogatewayapiapisv1.HTTPRoute:
@@ -119,6 +121,8 @@ func GetClient[T, TL runtime.Object](c ClientGetter,
namespace string) ktypes.Re
return
c.Kube().DiscoveryV1().EndpointSlices(namespace).(ktypes.ReadWriteAPI[T, TL])
case *k8sioapicorev1.Endpoints:
return
c.Kube().CoreV1().Endpoints(namespace).(ktypes.ReadWriteAPI[T, TL])
+ case
*apigithubcomapachedubbokubernetesapinetworkingv1alpha3.FaultInjectionPolicy:
+ return
c.Dubbo().NetworkingV1alpha3().FaultInjectionPolicies(namespace).(ktypes.ReadWriteAPI[T,
TL])
case *sigsk8siogatewayapiapisv1.GatewayClass:
return
c.GatewayAPI().GatewayV1().GatewayClasses().(ktypes.ReadWriteAPI[T, TL])
case *sigsk8siogatewayapiapisv1.HTTPRoute:
@@ -186,6 +190,8 @@ func gvrToObject(g schema.GroupVersionResource)
runtime.Object {
return &k8sioapidiscoveryv1.EndpointSlice{}
case gvr.Endpoints:
return &k8sioapicorev1.Endpoints{}
+ case gvr.FaultInjectionPolicy:
+ return
&apigithubcomapachedubbokubernetesapinetworkingv1alpha3.FaultInjectionPolicy{}
case gvr.GatewayClass:
return &sigsk8siogatewayapiapisv1.GatewayClass{}
case gvr.HTTPRoute:
@@ -301,6 +307,13 @@ func getInformerFiltered(c ClientGetter, opts
ktypes.InformerOptions, g schema.G
w = func(options metav1.ListOptions) (watch.Interface, error) {
return
c.Kube().CoreV1().Endpoints(opts.Namespace).Watch(context.Background(), options)
}
+ case gvr.FaultInjectionPolicy:
+ l = func(options metav1.ListOptions) (runtime.Object, error) {
+ return
c.Dubbo().NetworkingV1alpha3().FaultInjectionPolicies(opts.Namespace).List(context.Background(),
options)
+ }
+ w = func(options metav1.ListOptions) (watch.Interface, error) {
+ return
c.Dubbo().NetworkingV1alpha3().FaultInjectionPolicies(opts.Namespace).Watch(context.Background(),
options)
+ }
case gvr.GatewayClass:
l = func(options metav1.ListOptions) (runtime.Object, error) {
return
c.GatewayAPI().GatewayV1().GatewayClasses().List(context.Background(), options)
diff --git a/pkg/config/schema/kubetypes/resources.gen.go
b/pkg/config/schema/kubetypes/resources.gen.go
index 864ee1ad..544ef07a 100755
--- a/pkg/config/schema/kubetypes/resources.gen.go
+++ b/pkg/config/schema/kubetypes/resources.gen.go
@@ -48,6 +48,10 @@ func getGvk(obj any) (config.GroupVersionKind, bool) {
return gvk.EndpointSlice, true
case *k8sioapicorev1.Endpoints:
return gvk.Endpoints, true
+ case *githubcomkdubboapinetworkingv1alpha3.FaultInjectionPolicy:
+ return gvk.FaultInjectionPolicy, true
+ case
*apigithubcomapachedubbokubernetesapinetworkingv1alpha3.FaultInjectionPolicy:
+ return gvk.FaultInjectionPolicy, true
case *sigsk8siogatewayapiapisv1.GatewayClass:
return gvk.GatewayClass, true
case *sigsk8siogatewayapiapisv1.HTTPRoute:
diff --git a/pkg/config/schema/metadata.yaml b/pkg/config/schema/metadata.yaml
index 68d53812..a81a3356 100644
--- a/pkg/config/schema/metadata.yaml
+++ b/pkg/config/schema/metadata.yaml
@@ -230,6 +230,16 @@ resources:
statusProto: "dubbo.meta.v1alpha1.DubboStatus"
statusProtoPackage: "github.com/kdubbo/api/meta/v1alpha1"
+ - kind: FaultInjectionPolicy
+ plural: "faultinjectionpolicies"
+ group: "networking.dubbo.apache.org"
+ version: "v1alpha3"
+ proto: "dubbo.networking.v1alpha3.FaultInjectionPolicy"
+ protoPackage: "github.com/kdubbo/api/networking/v1alpha3"
+ validate: "validation.ValidateFaultInjectionPolicy"
+ statusProto: "dubbo.meta.v1alpha1.DubboStatus"
+ statusProtoPackage: "github.com/kdubbo/api/meta/v1alpha1"
+
- kind: ServiceEntry
plural: "serviceentries"
group: "networking.dubbo.apache.org"
diff --git a/pkg/config/validation/helpers.go b/pkg/config/validation/helpers.go
index ee624bfe..0c80f9c0 100644
--- a/pkg/config/validation/helpers.go
+++ b/pkg/config/validation/helpers.go
@@ -21,6 +21,7 @@ import (
"net/url"
"google.golang.org/protobuf/types/known/durationpb"
+ "google.golang.org/protobuf/types/known/wrapperspb"
"k8s.io/apimachinery/pkg/util/validation"
"github.com/apache/dubbo-kubernetes/pkg/config"
@@ -135,6 +136,16 @@ func validatePercent(name string, val int32) error {
return nil
}
+func validateOptionalUInt32Percent(name string, val *wrapperspb.UInt32Value)
error {
+ if val == nil {
+ return nil
+ }
+ if value := val.GetValue(); value > 100 {
+ return fmt.Errorf("%s must be in range [0, 100], got %d", name,
value)
+ }
+ return nil
+}
+
func validateNonNegativeInt32(name string, val int32) error {
if val < 0 {
return fmt.Errorf("%s must not be negative, got %d", name, val)
diff --git a/pkg/config/validation/validators.go
b/pkg/config/validation/validators.go
index 7d33922f..b5b61739 100644
--- a/pkg/config/validation/validators.go
+++ b/pkg/config/validation/validators.go
@@ -21,6 +21,7 @@ import (
"math"
"net"
"strings"
+ "time"
"github.com/apache/dubbo-kubernetes/pkg/config"
"github.com/apache/dubbo-kubernetes/pkg/config/constants"
@@ -213,6 +214,55 @@ var ValidateCircuitBreakerPolicy = validateFunc(
return v.Unwrap()
})
+// ValidateFaultInjectionPolicy checks that a FaultInjectionPolicy is safe and
executable.
+var ValidateFaultInjectionPolicy =
RegisterValidateFunc("ValidateFaultInjectionPolicy",
+ func(cfg config.Config) (Warning, error) {
+ spec, ok := cfg.Spec.(*networking.FaultInjectionPolicy)
+ if !ok {
+ return nil, fmt.Errorf("cannot cast to
FaultInjectionPolicy")
+ }
+ v := Validation{}
+ if len(spec.GetTargetRefs()) == 0 {
+ v = appendValidation(v, fmt.Errorf("targetRefs must not
be empty"))
+ }
+ for i, ref := range spec.GetTargetRefs() {
+ if ref == nil {
+ v = appendValidation(v,
fmt.Errorf("targetRefs[%d] must not be null", i))
+ continue
+ }
+ if ref.GetKind() != "Service" {
+ v = appendValidation(v,
fmt.Errorf("targetRefs[%d].kind %q is not supported; only Service targets are
applied", i, ref.GetKind()))
+ }
+ if group := strings.TrimSpace(ref.GetGroup()); group !=
"" && group != "core" {
+ v = appendValidation(v,
fmt.Errorf("targetRefs[%d].group %q is not supported; use the core API group",
i, group))
+ }
+ if ref.GetName() == "" {
+ v = appendValidation(v,
fmt.Errorf("targetRefs[%d].name must not be empty", i))
+ }
+ }
+ if spec.GetDelay() == nil && spec.GetAbort() == nil {
+ v = appendValidation(v, fmt.Errorf("at least one of
delay or abort must be set"))
+ }
+ if delay := spec.GetDelay(); delay != nil {
+ v = appendValidation(v,
+ validatePositiveDuration("delay.fixedDelay",
delay.GetFixedDelay()),
+
validateOptionalUInt32Percent("delay.percentage", delay.GetPercentage()),
+ )
+ if delay.GetFixedDelay() == nil {
+ v = appendValidation(v,
fmt.Errorf("delay.fixedDelay must be set"))
+ } else if delay.GetFixedDelay().CheckValid() == nil &&
delay.GetFixedDelay().AsDuration() < time.Millisecond {
+ v = appendValidation(v,
fmt.Errorf("delay.fixedDelay must be at least 1ms"))
+ }
+ }
+ if abort := spec.GetAbort(); abort != nil {
+ if status := abort.GetHttpStatus(); status < 400 ||
status > 599 {
+ v = appendValidation(v,
fmt.Errorf("abort.httpStatus must be in range [400, 599], got %d", status))
+ }
+ v = appendValidation(v,
validateOptionalUInt32Percent("abort.percentage", abort.GetPercentage()))
+ }
+ return v.Unwrap()
+ })
+
// ValidateServiceEntry checks that a ServiceEntry can be converted into
services and endpoints.
var ValidateServiceEntry = RegisterValidateFunc("ValidateServiceEntry",
func(cfg config.Config) (Warning, error) {
spec, ok := cfg.Spec.(*networking.ServiceEntry)
diff --git a/pkg/config/validation/validators_test.go
b/pkg/config/validation/validators_test.go
index b661f062..08dce1af 100644
--- a/pkg/config/validation/validators_test.go
+++ b/pkg/config/validation/validators_test.go
@@ -17,7 +17,9 @@
package validation
import (
+ "net/http"
"testing"
+ "time"
"google.golang.org/protobuf/types/known/durationpb"
"google.golang.org/protobuf/types/known/wrapperspb"
@@ -295,6 +297,91 @@ func TestValidateCircuitBreakerPolicy(t *testing.T) {
}
}
+func TestValidateFaultInjectionPolicy(t *testing.T) {
+ validRef := []*networking.PolicyTargetReference{{Kind: "Service", Name:
"backend"}}
+ cases := []struct {
+ name string
+ spec *networking.FaultInjectionPolicy
+ wantErr bool
+ }{
+ {
+ name: "delay and abort",
+ spec: &networking.FaultInjectionPolicy{
+ TargetRefs: validRef,
+ Delay: &networking.FaultDelay{
+ FixedDelay: durationpb.New(time.Second),
+ Percentage: wrapperspb.UInt32(25),
+ },
+ Abort: &networking.FaultAbort{
+ HttpStatus:
http.StatusServiceUnavailable,
+ Percentage: wrapperspb.UInt32(10),
+ },
+ },
+ },
+ {
+ name: "no target refs",
+ spec: &networking.FaultInjectionPolicy{Delay:
&networking.FaultDelay{FixedDelay: durationpb.New(time.Second)}},
+ wantErr: true,
+ },
+ {
+ name: "unsupported target",
+ spec: &networking.FaultInjectionPolicy{
+ TargetRefs:
[]*networking.PolicyTargetReference{{Kind: "Gateway", Name: "edge"}},
+ Abort: &networking.FaultAbort{HttpStatus:
http.StatusServiceUnavailable},
+ },
+ wantErr: true,
+ },
+ {
+ name: "no fault",
+ spec: &networking.FaultInjectionPolicy{TargetRefs:
validRef},
+ wantErr: true,
+ },
+ {
+ name: "invalid delay",
+ spec: &networking.FaultInjectionPolicy{
+ TargetRefs: validRef,
+ Delay: &networking.FaultDelay{FixedDelay:
durationpb.New(-time.Second)},
+ },
+ wantErr: true,
+ },
+ {
+ name: "delay below CRD minimum",
+ spec: &networking.FaultInjectionPolicy{
+ TargetRefs: validRef,
+ Delay: &networking.FaultDelay{FixedDelay:
durationpb.New(time.Microsecond)},
+ },
+ wantErr: true,
+ },
+ {
+ name: "invalid abort status",
+ spec: &networking.FaultInjectionPolicy{
+ TargetRefs: validRef,
+ Abort: &networking.FaultAbort{HttpStatus:
http.StatusOK},
+ },
+ wantErr: true,
+ },
+ {
+ name: "invalid percentage",
+ spec: &networking.FaultInjectionPolicy{
+ TargetRefs: validRef,
+ Abort: &networking.FaultAbort{
+ HttpStatus:
http.StatusServiceUnavailable,
+ Percentage: wrapperspb.UInt32(101),
+ },
+ },
+ wantErr: true,
+ },
+ }
+ for _, tc := range cases {
+ t.Run(tc.name, func(t *testing.T) {
+ _, err :=
ValidateFaultInjectionPolicy(makeConfig(tc.spec))
+ if (err != nil) != tc.wantErr {
+ t.Fatalf("got err=%v, wantErr=%v", err,
tc.wantErr)
+ }
+ })
+ }
+}
+
func TestValidateServiceEntry(t *testing.T) {
valid := func() *networking.ServiceEntry {
return &networking.ServiceEntry{
diff --git a/samples/httpbin/fault-injection.yaml
b/samples/httpbin/fault-injection.yaml
new file mode 100644
index 00000000..1560dabc
--- /dev/null
+++ b/samples/httpbin/fault-injection.yaml
@@ -0,0 +1,15 @@
+apiVersion: networking.dubbo.apache.org/v1alpha3
+kind: FaultInjectionPolicy
+metadata:
+ name: httpbin-fault-injection
+spec:
+ targetRefs:
+ - group: ""
+ kind: Service
+ name: httpbin
+ delay:
+ fixedDelay: 250ms
+ percentage: 20
+ abort:
+ httpStatus: 503
+ percentage: 10