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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
@@ -0,0 +1,56 @@
parameters:
clusterName: ""
os: "linux"
# linux: cilium, cniv1, cniv2, dualstack; windows: cniv1, cniv2, stateless
cni: "cilium"
# all, add-before-endpoint-commit, delete-after-intent-commit, endpoint-patch, restart-during-scale
scenario: "all"
scaleReplicas: 20
# Keep the test and task timeouts larger so context cancellation can restore the CNS DaemonSet.
timeoutMinutes: 60
testTimeoutMinutes: 80
taskTimeoutMinutes: 95
runID: "$(Build.BuildId)-$(System.JobId)-$(System.JobAttempt)"
artifactName: "migration-fault-injection-$(System.JobId)"
workloadImage: "mcr.microsoft.com/oss/kubernetes/pause:3.6"

steps:
- task: AzureCLI@2
displayName: "Run CNS migration fault injection"
timeoutInMinutes: ${{ parameters.taskTimeoutMinutes }}
inputs:
azureSubscription: $(BUILD_VALIDATIONS_SERVICE_CONNECTION)
scriptLocation: "inlineScript"
scriptType: "bash"
addSpnToEnvironment: true
inlineScript: |
set -euo pipefail

artifactDir="$(Build.ArtifactStagingDirectory)/migration-fault-injection"
mkdir -p "$artifactDir"
export KUBECONFIG="$(Build.ArtifactStagingDirectory)/migration-fault-injection-kubeconfig"
trap 'rm -f "$KUBECONFIG"' EXIT
make -C ./hack/aks set-kubeconf AZCLI=az CLUSTER=${{ parameters.clusterName }}

MIGRATION_FAULT_SCENARIO="${{ parameters.scenario }}" \
MIGRATION_FAULT_OS="${{ parameters.os }}" \
MIGRATION_FAULT_CNI="${{ parameters.cni }}" \
MIGRATION_FAULT_RUN_ID="${{ parameters.runID }}" \
MIGRATION_FAULT_ARTIFACT_DIR="$artifactDir" \
MIGRATION_FAULT_SCALE_REPLICAS="${{ parameters.scaleReplicas }}" \
MIGRATION_FAULT_TIMEOUT_MINUTES="${{ parameters.timeoutMinutes }}" \
MIGRATION_FAULT_WORKLOAD_IMAGE="${{ parameters.workloadImage }}" \
VALIDATE_STATE_BACKEND=bolt \
VALIDATE_CONVERGENCE_ATTEMPTS=12 \
VALIDATE_CONVERGENCE_INTERVAL_SECONDS=10 \
go test -mod=readonly -count=1 -timeout ${{ parameters.testTimeoutMinutes }}m -tags load \
./test/integration/state -run '^TestMigrationFaultInjection$' -v \
-args -test-kubeconfig="$KUBECONFIG" 2>&1 |
tee "$artifactDir/go-test.log"

- task: PublishPipelineArtifact@1
displayName: "Publish migration fault artifacts"
condition: always()
inputs:
targetPath: "$(Build.ArtifactStagingDirectory)/migration-fault-injection"
artifact: "${{ parameters.artifactName }}"
213 changes: 213 additions & 0 deletions cns/restserver/fault_injection.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,213 @@
package restserver

import (
"crypto/subtle"
"encoding/json"
"errors"
"net/http"
"os"
"sync"
"time"
)

const (
faultInjectionPath = "/debug/faultinjection"
faultInjectionTokenHeader = "X-CNS-Test-Fault-Token" //nolint:gosec // This is a header name, not a credential.
faultInjectionTokenEnv = "CNS_TEST_FAULT_INJECTION_TOKEN" //nolint:gosec // This is an environment variable name, not a credential.
defaultFaultInjectionTimeout = 10 * time.Minute
)

type faultPoint string

const (
faultPointAddBeforeEndpointCommit faultPoint = "add-before-endpoint-commit"
faultPointDeleteAfterIntentCommit faultPoint = "delete-after-intent-commit"
faultPointPatchBeforeEndpointCommit faultPoint = "patch-before-endpoint-commit"
)

type faultState string

const (
faultStateIdle faultState = "idle"
faultStateArmed faultState = "armed"
faultStateReached faultState = "reached"
)

var errFaultAlreadyArmed = errors.New("fault injection point is already armed")

type faultInjector struct {
mu sync.Mutex
token string
timeout time.Duration
point faultPoint
target faultInjectionTarget
state faultState
release chan struct{}
}

type faultInjectionRequest struct {
Point faultPoint `json:"point"`
Target faultInjectionTarget `json:"target"`
}

type faultInjectionTarget struct {
PodName string `json:"podName,omitempty"`
PodNamespace string `json:"podNamespace,omitempty"`
}

type faultInjectionStatus struct {
Point faultPoint `json:"point,omitempty"`
Target faultInjectionTarget `json:"target"`
State faultState `json:"state"`
}

func newFaultInjectorFromEnv() *faultInjector {
token := os.Getenv(faultInjectionTokenEnv)
if token == "" {
return nil
}
return newFaultInjector(token, defaultFaultInjectionTimeout)
}

func newFaultInjector(token string, timeout time.Duration) *faultInjector {
return &faultInjector{
token: token,
timeout: timeout,
state: faultStateIdle,
}
}

func validFaultPoint(point faultPoint) bool {
switch point {
case faultPointAddBeforeEndpointCommit,
faultPointDeleteAfterIntentCommit,
faultPointPatchBeforeEndpointCommit:
return true
default:
return false
}
}

func (injector *faultInjector) arm(point faultPoint, target faultInjectionTarget) error {
injector.mu.Lock()
defer injector.mu.Unlock()

if injector.state != faultStateIdle {
return errFaultAlreadyArmed
}
injector.point = point
injector.target = target
injector.state = faultStateArmed
injector.release = make(chan struct{})
return nil
}

func (injector *faultInjector) disarm() {
injector.mu.Lock()
defer injector.mu.Unlock()

if injector.release != nil {
close(injector.release)
}
injector.point = ""
injector.target = faultInjectionTarget{}
injector.state = faultStateIdle
injector.release = nil
}

func (injector *faultInjector) checkpoint(point faultPoint, target faultInjectionTarget) {
injector.mu.Lock()
if injector.point != point || injector.state != faultStateArmed || !injector.target.matches(target) {
injector.mu.Unlock()
return
}
release := injector.release
injector.state = faultStateReached
injector.mu.Unlock()

timer := time.NewTimer(injector.timeout)
defer timer.Stop()
select {
case <-release:
case <-timer.C:
}

injector.mu.Lock()
if injector.release == release {
injector.point = ""
injector.target = faultInjectionTarget{}
injector.state = faultStateIdle
injector.release = nil
}
injector.mu.Unlock()
}

func (injector *faultInjector) status() faultInjectionStatus {
injector.mu.Lock()
defer injector.mu.Unlock()
return faultInjectionStatus{
Point: injector.point,
Target: injector.target,
State: injector.state,
}
}

func (target faultInjectionTarget) matches(candidate faultInjectionTarget) bool {
return target.PodNamespace == candidate.PodNamespace &&
(target.PodName == "" || target.PodName == candidate.PodName)
}

func (injector *faultInjector) handle(w http.ResponseWriter, r *http.Request) {
if subtle.ConstantTimeCompare([]byte(r.Header.Get(faultInjectionTokenHeader)), []byte(injector.token)) != 1 {
http.Error(w, "forbidden", http.StatusForbidden)
return
}

switch r.Method {
case http.MethodGet:
writeFaultInjectionStatus(w, injector.status())
case http.MethodPut:
var request faultInjectionRequest
decoder := json.NewDecoder(http.MaxBytesReader(w, r.Body, 1024))
decoder.DisallowUnknownFields()
if err := decoder.Decode(&request); err != nil {
http.Error(w, "invalid request", http.StatusBadRequest)
return
}
if !validFaultPoint(request.Point) {
http.Error(w, "invalid fault injection point", http.StatusBadRequest)
return
}
if request.Target.PodNamespace == "" {
http.Error(w, "fault injection target namespace is required", http.StatusBadRequest)
return
}
if err := injector.arm(request.Point, request.Target); err != nil {
http.Error(w, err.Error(), http.StatusConflict)
return
}
writeFaultInjectionStatus(w, injector.status())
case http.MethodDelete:
injector.disarm()
w.WriteHeader(http.StatusNoContent)
default:
w.Header().Set("Allow", "DELETE, GET, PUT")
http.Error(w, "method not allowed", http.StatusMethodNotAllowed)
}
}

func writeFaultInjectionStatus(w http.ResponseWriter, status faultInjectionStatus) {
w.Header().Set("Content-Type", "application/json")
if err := json.NewEncoder(w).Encode(status); err != nil {
http.Error(w, "encoding response", http.StatusInternalServerError)
}
}

func (service *HTTPRestService) reachFaultPoint(point faultPoint, podName, podNamespace string) {
if service.faultInjector != nil {
service.faultInjector.checkpoint(point, faultInjectionTarget{
PodName: podName,
PodNamespace: podNamespace,
})
}
}
Loading
Loading