summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
-rw-r--r--.github/workflows/main.yaml6
-rw-r--r--.travis.yml2
-rw-r--r--common/constants/default.go2
-rw-r--r--edge/pkg/edged/kubeclientbridge/typed/core/v1/serviceaccount_bridge.go6
-rw-r--r--edge/pkg/metamanager/client/serviceaccount.go2
-rw-r--r--go.mod98
-rw-r--r--go.sum99
-rw-r--r--vendor/github.com/google/cadvisor/container/containerd/client.go9
-rw-r--r--vendor/github.com/google/cadvisor/container/containerd/handler.go9
-rw-r--r--vendor/k8s.io/api/policy/v1beta1/generated.proto1
-rw-r--r--vendor/k8s.io/api/policy/v1beta1/types.go3
-rw-r--r--vendor/k8s.io/apimachinery/pkg/util/proxy/upgradeaware.go27
-rw-r--r--vendor/k8s.io/apiserver/pkg/audit/union.go1
-rw-r--r--vendor/k8s.io/apiserver/pkg/authentication/group/authenticated_group_adder.go11
-rw-r--r--vendor/k8s.io/apiserver/pkg/authentication/group/group_adder.go12
-rw-r--r--vendor/k8s.io/apiserver/pkg/authentication/group/token_group_adder.go12
-rw-r--r--vendor/k8s.io/apiserver/pkg/endpoints/handlers/responsewriters/writers.go6
-rw-r--r--vendor/k8s.io/apiserver/pkg/endpoints/metrics/metrics.go2
-rw-r--r--vendor/k8s.io/apiserver/pkg/server/filters/timeout.go6
-rw-r--r--vendor/k8s.io/apiserver/pkg/server/options/audit.go4
-rw-r--r--vendor/k8s.io/apiserver/pkg/storage/etcd3/store.go148
-rw-r--r--vendor/k8s.io/apiserver/pkg/storage/storagebackend/factory/etcd3.go48
-rw-r--r--vendor/k8s.io/client-go/plugin/pkg/client/auth/exec/exec.go37
-rw-r--r--vendor/k8s.io/client-go/rest/request.go38
-rw-r--r--vendor/k8s.io/client-go/tools/clientcmd/auth_loaders.go2
-rw-r--r--vendor/k8s.io/client-go/transport/cache.go25
-rw-r--r--vendor/k8s.io/client-go/transport/config.go21
-rw-r--r--vendor/k8s.io/client-go/transport/transport.go25
-rw-r--r--vendor/k8s.io/cri-api/pkg/errors/doc.go19
-rw-r--r--vendor/k8s.io/cri-api/pkg/errors/errors.go38
-rw-r--r--vendor/k8s.io/csi-translation-lib/plugins/azure_disk.go2
-rw-r--r--vendor/k8s.io/csi-translation-lib/plugins/azure_file.go35
-rw-r--r--vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/types.go10
-rw-r--r--vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/defaults.go3
-rw-r--r--vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/doc.go4
-rw-r--r--vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/types.go4
-rw-r--r--vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/zz_generated.conversion.go1
-rw-r--r--vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/defaults.go3
-rw-r--r--vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/doc.go4
-rw-r--r--vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/types.go4
-rw-r--r--vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/zz_generated.conversion.go1
-rw-r--r--vendor/k8s.io/kubernetes/cmd/kubeadm/app/constants/constants.go9
-rw-r--r--vendor/k8s.io/kubernetes/cmd/kubeadm/app/images/images.go2
-rw-r--r--vendor/k8s.io/kubernetes/cmd/kubeadm/app/phases/etcd/local.go35
-rw-r--r--vendor/k8s.io/kubernetes/cmd/kubeadm/app/util/pkiutil/pki_helpers.go18
-rw-r--r--vendor/k8s.io/kubernetes/cmd/kubelet/app/options/container_runtime.go2
-rw-r--r--vendor/k8s.io/kubernetes/cmd/kubelet/app/server.go2
-rw-r--r--vendor/k8s.io/kubernetes/cmd/kubelet/app/server_windows.go59
-rw-r--r--vendor/k8s.io/kubernetes/pkg/api/v1/pod/util.go23
-rw-r--r--vendor/k8s.io/kubernetes/pkg/apis/core/validation/events.go13
-rw-r--r--vendor/k8s.io/kubernetes/pkg/credentialprovider/azure/azure_acr_helper.go2
-rw-r--r--vendor/k8s.io/kubernetes/pkg/credentialprovider/azure/azure_credentials.go56
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/cm/containermap/container_map.go7
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/cm/cpumanager/cpu_manager.go10
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/cm/memorymanager/memory_manager.go7
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/dockershim/docker_sandbox.go2
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/kubelet.go99
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/kubelet_pods.go23
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/kuberuntime/kuberuntime_container.go25
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/kuberuntime/kuberuntime_manager.go15
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/kuberuntime/security_context_windows.go10
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/metrics/collectors/resource_metrics.go12
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/pod/pod_manager.go36
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/pod_workers.go180
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/prober/prober_manager.go25
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/runonce.go6
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/stats/cri_stats_provider.go19
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/status/generate.go45
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/status/status_manager.go134
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/token/token_manager.go9
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/util/manager/cache_based_manager.go6
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/volume_host.go14
-rw-r--r--vendor/k8s.io/kubernetes/pkg/kubelet/volumemanager/reconciler/reconciler.go38
-rw-r--r--vendor/k8s.io/kubernetes/pkg/printers/internalversion/printers.go10
-rw-r--r--vendor/k8s.io/kubernetes/pkg/proxy/endpointslicecache.go8
-rw-r--r--vendor/k8s.io/kubernetes/pkg/proxy/iptables/proxier.go108
-rw-r--r--vendor/k8s.io/kubernetes/pkg/proxy/util/utils.go2
-rw-r--r--vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/nodeaffinity/node_affinity.go2
-rw-r--r--vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/nodename/node_name.go2
-rw-r--r--vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/nodeports/node_ports.go2
-rw-r--r--vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/noderesources/fit.go2
-rw-r--r--vendor/k8s.io/kubernetes/pkg/volume/csi/csi_attacher.go38
-rw-r--r--vendor/k8s.io/kubernetes/pkg/volume/plugins.go2
-rw-r--r--vendor/k8s.io/kubernetes/pkg/volume/util/operationexecutor/operation_executor.go18
-rw-r--r--vendor/k8s.io/kubernetes/pkg/volume/util/operationexecutor/operation_generator.go29
-rw-r--r--vendor/k8s.io/mount-utils/go.mod2
-rw-r--r--vendor/k8s.io/mount-utils/go.sum4
-rw-r--r--vendor/k8s.io/utils/clock/clock.go17
-rw-r--r--vendor/k8s.io/utils/inotify/inotify_linux.go6
-rw-r--r--vendor/modules.txt107
90 files changed, 1340 insertions, 752 deletions
diff --git a/.github/workflows/main.yaml b/.github/workflows/main.yaml
index fe9f61c2d..7f31d64ca 100644
--- a/.github/workflows/main.yaml
+++ b/.github/workflows/main.yaml
@@ -161,7 +161,7 @@ jobs:
run: |
command -v ginkgo || go install github.com/onsi/ginkgo/v2/ginkgo@${{ env.GINKGO_VERSION }}
go install sigs.k8s.io/kind@v0.11.1
- curl -LO https://storage.googleapis.com/kubernetes-release/release/v1.22.6/bin/linux/amd64/kubectl && sudo install kubectl /usr/local/bin/kubectl
+ curl -LO https://storage.googleapis.com/kubernetes-release/release/v1.22.17/bin/linux/amd64/kubectl && sudo install kubectl /usr/local/bin/kubectl
- name: Checkout code
uses: actions/checkout@v2
@@ -219,7 +219,7 @@ jobs:
run: |
command -v ginkgo || go install github.com/onsi/ginkgo/v2/ginkgo@${{ env.GINKGO_VERSION }}
go install sigs.k8s.io/kind@v0.11.1
- curl -LO https://storage.googleapis.com/kubernetes-release/release/v1.22.6/bin/linux/amd64/kubectl && sudo install kubectl /usr/local/bin/kubectl
+ curl -LO https://storage.googleapis.com/kubernetes-release/release/v1.22.17/bin/linux/amd64/kubectl && sudo install kubectl /usr/local/bin/kubectl
- name: Checkout code
uses: actions/checkout@v2
@@ -266,7 +266,7 @@ jobs:
run: |
command -v ginkgo || go install github.com/onsi/ginkgo/v2/ginkgo@${{ env.GINKGO_VERSION }}
go install sigs.k8s.io/kind@v0.11.1
- curl -LO https://storage.googleapis.com/kubernetes-release/release/v1.22.6/bin/linux/amd64/kubectl && sudo install kubectl /usr/local/bin/kubectl
+ curl -LO https://storage.googleapis.com/kubernetes-release/release/v1.22.17/bin/linux/amd64/kubectl && sudo install kubectl /usr/local/bin/kubectl
- name: Checkout code
uses: actions/checkout@v2
diff --git a/.travis.yml b/.travis.yml
index 2571e6c67..cf1ff4744 100644
--- a/.travis.yml
+++ b/.travis.yml
@@ -59,6 +59,6 @@ jobs:
before_script:
- command -v ginkgo || go install github.com/onsi/ginkgo/ginkgo@latest
- go install sigs.k8s.io/kind@v0.11.1
- - curl -LO https://storage.googleapis.com/kubernetes-release/release/v1.22.6/bin/linux/arm64/kubectl && sudo install kubectl /usr/local/bin/kubectl
+ - curl -LO https://storage.googleapis.com/kubernetes-release/release/v1.22.17/bin/linux/arm64/kubectl && sudo install kubectl /usr/local/bin/kubectl
script:
- make e2e
diff --git a/common/constants/default.go b/common/constants/default.go
index 669023dec..70f9897f0 100644
--- a/common/constants/default.go
+++ b/common/constants/default.go
@@ -66,7 +66,7 @@ const (
DefaultVolumeStatsAggPeriod = time.Minute
DefaultTunnelPort = 10004
- CurrentSupportK8sVersion = "v1.22.6"
+ CurrentSupportK8sVersion = "v1.22.17"
// MetaManager
DefaultRemoteQueryTimeout = 60
diff --git a/edge/pkg/edged/kubeclientbridge/typed/core/v1/serviceaccount_bridge.go b/edge/pkg/edged/kubeclientbridge/typed/core/v1/serviceaccount_bridge.go
index 4f39c94f7..278f1c0a1 100644
--- a/edge/pkg/edged/kubeclientbridge/typed/core/v1/serviceaccount_bridge.go
+++ b/edge/pkg/edged/kubeclientbridge/typed/core/v1/serviceaccount_bridge.go
@@ -28,6 +28,7 @@ import (
authenticationv1 "k8s.io/api/authentication/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
+ "k8s.io/apimachinery/pkg/types"
fakecorev1 "k8s.io/client-go/kubernetes/typed/core/v1/fake"
"github.com/kubeedge/kubeedge/edge/pkg/metamanager/client"
@@ -44,3 +45,8 @@ type ServiceAccountsBridge struct {
func (c *ServiceAccountsBridge) CreateToken(ctx context.Context, serviceAccountName string, tokenRequest *authenticationv1.TokenRequest, opts metav1.CreateOptions) (result *authenticationv1.TokenRequest, err error) {
return c.MetaClient.ServiceAccountToken().GetServiceAccountToken(c.ns, serviceAccountName, tokenRequest)
}
+
+func (c *ServiceAccountsBridge) Delete(ctx context.Context, podUID string, opts metav1.DeleteOptions) error {
+ c.MetaClient.ServiceAccountToken().DeleteServiceAccountToken(types.UID(podUID))
+ return nil
+}
diff --git a/edge/pkg/metamanager/client/serviceaccount.go b/edge/pkg/metamanager/client/serviceaccount.go
index 0323fab40..121a64c43 100644
--- a/edge/pkg/metamanager/client/serviceaccount.go
+++ b/edge/pkg/metamanager/client/serviceaccount.go
@@ -48,7 +48,7 @@ func (c *serviceAccountToken) DeleteServiceAccountToken(podUID types.UID) {
for _, sa := range *svcAccounts {
var tr authenticationv1.TokenRequest
err = json.Unmarshal([]byte(sa.Value), &tr)
- if err != nil {
+ if err != nil || tr.Spec.BoundObjectRef == nil {
klog.Errorf("unmarshal resource %s token request failed: %v", sa.Key, err)
continue
}
diff --git a/go.mod b/go.mod
index 5b09e7d6c..15f2952f5 100644
--- a/go.mod
+++ b/go.mod
@@ -43,27 +43,27 @@ require (
google.golang.org/grpc v1.42.0
google.golang.org/protobuf v1.27.1
helm.sh/helm/v3 v3.7.2
- k8s.io/api v0.22.6
- k8s.io/apiextensions-apiserver v0.22.6
- k8s.io/apimachinery v0.22.6
- k8s.io/apiserver v0.22.6
- k8s.io/cli-runtime v0.22.6
- k8s.io/client-go v0.22.6
- k8s.io/cloud-provider v0.22.6 // indirect
- k8s.io/cluster-bootstrap v0.22.6 // indirect
- k8s.io/code-generator v0.22.6
- k8s.io/component-base v0.22.6
- k8s.io/cri-api v0.22.6
- k8s.io/csi-translation-lib v0.22.6 // indirect
+ k8s.io/api v0.22.17
+ k8s.io/apiextensions-apiserver v0.22.17
+ k8s.io/apimachinery v0.22.17
+ k8s.io/apiserver v0.22.17
+ k8s.io/cli-runtime v0.22.17
+ k8s.io/client-go v0.22.17
+ k8s.io/cloud-provider v0.22.17 // indirect
+ k8s.io/cluster-bootstrap v0.22.17 // indirect
+ k8s.io/code-generator v0.22.17
+ k8s.io/component-base v0.22.17
+ k8s.io/cri-api v0.22.17
+ k8s.io/csi-translation-lib v0.22.17 // indirect
k8s.io/klog/v2 v2.9.0
k8s.io/kube-openapi v0.0.0-20211115234752-e816edb12b65
- k8s.io/kube-scheduler v0.22.6 // indirect
- k8s.io/kubelet v0.22.6
- k8s.io/kubernetes v1.22.6
- k8s.io/mount-utils v0.22.6 // indirect
- k8s.io/utils v0.0.0-20210819203725-bdf08cb9a70a
+ k8s.io/kube-scheduler v0.22.17 // indirect
+ k8s.io/kubelet v0.22.17
+ k8s.io/kubernetes v1.22.17
+ k8s.io/mount-utils v0.22.17 // indirect
+ k8s.io/utils v0.0.0-20211116205334-6203023598ed
sigs.k8s.io/apiserver-network-proxy v0.0.27
- sigs.k8s.io/apiserver-network-proxy/konnectivity-client v0.0.27
+ sigs.k8s.io/apiserver-network-proxy/konnectivity-client v0.0.30
sigs.k8s.io/controller-runtime v0.10.3
sigs.k8s.io/yaml v1.2.0
)
@@ -84,38 +84,38 @@ replace (
go.etcd.io/etcd/pkg/v3 => go.etcd.io/etcd/pkg/v3 v3.5.0
go.etcd.io/etcd/raft/v3 => go.etcd.io/etcd/raft/v3 v3.5.0
go.etcd.io/etcd/server/v3 => go.etcd.io/etcd/server/v3 v3.5.0
- k8s.io/api => github.com/kubeedge/kubernetes/staging/src/k8s.io/api v1.22.6-kubeedge4
- k8s.io/apiextensions-apiserver => github.com/kubeedge/kubernetes/staging/src/k8s.io/apiextensions-apiserver v1.22.6-kubeedge4
- k8s.io/apimachinery => github.com/kubeedge/kubernetes/staging/src/k8s.io/apimachinery v1.22.6-kubeedge4
- k8s.io/apiserver => github.com/kubeedge/kubernetes/staging/src/k8s.io/apiserver v1.22.6-kubeedge4
- k8s.io/cli-runtime => github.com/kubeedge/kubernetes/staging/src/k8s.io/cli-runtime v1.22.6-kubeedge4
- k8s.io/client-go => github.com/kubeedge/kubernetes/staging/src/k8s.io/client-go v1.22.6-kubeedge4
- k8s.io/cloud-provider => github.com/kubeedge/kubernetes/staging/src/k8s.io/cloud-provider v1.22.6-kubeedge4
- k8s.io/cluster-bootstrap => github.com/kubeedge/kubernetes/staging/src/k8s.io/cluster-bootstrap v1.22.6-kubeedge4
- k8s.io/code-generator => github.com/kubeedge/kubernetes/staging/src/k8s.io/code-generator v1.22.6-kubeedge4
- k8s.io/component-base => github.com/kubeedge/kubernetes/staging/src/k8s.io/component-base v1.22.6-kubeedge4
- k8s.io/component-helpers => github.com/kubeedge/kubernetes/staging/src/k8s.io/component-helpers v1.22.6-kubeedge4
- k8s.io/controller-manager => github.com/kubeedge/kubernetes/staging/src/k8s.io/controller-manager v1.22.6-kubeedge4
- k8s.io/cri-api => github.com/kubeedge/kubernetes/staging/src/k8s.io/cri-api v1.22.6-kubeedge4
- k8s.io/csi-api => github.com/kubeedge/kubernetes/staging/src/k8s.io/csi-api v1.22.6-kubeedge4
- k8s.io/csi-translation-lib => github.com/kubeedge/kubernetes/staging/src/k8s.io/csi-translation-lib v1.22.6-kubeedge4
- k8s.io/gengo v0.0.0 => k8s.io/gengo v0.22.6
+ k8s.io/api => github.com/kubeedge/kubernetes/staging/src/k8s.io/api v1.22.17-kubeedge1
+ k8s.io/apiextensions-apiserver => github.com/kubeedge/kubernetes/staging/src/k8s.io/apiextensions-apiserver v1.22.17-kubeedge1
+ k8s.io/apimachinery => github.com/kubeedge/kubernetes/staging/src/k8s.io/apimachinery v1.22.17-kubeedge1
+ k8s.io/apiserver => github.com/kubeedge/kubernetes/staging/src/k8s.io/apiserver v1.22.17-kubeedge1
+ k8s.io/cli-runtime => github.com/kubeedge/kubernetes/staging/src/k8s.io/cli-runtime v1.22.17-kubeedge1
+ k8s.io/client-go => github.com/kubeedge/kubernetes/staging/src/k8s.io/client-go v1.22.17-kubeedge1
+ k8s.io/cloud-provider => github.com/kubeedge/kubernetes/staging/src/k8s.io/cloud-provider v1.22.17-kubeedge1
+ k8s.io/cluster-bootstrap => github.com/kubeedge/kubernetes/staging/src/k8s.io/cluster-bootstrap v1.22.17-kubeedge1
+ k8s.io/code-generator => github.com/kubeedge/kubernetes/staging/src/k8s.io/code-generator v1.22.17-kubeedge1
+ k8s.io/component-base => github.com/kubeedge/kubernetes/staging/src/k8s.io/component-base v1.22.17-kubeedge1
+ k8s.io/component-helpers => github.com/kubeedge/kubernetes/staging/src/k8s.io/component-helpers v1.22.17-kubeedge1
+ k8s.io/controller-manager => github.com/kubeedge/kubernetes/staging/src/k8s.io/controller-manager v1.22.17-kubeedge1
+ k8s.io/cri-api => github.com/kubeedge/kubernetes/staging/src/k8s.io/cri-api v1.22.17-kubeedge1
+ k8s.io/csi-api => github.com/kubeedge/kubernetes/staging/src/k8s.io/csi-api v1.22.17-kubeedge1
+ k8s.io/csi-translation-lib => github.com/kubeedge/kubernetes/staging/src/k8s.io/csi-translation-lib v1.22.17-kubeedge1
+ k8s.io/gengo v0.0.0 => k8s.io/gengo v0.22.17
k8s.io/heapster => k8s.io/heapster v1.2.0-beta.1 // indirect
k8s.io/klog/v2 => k8s.io/klog/v2 v2.8.0
- k8s.io/kube-aggregator => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-aggregator v1.22.6-kubeedge4
- k8s.io/kube-controller-manager => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-controller-manager v1.22.6-kubeedge4
- k8s.io/kube-proxy => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-proxy v1.22.6-kubeedge4
- k8s.io/kube-scheduler => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-scheduler v1.22.6-kubeedge4
- k8s.io/kubectl => github.com/kubeedge/kubernetes/staging/src/k8s.io/kubectl v1.22.6-kubeedge4
- k8s.io/kubelet => github.com/kubeedge/kubernetes/staging/src/k8s.io/kubelet v1.22.6-kubeedge4
- k8s.io/kubernetes => github.com/kubeedge/kubernetes v1.22.6-kubeedge4
- k8s.io/legacy-cloud-providers => github.com/kubeedge/kubernetes/staging/src/k8s.io/legacy-cloud-providers v1.22.6-kubeedge4
- k8s.io/metrics => github.com/kubeedge/kubernetes/staging/src/k8s.io/metrics v1.22.6-kubeedge4
- k8s.io/mount-utils => github.com/kubeedge/kubernetes/staging/src/k8s.io/mount-utils v1.22.6-kubeedge4
- k8s.io/node-api => github.com/kubeedge/kubernetes/staging/src/k8s.io/node-api v1.22.6-kubeedge4
- k8s.io/pod-security-admission => k8s.io/pod-security-admission v0.22.6
- k8s.io/repo-infra => github.com/kubeedge/kubernetes/staging/src/k8s.io/repo-infra v1.22.6-kubeedge4
- k8s.io/sample-apiserver => github.com/kubeedge/kubernetes/staging/src/k8s.io/sample-apiserver v1.22.6-kubeedge4
- k8s.io/utils v0.0.0 => k8s.io/utils v0.22.6
+ k8s.io/kube-aggregator => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-aggregator v1.22.17-kubeedge1
+ k8s.io/kube-controller-manager => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-controller-manager v1.22.17-kubeedge1
+ k8s.io/kube-proxy => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-proxy v1.22.17-kubeedge1
+ k8s.io/kube-scheduler => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-scheduler v1.22.17-kubeedge1
+ k8s.io/kubectl => github.com/kubeedge/kubernetes/staging/src/k8s.io/kubectl v1.22.17-kubeedge1
+ k8s.io/kubelet => github.com/kubeedge/kubernetes/staging/src/k8s.io/kubelet v1.22.17-kubeedge1
+ k8s.io/kubernetes => github.com/kubeedge/kubernetes v1.22.17-kubeedge1
+ k8s.io/legacy-cloud-providers => github.com/kubeedge/kubernetes/staging/src/k8s.io/legacy-cloud-providers v1.22.17-kubeedge1
+ k8s.io/metrics => github.com/kubeedge/kubernetes/staging/src/k8s.io/metrics v1.22.17-kubeedge1
+ k8s.io/mount-utils => github.com/kubeedge/kubernetes/staging/src/k8s.io/mount-utils v1.22.17-kubeedge1
+ k8s.io/node-api => github.com/kubeedge/kubernetes/staging/src/k8s.io/node-api v1.22.17-kubeedge1
+ k8s.io/pod-security-admission => k8s.io/pod-security-admission v0.22.17
+ k8s.io/repo-infra => github.com/kubeedge/kubernetes/staging/src/k8s.io/repo-infra v1.22.17-kubeedge1
+ k8s.io/sample-apiserver => github.com/kubeedge/kubernetes/staging/src/k8s.io/sample-apiserver v1.22.17-kubeedge1
+ k8s.io/utils v0.0.0 => k8s.io/utils v0.22.17
sigs.k8s.io/apiserver-network-proxy/konnectivity-client => sigs.k8s.io/apiserver-network-proxy/konnectivity-client v0.0.27
)
diff --git a/go.sum b/go.sum
index 2ceff2a42..f751946c1 100644
--- a/go.sum
+++ b/go.sum
@@ -556,8 +556,8 @@ github.com/google/btree v0.0.0-20180813153112-4030bb1f1f0c/go.mod h1:lNA+9X1NB3Z
github.com/google/btree v1.0.0/go.mod h1:lNA+9X1NB3Zf8V7Ke586lFgjr2dZNuvo3lPJSGZ5JPQ=
github.com/google/btree v1.0.1 h1:gK4Kx5IaGY9CD5sPJ36FHiBJ6ZXl0kilRiiCj+jdYp4=
github.com/google/btree v1.0.1/go.mod h1:xXMiIv4Fb/0kKde4SpL7qlzvu5cMJDRkFDxJfI9uaxA=
-github.com/google/cadvisor v0.39.3 h1:Z04WmZQWsdZKxAdtqE8h68zyfpCu9LimRyDg3GOmrbc=
-github.com/google/cadvisor v0.39.3/go.mod h1:kN93gpdevu+bpS227TyHVZyCU5bbqCzTj5T9drl34MI=
+github.com/google/cadvisor v0.39.4 h1:4Zep62r142bSYOmkY236DoUxSnWDAsGAytIiHxeSJpA=
+github.com/google/cadvisor v0.39.4/go.mod h1:kN93gpdevu+bpS227TyHVZyCU5bbqCzTj5T9drl34MI=
github.com/google/go-cmp v0.2.0/go.mod h1:oXzfMopK8JAjlY9xF4vHSVASa0yLyX7SntLO5aqRK0M=
github.com/google/go-cmp v0.3.0/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU=
github.com/google/go-cmp v0.3.1/go.mod h1:8QqcDgzrUqlUb/G2PQTWiueGozuR1884gddMywk6iLU=
@@ -722,51 +722,51 @@ github.com/kr/pty v1.1.5/go.mod h1:9r2w37qlBe7rQ6e1fg1S/9xpWHSnaqNdHD3WcMdbPDA=
github.com/kr/text v0.1.0/go.mod h1:4Jbv+DJW3UT/LiOwJeYQe1efqtUx/iVham/4vfdArNI=
github.com/kr/text v0.2.0 h1:5Nx0Ya0ZqY2ygV366QzturHI13Jq95ApcVaJBhpS+AY=
github.com/kr/text v0.2.0/go.mod h1:eLer722TekiGuMkidMxC/pM04lWEeraHUUmBw8l2grE=
-github.com/kubeedge/kubernetes v1.22.6-kubeedge4 h1:FJ8ptUbJlO0SJ+ifnRE2w4hFjnnQ7oqzAWFf2RFiF1c=
-github.com/kubeedge/kubernetes v1.22.6-kubeedge4/go.mod h1:l2ikQCpfvsMAXgL7FDtzgn/AVdjt4XGUYHMXn2vuzYI=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/api v1.22.6-kubeedge4 h1:sbjzRJJR8KNzrapAwiVuY6j3l3743U3BY3wcDlCF/io=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/api v1.22.6-kubeedge4/go.mod h1:IpPnJRE5t3olVaut5p67N16cZkWwwU5KVFM35xCKyxM=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/apiextensions-apiserver v1.22.6-kubeedge4 h1:JzjL1Wtuk2ZKcx41IDWTs2ctbMRd7XeM+1wqzLbC4sM=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/apiextensions-apiserver v1.22.6-kubeedge4/go.mod h1:p+50OD4/67BZHVF5arX0513Hv3Rdjcm3aAyyncaTOrM=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/apimachinery v1.22.6-kubeedge4 h1:fEkdJ5UHjvCNdjHmUbkwQCFoKrvIw7kpKiupUOkOv7o=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/apimachinery v1.22.6-kubeedge4/go.mod h1:/IrK/NrbS5lx7TWGM7uFr6a0nsKInBlU262XJahYCJA=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/apiserver v1.22.6-kubeedge4 h1:QweudAi3sl3pDLt+QuHHT68OYFipWslZG1kTT05Ujh4=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/apiserver v1.22.6-kubeedge4/go.mod h1:rjtkRXvtA12EAy1xGxy/mLQMCgG/+k3UEshKxF67l4k=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/cli-runtime v1.22.6-kubeedge4 h1:Fhcwbs2hia0wvNcUqn6jvllasELXTYTXB22RHTyNle8=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/cli-runtime v1.22.6-kubeedge4/go.mod h1:lfgLhMRSK39X7vRKyVDIHEUngwllg9H4TCGWbC6sCeA=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/client-go v1.22.6-kubeedge4 h1:akFrDvcbwEsIA/H7IfrlNN7PdS0A5JdQeZDnt5N3QWE=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/client-go v1.22.6-kubeedge4/go.mod h1:cm0LK9J9pqad2fptYbbmdQ8DFhziO7KRfymEWmmLM2w=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/cloud-provider v1.22.6-kubeedge4 h1:jkOZkZf19aapXn2/10T8f6sfbj/hEOO7Y/GQyNsFRKY=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/cloud-provider v1.22.6-kubeedge4/go.mod h1:YfjUcxHPiB9x/eHUrBtefZ61AuHGSDXfyXtsLS5UlMQ=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/cluster-bootstrap v1.22.6-kubeedge4 h1:oOASTjSlD+vL/85GXebc48NAod3OVVQIwnJYo9uUzxE=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/cluster-bootstrap v1.22.6-kubeedge4/go.mod h1:ppZJmhTukDTa5g/F0ksVMLM0Owbi9GeKhzuTXAVVJig=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/code-generator v1.22.6-kubeedge4 h1:zf3rCPzWhO7cuEWydmpsy/q5VpxZ7XyragFIsx0Y5H0=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/code-generator v1.22.6-kubeedge4/go.mod h1:NKIwVRZkqwXayEfmqQxbcDKfPU0diBpwnwTbhYI52WU=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/component-base v1.22.6-kubeedge4 h1:+4i1hyMQHrb/TcbOE1mCfLMYiX8xgPamZp5/ggSIDEI=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/component-base v1.22.6-kubeedge4/go.mod h1:cn9EB9A1wujtKWsHqB9lkYq8FL4dUuftmiqNyXIQEmE=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/component-helpers v1.22.6-kubeedge4 h1:QlszEqW2FGQemsAr2QmSJc7GVH8vEgNXvr4cJES6H6E=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/component-helpers v1.22.6-kubeedge4/go.mod h1:9Bx6HezI9sKzn5Boasw7vMT8FRgcXsExOoT87Wzdls4=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/controller-manager v1.22.6-kubeedge4/go.mod h1:aPin+82yKPEirDGBtNS/4fcc3a1QVOqdt6zzxOlrfc8=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/cri-api v1.22.6-kubeedge4 h1:5Slef+C/LXJsLTuUCAFP0lRw7FBjOCKApDlal4IQf8M=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/cri-api v1.22.6-kubeedge4/go.mod h1:GsmfuA0cekL0pkx1Ejxk3HZGHza02Ch2Lo09bNQLR8E=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/csi-translation-lib v1.22.6-kubeedge4 h1:Cyz0/Bo7sb1TTCXq07SXPWVW66xSpI64kSY3Ps+4lPE=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/csi-translation-lib v1.22.6-kubeedge4/go.mod h1:B1gPUSbK2PVSnkxCgw/fmDckzQU6UCuyl670XFbEw6Q=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-aggregator v1.22.6-kubeedge4/go.mod h1:SofO8m0VM1l/flYW/0v5+V60YXjujtX9UT1vmFbdqbQ=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-controller-manager v1.22.6-kubeedge4/go.mod h1:46iKO45TZat/zvPyqe8TjLLrTS/U/nGB92Ft63PEPF0=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-proxy v1.22.6-kubeedge4 h1:aDEuJKtPqAVFbJVpHN5ETeQCSW7J2jgVND0/4Vfr24o=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-proxy v1.22.6-kubeedge4/go.mod h1:6mEp02ABsuOeeBuUrrol78v9LYysX7Z8CZOMFlkPOOI=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-scheduler v1.22.6-kubeedge4 h1:iwumDiTtHyytZiT7Twko1OVg/ARHezFfIvOhM1qN4js=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-scheduler v1.22.6-kubeedge4/go.mod h1:xZnfOrGTta6rB9IWNKl82yzWKpMSUXVmyGHRilQ9kzM=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/kubectl v1.22.6-kubeedge4 h1:Z9N4Z4QKbAPvXUno5dNUPG3VlHVV6ak1B8Gu2vyWWv4=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/kubectl v1.22.6-kubeedge4/go.mod h1:tQ6sES/uhan72qymDDma4H0wrrF0gcZIYV8cczR20es=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/kubelet v1.22.6-kubeedge4 h1:M/5f4NSJDQMB3zMddGRv813f/NKQT8b4cTeHC7WZlgU=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/kubelet v1.22.6-kubeedge4/go.mod h1:WyrLJLrfabeu6M39g8+g4tuj+qjP70IVV8HQFp3EBb4=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/legacy-cloud-providers v1.22.6-kubeedge4 h1:9lxkbHxPq+tACBfshOrwBoXjI/x605a5o1b2mHYWdKE=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/legacy-cloud-providers v1.22.6-kubeedge4/go.mod h1:X8EaUY5K2IM/62KAMuHGuHyOhsJwvsoRwdvsyWjm++g=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/metrics v1.22.6-kubeedge4/go.mod h1:I5RbQZ+gj12KSgWzMyHaE0hudGajvT/Nc5jRE/WMJnI=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/mount-utils v1.22.6-kubeedge4 h1:YLyHq2Z0q0LRuUdVqtqVpv8dcSMvTOTD9GIsvOVSvYo=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/mount-utils v1.22.6-kubeedge4/go.mod h1:7UvmmOyjKl2RW0tgpT4l1z7dxVV4TMnAAlSN95cpUeM=
-github.com/kubeedge/kubernetes/staging/src/k8s.io/sample-apiserver v1.22.6-kubeedge4/go.mod h1:UxJ/6uQrndnWNUMK/16LYsFZK98WTolizwpjcIINGCg=
+github.com/kubeedge/kubernetes v1.22.17-kubeedge1 h1:/IMlJ/XmFCtbyHD7FrF9jL3aIKVTYFP1DM0LU02g8X4=
+github.com/kubeedge/kubernetes v1.22.17-kubeedge1/go.mod h1:4jiaBaIxIZcb9Y5IKxnavbaCCK623RcD4gEVFDXMLmY=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/api v1.22.17-kubeedge1 h1:wrazBXXmKNa4uFHQYDasYeaqaq2sTBXujfXqpRSL3/E=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/api v1.22.17-kubeedge1/go.mod h1:IpPnJRE5t3olVaut5p67N16cZkWwwU5KVFM35xCKyxM=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/apiextensions-apiserver v1.22.17-kubeedge1 h1:G+HuaZQ1cZENkZ3Ie57TkRfY7quJUzmpN9kUmN1oy7A=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/apiextensions-apiserver v1.22.17-kubeedge1/go.mod h1:mWBMLYRUDXJg1TwanLNwmWPiqldeUjJot9IdAPORiKk=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/apimachinery v1.22.17-kubeedge1 h1:Iz3ZdolM0l8tJ+Po/52WIsExU/b67vZtaFnQBHltKEs=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/apimachinery v1.22.17-kubeedge1/go.mod h1:/IrK/NrbS5lx7TWGM7uFr6a0nsKInBlU262XJahYCJA=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/apiserver v1.22.17-kubeedge1 h1:iI2+u8jKXfhhSd1zVj0NKNgZuzrCmyWhP1X/Ddao8RQ=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/apiserver v1.22.17-kubeedge1/go.mod h1:bGe2Wes3DLHrLH5ALNbxAZ07M1LBNpteVMTMGTTH6Ws=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/cli-runtime v1.22.17-kubeedge1 h1:PLv9TdVH6kNaZQi5ldq3MgrZXSe7lKdnm6Y0fVrtlj4=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/cli-runtime v1.22.17-kubeedge1/go.mod h1:lfgLhMRSK39X7vRKyVDIHEUngwllg9H4TCGWbC6sCeA=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/client-go v1.22.17-kubeedge1 h1:/itwI14WfZSLeob2JhSkzJ8W6KxriciWKw4oBPZCko4=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/client-go v1.22.17-kubeedge1/go.mod h1:nxkuZlMqecxZ2kmKjLghyn6Yyyq0g+tJ+qoCDqEFr/E=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/cloud-provider v1.22.17-kubeedge1 h1:oMdV6uCbhfD9EU3tHfFjWPmdgO1eEseT8QpAuQP31/M=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/cloud-provider v1.22.17-kubeedge1/go.mod h1:OmGl0zrZ/2w3B0A1izHeK+K6vUWA5SpRZtghMCUQ2yY=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/cluster-bootstrap v1.22.17-kubeedge1 h1:j98bs156vfDxToArMt8lz8fGVEB3ViqwaPEEk3DbApI=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/cluster-bootstrap v1.22.17-kubeedge1/go.mod h1:ppZJmhTukDTa5g/F0ksVMLM0Owbi9GeKhzuTXAVVJig=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/code-generator v1.22.17-kubeedge1 h1:GwmuYmCrUOLl8qCI9Cd2nk2A1t1b/cn2MU0gjdn/G8c=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/code-generator v1.22.17-kubeedge1/go.mod h1:NKIwVRZkqwXayEfmqQxbcDKfPU0diBpwnwTbhYI52WU=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/component-base v1.22.17-kubeedge1 h1:zfaPxaXXnJXaYwJADjOqK+rtWo/hpjZorDsTMiASDqM=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/component-base v1.22.17-kubeedge1/go.mod h1:EiX/DP8EO++EqDAikSg/z8YinZs98moaA/sEDqsXhyM=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/component-helpers v1.22.17-kubeedge1 h1:Pez8QXM6ZEFyye0V0pXopzqTDVqleSeOozSmm7dVBEY=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/component-helpers v1.22.17-kubeedge1/go.mod h1:VgE1xZZ1WqClZk4W7bZJDM1LKbG1voCUV8fXMZsBDtw=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/controller-manager v1.22.17-kubeedge1/go.mod h1:dXNuo21WVfljYJMkWvc+TjsSS8cmYHyulrk0/TAH1So=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/cri-api v1.22.17-kubeedge1 h1:p5l4fmx09fYD9QK5+PDEkFMMCtc7XzLYuM0nTXjZYdk=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/cri-api v1.22.17-kubeedge1/go.mod h1:GsmfuA0cekL0pkx1Ejxk3HZGHza02Ch2Lo09bNQLR8E=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/csi-translation-lib v1.22.17-kubeedge1 h1:WV/Oi2K1NFDIz8hdmY8Nl4bb+HAOQsDocsocBw8fAQE=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/csi-translation-lib v1.22.17-kubeedge1/go.mod h1:B1gPUSbK2PVSnkxCgw/fmDckzQU6UCuyl670XFbEw6Q=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-aggregator v1.22.17-kubeedge1/go.mod h1:L/PoLCUoh1YkpgmCzLtpgqzI55I2Tqqtpyji/Gf63ik=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-controller-manager v1.22.17-kubeedge1/go.mod h1:46iKO45TZat/zvPyqe8TjLLrTS/U/nGB92Ft63PEPF0=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-proxy v1.22.17-kubeedge1 h1:9KdlCLju0OOXMUbEBWK+7Ybvbp7MkItbHE0aY16Aya8=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-proxy v1.22.17-kubeedge1/go.mod h1:6mEp02ABsuOeeBuUrrol78v9LYysX7Z8CZOMFlkPOOI=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-scheduler v1.22.17-kubeedge1 h1:XQ/HMtR0YCc0rOt++L18MFJ28Cdd1vcM+jJXrE1irqw=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-scheduler v1.22.17-kubeedge1/go.mod h1:xZnfOrGTta6rB9IWNKl82yzWKpMSUXVmyGHRilQ9kzM=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/kubectl v1.22.17-kubeedge1 h1:eF4dRBXjCVpXFSSGWt+q8LFpq3tIalRA1JrALk2vt3w=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/kubectl v1.22.17-kubeedge1/go.mod h1:cnrmkfHvCJUdvBolwwZ+dBRcs2LH9Pz3CcX+5X4JwdQ=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/kubelet v1.22.17-kubeedge1 h1:L7z21aoSp2AEOVDcR+RjpSpxhIogBR/yPHbtAvz0kf4=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/kubelet v1.22.17-kubeedge1/go.mod h1:WyrLJLrfabeu6M39g8+g4tuj+qjP70IVV8HQFp3EBb4=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/legacy-cloud-providers v1.22.17-kubeedge1 h1:poiZYHIGkbVToIOJsqOu7PJAd0OtWhqQqxxmLheatiA=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/legacy-cloud-providers v1.22.17-kubeedge1/go.mod h1:Pb1qL1ouXy+fu+FJ7x0eDYG1EZaAKKz4K/TAD84vKLA=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/metrics v1.22.17-kubeedge1/go.mod h1:I5RbQZ+gj12KSgWzMyHaE0hudGajvT/Nc5jRE/WMJnI=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/mount-utils v1.22.17-kubeedge1 h1:Lo0LnXdfIlHNRFxjQTQs1/7MsxEHbDT3uAuaEgenP+w=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/mount-utils v1.22.17-kubeedge1/go.mod h1:qMu6VLq2u4wHl7+IQDnlt/BqGUshm9qv4dzMrx7GaFU=
+github.com/kubeedge/kubernetes/staging/src/k8s.io/sample-apiserver v1.22.17-kubeedge1/go.mod h1:UxJ/6uQrndnWNUMK/16LYsFZK98WTolizwpjcIINGCg=
github.com/kubernetes-csi/csi-lib-utils v0.6.1 h1:+AZ58SRSRWh2vmMoWAAGcv7x6fIyBMpyCXAgIc9kT28=
github.com/kubernetes-csi/csi-lib-utils v0.6.1/go.mod h1:GVmlUmxZ+SUjVLXicRFjqWUUvWez0g0Y78zNV9t7KfQ=
github.com/lann/builder v0.0.0-20180802200727-47ae307949d0 h1:SOEGU9fKiNWd/HOJuq6+3iTQz8KNCLtVX6idSoTLdUw=
@@ -1730,13 +1730,14 @@ k8s.io/kube-openapi v0.0.0-20210421082810-95288971da7e/go.mod h1:vHXdDvt9+2spS2R
k8s.io/kube-openapi v0.0.0-20211109043538-20434351676c/go.mod h1:vHXdDvt9+2spS2Rx9ql3I8tycm3H9FDfdUoIuKCefvw=
k8s.io/kube-openapi v0.0.0-20211115234752-e816edb12b65 h1:E3J9oCLlaobFUqsjG9DfKbP2BmgwBL2p7pn0A3dG9W4=
k8s.io/kube-openapi v0.0.0-20211115234752-e816edb12b65/go.mod h1:sX9MT8g7NVZM5lVL/j8QyCCJe8YSMW30QvGZWaCIDIk=
-k8s.io/pod-security-admission v0.22.6/go.mod h1:aBlsoKgqpuixDMaCu7U4d8IywNzcmGAKi7/fmfr840M=
+k8s.io/pod-security-admission v0.22.17/go.mod h1:vS8kM94GP4rGrgSDi5UOYvo4RAKIo66zQnbOfbdOK7I=
k8s.io/system-validators v1.5.0 h1:gGgluCTkpKc/zUszjamp4LFfMVM0wuYG2qjIFL4MMeQ=
k8s.io/system-validators v1.5.0/go.mod h1:bPldcLgkIUK22ALflnsXk8pvkTEndYdNuaHH6gRrl0Q=
k8s.io/utils v0.0.0-20201110183641-67b214c5f920/go.mod h1:jPW/WVKK9YHAvNhRxK0md/EJ228hCsBRufyofKtW8HA=
k8s.io/utils v0.0.0-20210802155522-efc7438f0176/go.mod h1:jPW/WVKK9YHAvNhRxK0md/EJ228hCsBRufyofKtW8HA=
-k8s.io/utils v0.0.0-20210819203725-bdf08cb9a70a h1:8dYfu/Fc9Gz2rNJKB9IQRGgQOh2clmRzNIPPY1xLY5g=
k8s.io/utils v0.0.0-20210819203725-bdf08cb9a70a/go.mod h1:jPW/WVKK9YHAvNhRxK0md/EJ228hCsBRufyofKtW8HA=
+k8s.io/utils v0.0.0-20211116205334-6203023598ed h1:ck1fRPWPJWsMd8ZRFsWc6mh/zHp5fZ/shhbrgPUxDAE=
+k8s.io/utils v0.0.0-20211116205334-6203023598ed/go.mod h1:jPW/WVKK9YHAvNhRxK0md/EJ228hCsBRufyofKtW8HA=
modernc.org/cc v1.0.0/go.mod h1:1Sk4//wdnYJiUIxnW8ddKpaOJCF37yAdqYnkxUpaYxw=
modernc.org/golex v1.0.0/go.mod h1:b/QX9oBD/LhixY6NDh+IdGv17hgB+51fET1i2kPSmvk=
modernc.org/mathutil v1.0.0/go.mod h1:wU0vUrJsVWBZ4P6e7xtFJEhFSNsfRLJ8H458uRjg03k=
diff --git a/vendor/github.com/google/cadvisor/container/containerd/client.go b/vendor/github.com/google/cadvisor/container/containerd/client.go
index 47eaffad7..d67f4aee0 100644
--- a/vendor/github.com/google/cadvisor/container/containerd/client.go
+++ b/vendor/github.com/google/cadvisor/container/containerd/client.go
@@ -16,6 +16,7 @@ package containerd
import (
"context"
+ "errors"
"fmt"
"net"
"sync"
@@ -24,6 +25,7 @@ import (
containersapi "github.com/containerd/containerd/api/services/containers/v1"
tasksapi "github.com/containerd/containerd/api/services/tasks/v1"
versionapi "github.com/containerd/containerd/api/services/version/v1"
+ tasktypes "github.com/containerd/containerd/api/types/task"
"github.com/containerd/containerd/containers"
"github.com/containerd/containerd/errdefs"
"github.com/containerd/containerd/pkg/dialer"
@@ -44,6 +46,10 @@ type ContainerdClient interface {
Version(ctx context.Context) (string, error)
}
+var (
+ ErrTaskIsInUnknownState = errors.New("containerd task is in unknown state") // used when process reported in containerd task is in Unknown State
+)
+
var once sync.Once
var ctrdClient ContainerdClient = nil
@@ -114,6 +120,9 @@ func (c *client) TaskPid(ctx context.Context, id string) (uint32, error) {
if err != nil {
return 0, errdefs.FromGRPC(err)
}
+ if response.Process.Status == tasktypes.StatusUnknown {
+ return 0, ErrTaskIsInUnknownState
+ }
return response.Process.Pid, nil
}
diff --git a/vendor/github.com/google/cadvisor/container/containerd/handler.go b/vendor/github.com/google/cadvisor/container/containerd/handler.go
index c21397902..0c8c80432 100644
--- a/vendor/github.com/google/cadvisor/container/containerd/handler.go
+++ b/vendor/github.com/google/cadvisor/container/containerd/handler.go
@@ -17,6 +17,7 @@ package containerd
import (
"encoding/json"
+ "errors"
"fmt"
"strings"
"time"
@@ -99,10 +100,14 @@ func newContainerdContainerHandler(
if err == nil {
break
}
- retry--
- if !errdefs.IsNotFound(err) || retry == 0 {
+
+ // Retry when task is not created yet or task is in unknown state (likely in process of initializing)
+ isRetriableError := errdefs.IsNotFound(err) || errors.Is(err, ErrTaskIsInUnknownState)
+ if !isRetriableError || retry == 0 {
return nil, err
}
+
+ retry--
time.Sleep(backoff)
backoff *= 2
}
diff --git a/vendor/k8s.io/api/policy/v1beta1/generated.proto b/vendor/k8s.io/api/policy/v1beta1/generated.proto
index 133ba0493..8a2824b51 100644
--- a/vendor/k8s.io/api/policy/v1beta1/generated.proto
+++ b/vendor/k8s.io/api/policy/v1beta1/generated.proto
@@ -144,7 +144,6 @@ message PodDisruptionBudgetSpec {
// A null selector selects no pods.
// An empty selector ({}) also selects no pods, which differs from standard behavior of selecting all pods.
// In policy/v1, an empty selector will select all pods in the namespace.
- // +patchStrategy=replace
// +optional
optional k8s.io.apimachinery.pkg.apis.meta.v1.LabelSelector selector = 2;
diff --git a/vendor/k8s.io/api/policy/v1beta1/types.go b/vendor/k8s.io/api/policy/v1beta1/types.go
index 553cb316d..486f93461 100644
--- a/vendor/k8s.io/api/policy/v1beta1/types.go
+++ b/vendor/k8s.io/api/policy/v1beta1/types.go
@@ -36,9 +36,8 @@ type PodDisruptionBudgetSpec struct {
// A null selector selects no pods.
// An empty selector ({}) also selects no pods, which differs from standard behavior of selecting all pods.
// In policy/v1, an empty selector will select all pods in the namespace.
- // +patchStrategy=replace
// +optional
- Selector *metav1.LabelSelector `json:"selector,omitempty" patchStrategy:"replace" protobuf:"bytes,2,opt,name=selector"`
+ Selector *metav1.LabelSelector `json:"selector,omitempty" protobuf:"bytes,2,opt,name=selector"`
// An eviction is allowed if at most "maxUnavailable" pods selected by
// "selector" are unavailable after the eviction, i.e. even in absence of
diff --git a/vendor/k8s.io/apimachinery/pkg/util/proxy/upgradeaware.go b/vendor/k8s.io/apimachinery/pkg/util/proxy/upgradeaware.go
index 8ef16eeb6..2c3073672 100644
--- a/vendor/k8s.io/apimachinery/pkg/util/proxy/upgradeaware.go
+++ b/vendor/k8s.io/apimachinery/pkg/util/proxy/upgradeaware.go
@@ -86,6 +86,8 @@ type UpgradeAwareHandler struct {
MaxBytesPerSec int64
// Responder is passed errors that occur while setting up proxying.
Responder ErrorResponder
+ // Reject to forward redirect response
+ RejectForwardingRedirects bool
}
const defaultFlushInterval = 200 * time.Millisecond
@@ -243,6 +245,31 @@ func (h *UpgradeAwareHandler) ServeHTTP(w http.ResponseWriter, req *http.Request
proxy.Transport = h.Transport
proxy.FlushInterval = h.FlushInterval
proxy.ErrorLog = log.New(noSuppressPanicError{}, "", log.LstdFlags)
+ if h.RejectForwardingRedirects {
+ oldModifyResponse := proxy.ModifyResponse
+ proxy.ModifyResponse = func(response *http.Response) error {
+ code := response.StatusCode
+ if code >= 300 && code <= 399 && len(response.Header.Get("Location")) > 0 {
+ // close the original response
+ response.Body.Close()
+ msg := "the backend attempted to redirect this request, which is not permitted"
+ // replace the response
+ *response = http.Response{
+ StatusCode: http.StatusBadGateway,
+ Status: fmt.Sprintf("%d %s", response.StatusCode, http.StatusText(response.StatusCode)),
+ Body: io.NopCloser(strings.NewReader(msg)),
+ ContentLength: int64(len(msg)),
+ }
+ } else {
+ if oldModifyResponse != nil {
+ if err := oldModifyResponse(response); err != nil {
+ return err
+ }
+ }
+ }
+ return nil
+ }
+ }
if h.Responder != nil {
// if an optional error interceptor/responder was provided wire it
// the custom responder might be used for providing a unified error reporting
diff --git a/vendor/k8s.io/apiserver/pkg/audit/union.go b/vendor/k8s.io/apiserver/pkg/audit/union.go
index 39dd74f74..0766a9207 100644
--- a/vendor/k8s.io/apiserver/pkg/audit/union.go
+++ b/vendor/k8s.io/apiserver/pkg/audit/union.go
@@ -48,6 +48,7 @@ func (u union) ProcessEvents(events ...*auditinternal.Event) bool {
func (u union) Run(stopCh <-chan struct{}) error {
var funcs []func() error
for _, backend := range u.backends {
+ backend := backend
funcs = append(funcs, func() error {
return backend.Run(stopCh)
})
diff --git a/vendor/k8s.io/apiserver/pkg/authentication/group/authenticated_group_adder.go b/vendor/k8s.io/apiserver/pkg/authentication/group/authenticated_group_adder.go
index 5ac6b2ddf..2f2cc6ec5 100644
--- a/vendor/k8s.io/apiserver/pkg/authentication/group/authenticated_group_adder.go
+++ b/vendor/k8s.io/apiserver/pkg/authentication/group/authenticated_group_adder.go
@@ -51,11 +51,16 @@ func (g *AuthenticatedGroupAdder) AuthenticateRequest(req *http.Request) (*authe
}
}
- r.User = &user.DefaultInfo{
+ newGroups := make([]string, 0, len(r.User.GetGroups())+1)
+ newGroups = append(newGroups, r.User.GetGroups()...)
+ newGroups = append(newGroups, user.AllAuthenticated)
+
+ ret := *r // shallow copy
+ ret.User = &user.DefaultInfo{
Name: r.User.GetName(),
UID: r.User.GetUID(),
- Groups: append(r.User.GetGroups(), user.AllAuthenticated),
+ Groups: newGroups,
Extra: r.User.GetExtra(),
}
- return r, true, nil
+ return &ret, true, nil
}
diff --git a/vendor/k8s.io/apiserver/pkg/authentication/group/group_adder.go b/vendor/k8s.io/apiserver/pkg/authentication/group/group_adder.go
index 3079dad62..1234c595a 100644
--- a/vendor/k8s.io/apiserver/pkg/authentication/group/group_adder.go
+++ b/vendor/k8s.io/apiserver/pkg/authentication/group/group_adder.go
@@ -41,11 +41,17 @@ func (g *GroupAdder) AuthenticateRequest(req *http.Request) (*authenticator.Resp
if err != nil || !ok {
return nil, ok, err
}
- r.User = &user.DefaultInfo{
+
+ newGroups := make([]string, 0, len(r.User.GetGroups())+len(g.Groups))
+ newGroups = append(newGroups, r.User.GetGroups()...)
+ newGroups = append(newGroups, g.Groups...)
+
+ ret := *r // shallow copy
+ ret.User = &user.DefaultInfo{
Name: r.User.GetName(),
UID: r.User.GetUID(),
- Groups: append(r.User.GetGroups(), g.Groups...),
+ Groups: newGroups,
Extra: r.User.GetExtra(),
}
- return r, true, nil
+ return &ret, true, nil
}
diff --git a/vendor/k8s.io/apiserver/pkg/authentication/group/token_group_adder.go b/vendor/k8s.io/apiserver/pkg/authentication/group/token_group_adder.go
index 0ed5ee559..51e27a67d 100644
--- a/vendor/k8s.io/apiserver/pkg/authentication/group/token_group_adder.go
+++ b/vendor/k8s.io/apiserver/pkg/authentication/group/token_group_adder.go
@@ -41,11 +41,17 @@ func (g *TokenGroupAdder) AuthenticateToken(ctx context.Context, token string) (
if err != nil || !ok {
return nil, ok, err
}
- r.User = &user.DefaultInfo{
+
+ newGroups := make([]string, 0, len(r.User.GetGroups())+len(g.Groups))
+ newGroups = append(newGroups, r.User.GetGroups()...)
+ newGroups = append(newGroups, g.Groups...)
+
+ ret := *r // shallow copy
+ ret.User = &user.DefaultInfo{
Name: r.User.GetName(),
UID: r.User.GetUID(),
- Groups: append(r.User.GetGroups(), g.Groups...),
+ Groups: newGroups,
Extra: r.User.GetExtra(),
}
- return r, true, nil
+ return &ret, true, nil
}
diff --git a/vendor/k8s.io/apiserver/pkg/endpoints/handlers/responsewriters/writers.go b/vendor/k8s.io/apiserver/pkg/endpoints/handlers/responsewriters/writers.go
index 16e16c353..2a50c1c58 100644
--- a/vendor/k8s.io/apiserver/pkg/endpoints/handlers/responsewriters/writers.go
+++ b/vendor/k8s.io/apiserver/pkg/endpoints/handlers/responsewriters/writers.go
@@ -143,8 +143,10 @@ var gzipPool = &sync.Pool{
}
const (
- // defaultGzipContentEncodingLevel is set to 4 which uses less CPU than the default level
- defaultGzipContentEncodingLevel = 4
+ // defaultGzipContentEncodingLevel is set to 1 which uses least CPU compared to higher levels, yet offers
+ // similar compression ratios (off by at most 1.5x, but typically within 1.1x-1.3x). For further details see -
+ // https://github.com/kubernetes/kubernetes/issues/112296
+ defaultGzipContentEncodingLevel = 1
// defaultGzipThresholdBytes is compared to the size of the first write from the stream
// (usually the entire object), and if the size is smaller no gzipping will be performed
// if the client requests it.
diff --git a/vendor/k8s.io/apiserver/pkg/endpoints/metrics/metrics.go b/vendor/k8s.io/apiserver/pkg/endpoints/metrics/metrics.go
index 5efc0a834..406d219a2 100644
--- a/vendor/k8s.io/apiserver/pkg/endpoints/metrics/metrics.go
+++ b/vendor/k8s.io/apiserver/pkg/endpoints/metrics/metrics.go
@@ -514,7 +514,7 @@ func InstrumentHandlerFunc(verb, group, version, resource, subresource, scope, c
// CleanScope returns the scope of the request.
func CleanScope(requestInfo *request.RequestInfo) string {
- if requestInfo.Name != "" {
+ if requestInfo.Name != "" || requestInfo.Verb == "create" {
return "resource"
}
if requestInfo.Namespace != "" {
diff --git a/vendor/k8s.io/apiserver/pkg/server/filters/timeout.go b/vendor/k8s.io/apiserver/pkg/server/filters/timeout.go
index 3f07306a9..31ec59354 100644
--- a/vendor/k8s.io/apiserver/pkg/server/filters/timeout.go
+++ b/vendor/k8s.io/apiserver/pkg/server/filters/timeout.go
@@ -91,6 +91,10 @@ func (t *timeoutHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
// resultCh is used as both errCh and stopCh
resultCh := make(chan interface{})
tw := newTimeoutWriter(w)
+
+ // Make a copy of request and work on it in new goroutine
+ // to avoid race condition when accessing/modifying request (e.g. headers)
+ rCopy := r.Clone(r.Context())
go func() {
defer func() {
err := recover()
@@ -105,7 +109,7 @@ func (t *timeoutHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
}
resultCh <- err
}()
- t.handler.ServeHTTP(tw, r)
+ t.handler.ServeHTTP(tw, rCopy)
}()
select {
case err := <-resultCh:
diff --git a/vendor/k8s.io/apiserver/pkg/server/options/audit.go b/vendor/k8s.io/apiserver/pkg/server/options/audit.go
index 6e062c987..017b910b1 100644
--- a/vendor/k8s.io/apiserver/pkg/server/options/audit.go
+++ b/vendor/k8s.io/apiserver/pkg/server/options/audit.go
@@ -20,6 +20,7 @@ import (
"fmt"
"io"
"os"
+ "path/filepath"
"strings"
"time"
@@ -529,6 +530,9 @@ func (o *AuditLogOptions) getWriter() (io.Writer, error) {
}
func (o *AuditLogOptions) ensureLogFile() error {
+ if err := os.MkdirAll(filepath.Dir(o.Path), 0700); err != nil {
+ return err
+ }
mode := os.FileMode(0600)
f, err := os.OpenFile(o.Path, os.O_CREATE|os.O_APPEND|os.O_RDWR, mode)
if err != nil {
diff --git a/vendor/k8s.io/apiserver/pkg/storage/etcd3/store.go b/vendor/k8s.io/apiserver/pkg/storage/etcd3/store.go
index 7f558a588..ecb8b179c 100644
--- a/vendor/k8s.io/apiserver/pkg/storage/etcd3/store.go
+++ b/vendor/k8s.io/apiserver/pkg/storage/etcd3/store.go
@@ -89,18 +89,23 @@ func New(c *clientv3.Client, codec runtime.Codec, newFunc func() runtime.Object,
func newStore(c *clientv3.Client, codec runtime.Codec, newFunc func() runtime.Object, prefix string, transformer value.Transformer, pagingEnabled bool, leaseManagerConfig LeaseManagerConfig) *store {
versioner := APIObjectVersioner{}
+ // for compatibility with etcd2 impl.
+ // no-op for default prefix of '/registry'.
+ // keeps compatibility with etcd2 impl for custom prefixes that don't start with '/'
+ pathPrefix := path.Join("/", prefix)
+ if !strings.HasSuffix(pathPrefix, "/") {
+ // Ensure the pathPrefix ends in "/" here to simplify key concatenation later.
+ pathPrefix += "/"
+ }
result := &store{
client: c,
codec: codec,
versioner: versioner,
transformer: transformer,
pagingEnabled: pagingEnabled,
- // for compatibility with etcd2 impl.
- // no-op for default prefix of '/registry'.
- // keeps compatibility with etcd2 impl for custom prefixes that don't start with '/'
- pathPrefix: path.Join("/", prefix),
- watcher: newWatcher(c, codec, newFunc, versioner, transformer),
- leaseManager: newDefaultLeaseManager(c, leaseManagerConfig),
+ pathPrefix: pathPrefix,
+ watcher: newWatcher(c, codec, newFunc, versioner, transformer),
+ leaseManager: newDefaultLeaseManager(c, leaseManagerConfig),
}
return result
}
@@ -112,9 +117,12 @@ func (s *store) Versioner() storage.Versioner {
// Get implements storage.Interface.Get.
func (s *store) Get(ctx context.Context, key string, opts storage.GetOptions, out runtime.Object) error {
- key = path.Join(s.pathPrefix, key)
+ preparedKey, err := s.prepareKey(key)
+ if err != nil {
+ return err
+ }
startTime := time.Now()
- getResp, err := s.client.KV.Get(ctx, key)
+ getResp, err := s.client.KV.Get(ctx, preparedKey)
metrics.RecordEtcdRequestLatency("get", getTypeName(out), startTime)
if err != nil {
return err
@@ -127,11 +135,11 @@ func (s *store) Get(ctx context.Context, key string, opts storage.GetOptions, ou
if opts.IgnoreNotFound {
return runtime.SetZeroValue(out)
}
- return storage.NewKeyNotFoundError(key, 0)
+ return storage.NewKeyNotFoundError(preparedKey, 0)
}
kv := getResp.Kvs[0]
- data, _, err := s.transformer.TransformFromStorage(kv.Value, authenticatedDataString(key))
+ data, _, err := s.transformer.TransformFromStorage(kv.Value, authenticatedDataString(preparedKey))
if err != nil {
return storage.NewInternalError(err.Error())
}
@@ -141,6 +149,11 @@ func (s *store) Get(ctx context.Context, key string, opts storage.GetOptions, ou
// Create implements storage.Interface.Create.
func (s *store) Create(ctx context.Context, key string, obj, out runtime.Object, ttl uint64) error {
+ preparedKey, err := s.prepareKey(key)
+ if err != nil {
+ return err
+ }
+
if version, err := s.versioner.ObjectResourceVersion(obj); err == nil && version != 0 {
return errors.New("resourceVersion should not be set on objects to be created")
}
@@ -151,30 +164,29 @@ func (s *store) Create(ctx context.Context, key string, obj, out runtime.Object,
if err != nil {
return err
}
- key = path.Join(s.pathPrefix, key)
opts, err := s.ttlOpts(ctx, int64(ttl))
if err != nil {
return err
}
- newData, err := s.transformer.TransformToStorage(data, authenticatedDataString(key))
+ newData, err := s.transformer.TransformToStorage(data, authenticatedDataString(preparedKey))
if err != nil {
return storage.NewInternalError(err.Error())
}
startTime := time.Now()
txnResp, err := s.client.KV.Txn(ctx).If(
- notFound(key),
+ notFound(preparedKey),
).Then(
- clientv3.OpPut(key, string(newData), opts...),
+ clientv3.OpPut(preparedKey, string(newData), opts...),
).Commit()
metrics.RecordEtcdRequestLatency("create", getTypeName(obj), startTime)
if err != nil {
return err
}
if !txnResp.Succeeded {
- return storage.NewKeyExistsError(key, 0)
+ return storage.NewKeyExistsError(preparedKey, 0)
}
if out != nil {
@@ -188,12 +200,15 @@ func (s *store) Create(ctx context.Context, key string, obj, out runtime.Object,
func (s *store) Delete(
ctx context.Context, key string, out runtime.Object, preconditions *storage.Preconditions,
validateDeletion storage.ValidateObjectFunc, cachedExistingObject runtime.Object) error {
+ preparedKey, err := s.prepareKey(key)
+ if err != nil {
+ return err
+ }
v, err := conversion.EnforcePtr(out)
if err != nil {
return fmt.Errorf("unable to convert output object to pointer: %v", err)
}
- key = path.Join(s.pathPrefix, key)
- return s.conditionalDelete(ctx, key, out, v, preconditions, validateDeletion, cachedExistingObject)
+ return s.conditionalDelete(ctx, preparedKey, out, v, preconditions, validateDeletion, cachedExistingObject)
}
func (s *store) conditionalDelete(
@@ -306,6 +321,10 @@ func (s *store) conditionalDelete(
func (s *store) GuaranteedUpdate(
ctx context.Context, key string, out runtime.Object, ignoreNotFound bool,
preconditions *storage.Preconditions, tryUpdate storage.UpdateFunc, cachedExistingObject runtime.Object) error {
+ preparedKey, err := s.prepareKey(key)
+ if err != nil {
+ return err
+ }
trace := utiltrace.New("GuaranteedUpdate etcd3", utiltrace.Field{"type", getTypeName(out)})
defer trace.LogIfLong(500 * time.Millisecond)
@@ -313,16 +332,15 @@ func (s *store) GuaranteedUpdate(
if err != nil {
return fmt.Errorf("unable to convert output object to pointer: %v", err)
}
- key = path.Join(s.pathPrefix, key)
getCurrentState := func() (*objState, error) {
startTime := time.Now()
- getResp, err := s.client.KV.Get(ctx, key)
+ getResp, err := s.client.KV.Get(ctx, preparedKey)
metrics.RecordEtcdRequestLatency("get", getTypeName(out), startTime)
if err != nil {
return nil, err
}
- return s.getState(getResp, key, v, ignoreNotFound)
+ return s.getState(getResp, preparedKey, v, ignoreNotFound)
}
var origState *objState
@@ -338,9 +356,9 @@ func (s *store) GuaranteedUpdate(
}
trace.Step("initial value restored")
- transformContext := authenticatedDataString(key)
+ transformContext := authenticatedDataString(preparedKey)
for {
- if err := preconditions.Check(key, origState.obj); err != nil {
+ if err := preconditions.Check(preparedKey, origState.obj); err != nil {
// If our data is already up to date, return the error
if origStateIsCurrent {
return err
@@ -423,11 +441,11 @@ func (s *store) GuaranteedUpdate(
startTime := time.Now()
txnResp, err := s.client.KV.Txn(ctx).If(
- clientv3.Compare(clientv3.ModRevision(key), "=", origState.rev),
+ clientv3.Compare(clientv3.ModRevision(preparedKey), "=", origState.rev),
).Then(
- clientv3.OpPut(key, string(newData), opts...),
+ clientv3.OpPut(preparedKey, string(newData), opts...),
).Else(
- clientv3.OpGet(key),
+ clientv3.OpGet(preparedKey),
).Commit()
metrics.RecordEtcdRequestLatency("update", getTypeName(out), startTime)
if err != nil {
@@ -436,8 +454,8 @@ func (s *store) GuaranteedUpdate(
trace.Step("Transaction committed")
if !txnResp.Succeeded {
getResp := (*clientv3.GetResponse)(txnResp.Responses[0].GetResponseRange())
- klog.V(4).Infof("GuaranteedUpdate of %s failed because of a conflict, going to retry", key)
- origState, err = s.getState(getResp, key, v, ignoreNotFound)
+ klog.V(4).Infof("GuaranteedUpdate of %s failed because of a conflict, going to retry", preparedKey)
+ origState, err = s.getState(getResp, preparedKey, v, ignoreNotFound)
if err != nil {
return err
}
@@ -453,6 +471,11 @@ func (s *store) GuaranteedUpdate(
// GetToList implements storage.Interface.GetToList.
func (s *store) GetToList(ctx context.Context, key string, listOpts storage.ListOptions, listObj runtime.Object) error {
+ preparedKey, err := s.prepareKey(key)
+ if err != nil {
+ return err
+ }
+
resourceVersion := listOpts.ResourceVersion
match := listOpts.ResourceVersionMatch
pred := listOpts.Predicate
@@ -474,7 +497,6 @@ func (s *store) GetToList(ctx context.Context, key string, listOpts storage.List
newItemFunc := getNewItemFunc(listObj, v)
- key = path.Join(s.pathPrefix, key)
startTime := time.Now()
var opts []clientv3.OpOption
if len(resourceVersion) > 0 && match == metav1.ResourceVersionMatchExact {
@@ -485,7 +507,7 @@ func (s *store) GetToList(ctx context.Context, key string, listOpts storage.List
opts = append(opts, clientv3.WithRev(int64(rv)))
}
- getResp, err := s.client.KV.Get(ctx, key, opts...)
+ getResp, err := s.client.KV.Get(ctx, preparedKey, opts...)
metrics.RecordEtcdRequestLatency("get", getTypeName(listPtr), startTime)
if err != nil {
return err
@@ -495,7 +517,7 @@ func (s *store) GetToList(ctx context.Context, key string, listOpts storage.List
}
if len(getResp.Kvs) > 0 {
- data, _, err := s.transformer.TransformFromStorage(getResp.Kvs[0].Value, authenticatedDataString(key))
+ data, _, err := s.transformer.TransformFromStorage(getResp.Kvs[0].Value, authenticatedDataString(preparedKey))
if err != nil {
return storage.NewInternalError(err.Error())
}
@@ -525,18 +547,21 @@ func getNewItemFunc(listObj runtime.Object, v reflect.Value) func() runtime.Obje
}
func (s *store) Count(key string) (int64, error) {
- key = path.Join(s.pathPrefix, key)
+ preparedKey, err := s.prepareKey(key)
+ if err != nil {
+ return 0, err
+ }
// We need to make sure the key ended with "/" so that we only get children "directories".
// e.g. if we have key "/a", "/a/b", "/ab", getting keys with prefix "/a" will return all three,
// while with prefix "/a/" will return only "/a/b" which is the correct answer.
- if !strings.HasSuffix(key, "/") {
- key += "/"
+ if !strings.HasSuffix(preparedKey, "/") {
+ preparedKey += "/"
}
startTime := time.Now()
- getResp, err := s.client.KV.Get(context.Background(), key, clientv3.WithRange(clientv3.GetPrefixRangeEnd(key)), clientv3.WithCountOnly())
- metrics.RecordEtcdRequestLatency("listWithCount", key, startTime)
+ getResp, err := s.client.KV.Get(context.Background(), preparedKey, clientv3.WithRange(clientv3.GetPrefixRangeEnd(preparedKey)), clientv3.WithCountOnly())
+ metrics.RecordEtcdRequestLatency("listWithCount", preparedKey, startTime)
if err != nil {
return 0, err
}
@@ -605,6 +630,11 @@ func encodeContinue(key, keyPrefix string, resourceVersion int64) (string, error
// List implements storage.Interface.List.
func (s *store) List(ctx context.Context, key string, opts storage.ListOptions, listObj runtime.Object) error {
+ preparedKey, err := s.prepareKey(key)
+ if err != nil {
+ return err
+ }
+
resourceVersion := opts.ResourceVersion
match := opts.ResourceVersionMatch
pred := opts.Predicate
@@ -624,16 +654,13 @@ func (s *store) List(ctx context.Context, key string, opts storage.ListOptions,
return fmt.Errorf("need ptr to slice: %v", err)
}
- if s.pathPrefix != "" {
- key = path.Join(s.pathPrefix, key)
- }
// We need to make sure the key ended with "/" so that we only get children "directories".
// e.g. if we have key "/a", "/a/b", "/ab", getting keys with prefix "/a" will return all three,
// while with prefix "/a/" will return only "/a/b" which is the correct answer.
- if !strings.HasSuffix(key, "/") {
- key += "/"
+ if !strings.HasSuffix(preparedKey, "/") {
+ preparedKey += "/"
}
- keyPrefix := key
+ keyPrefix := preparedKey
// set the appropriate clientv3 options to filter the returned data set
var paging bool
@@ -669,7 +696,7 @@ func (s *store) List(ctx context.Context, key string, opts storage.ListOptions,
rangeEnd := clientv3.GetPrefixRangeEnd(keyPrefix)
options = append(options, clientv3.WithRange(rangeEnd))
- key = continueKey
+ preparedKey = continueKey
// If continueRV > 0, the LIST request needs a specific resource version.
// continueRV==0 is invalid.
@@ -726,7 +753,7 @@ func (s *store) List(ctx context.Context, key string, opts storage.ListOptions,
var getResp *clientv3.GetResponse
for {
startTime := time.Now()
- getResp, err = s.client.KV.Get(ctx, key, options...)
+ getResp, err = s.client.KV.Get(ctx, preparedKey, options...)
metrics.RecordEtcdRequestLatency("list", getTypeName(listPtr), startTime)
if err != nil {
return interpretListError(err, len(pred.Continue) > 0, continueKey, keyPrefix)
@@ -779,7 +806,7 @@ func (s *store) List(ctx context.Context, key string, opts storage.ListOptions,
if int64(v.Len()) >= pred.Limit {
break
}
- key = string(lastKey) + "\x00"
+ preparedKey = string(lastKey) + "\x00"
if withRev == 0 {
withRev = returnedRV
options = append(options, clientv3.WithRev(withRev))
@@ -833,7 +860,7 @@ func growSlice(v reflect.Value, maxCapacity int, sizes ...int) {
return
}
if v.Len() > 0 {
- extra := reflect.MakeSlice(v.Type(), 0, max)
+ extra := reflect.MakeSlice(v.Type(), v.Len(), max)
reflect.Copy(extra, v)
v.Set(extra)
} else {
@@ -853,12 +880,15 @@ func (s *store) WatchList(ctx context.Context, key string, opts storage.ListOpti
}
func (s *store) watch(ctx context.Context, key string, opts storage.ListOptions, recursive bool) (watch.Interface, error) {
+ preparedKey, err := s.prepareKey(key)
+ if err != nil {
+ return nil, err
+ }
rev, err := s.versioner.ParseResourceVersion(opts.ResourceVersion)
if err != nil {
return nil, err
}
- key = path.Join(s.pathPrefix, key)
- return s.watcher.Watch(ctx, key, int64(rev), recursive, opts.ProgressNotify, opts.Predicate)
+ return s.watcher.Watch(ctx, preparedKey, int64(rev), recursive, opts.ProgressNotify, opts.Predicate)
}
func (s *store) getState(getResp *clientv3.GetResponse, key string, v reflect.Value, ignoreNotFound bool) (*objState, error) {
@@ -970,6 +1000,30 @@ func (s *store) validateMinimumResourceVersion(minimumResourceVersion string, ac
return nil
}
+func (s *store) prepareKey(key string) (string, error) {
+ if key == ".." ||
+ strings.HasPrefix(key, "../") ||
+ strings.HasSuffix(key, "/..") ||
+ strings.Contains(key, "/../") {
+ return "", fmt.Errorf("invalid key: %q", key)
+ }
+ if key == "." ||
+ strings.HasPrefix(key, "./") ||
+ strings.HasSuffix(key, "/.") ||
+ strings.Contains(key, "/./") {
+ return "", fmt.Errorf("invalid key: %q", key)
+ }
+ if key == "" || key == "/" {
+ return "", fmt.Errorf("empty key: %q", key)
+ }
+ // We ensured that pathPrefix ends in '/' in construction, so skip any leading '/' in the key now.
+ startIndex := 0
+ if key[0] == '/' {
+ startIndex = 1
+ }
+ return s.pathPrefix + key[startIndex:], nil
+}
+
// decode decodes value of bytes into object. It will also set the object resource version to rev.
// On success, objPtr would be set to the object.
func decode(codec runtime.Codec, versioner storage.Versioner, value []byte, objPtr runtime.Object, rev int64) error {
diff --git a/vendor/k8s.io/apiserver/pkg/storage/storagebackend/factory/etcd3.go b/vendor/k8s.io/apiserver/pkg/storage/storagebackend/factory/etcd3.go
index 74ebea655..eef9bd586 100644
--- a/vendor/k8s.io/apiserver/pkg/storage/storagebackend/factory/etcd3.go
+++ b/vendor/k8s.io/apiserver/pkg/storage/storagebackend/factory/etcd3.go
@@ -19,8 +19,10 @@ package factory
import (
"context"
"fmt"
+ "log"
"net"
"net/url"
+ "os"
"path"
"strings"
"sync"
@@ -28,9 +30,12 @@ import (
"time"
grpcprom "github.com/grpc-ecosystem/go-grpc-prometheus"
+ "go.etcd.io/etcd/client/pkg/v3/logutil"
"go.etcd.io/etcd/client/pkg/v3/transport"
clientv3 "go.etcd.io/etcd/client/v3"
"go.opentelemetry.io/contrib/instrumentation/google.golang.org/grpc/otelgrpc"
+ "go.uber.org/zap"
+ "go.uber.org/zap/zapcore"
"google.golang.org/grpc"
"k8s.io/apimachinery/pkg/runtime"
@@ -63,6 +68,14 @@ const (
dbMetricsMonitorJitter = 0.5
)
+// TODO(negz): Stop using a package scoped logger. At the time of writing we're
+// creating an etcd client for each CRD. We need to pass each etcd client a
+// logger or each client will create its own, which comes with a significant
+// memory cost (around 20% of the API server's memory when hundreds of CRDs are
+// present). The correct fix here is to not create a client per CRD. See
+// https://github.com/kubernetes/kubernetes/issues/111476 for more.
+var etcd3ClientLogger *zap.Logger
+
func init() {
// grpcprom auto-registers (via an init function) their client metrics, since we are opting out of
// using the global prometheus registry and using our own wrapped global registry,
@@ -70,6 +83,30 @@ func init() {
// For reference: https://github.com/kubernetes/kubernetes/pull/81387
legacyregistry.RawMustRegister(grpcprom.DefaultClientMetrics)
dbMetricsMonitors = make(map[string]struct{})
+
+ c := logutil.DefaultZapLoggerConfig
+ c.Level = zap.NewAtomicLevelAt(etcdClientDebugLevel())
+ l, err := c.Build()
+ if err != nil {
+ l = zap.NewNop()
+ }
+ etcd3ClientLogger = l.Named("etcd-client")
+}
+
+// etcdClientDebugLevel translates ETCD_CLIENT_DEBUG into zap log level.
+// NOTE(negz): This is a copy of a private etcd client function:
+// https://github.com/etcd-io/etcd/blob/v3.5.4/client/v3/logger.go#L47
+func etcdClientDebugLevel() zapcore.Level {
+ envLevel := os.Getenv("ETCD_CLIENT_DEBUG")
+ if envLevel == "" || envLevel == "true" {
+ return zapcore.InfoLevel
+ }
+ var l zapcore.Level
+ if err := l.Set(envLevel); err == nil {
+ log.Printf("Deprecated env ETCD_CLIENT_DEBUG value. Using default level: 'info'")
+ return zapcore.InfoLevel
+ }
+ return l
}
func newETCD3HealthCheck(c storagebackend.Config) (func() error, error) {
@@ -137,8 +174,13 @@ func newETCD3Client(c storagebackend.TransportConfig) (*clientv3.Client, error)
}
dialOptions := []grpc.DialOption{
grpc.WithBlock(), // block until the underlying connection is up
- grpc.WithUnaryInterceptor(grpcprom.UnaryClientInterceptor),
- grpc.WithStreamInterceptor(grpcprom.StreamClientInterceptor),
+ // use chained interceptors so that the default (retry and backoff) interceptors are added.
+ // otherwise they will be overwritten by the metric interceptor.
+ //
+ // these optional interceptors will be placed after the default ones.
+ // which seems to be what we want as the metrics will be collected on each attempt (retry)
+ grpc.WithChainUnaryInterceptor(grpcprom.UnaryClientInterceptor),
+ grpc.WithChainStreamInterceptor(grpcprom.StreamClientInterceptor),
}
if utilfeature.DefaultFeatureGate.Enabled(genericfeatures.APIServerTracing) {
tracingOpts := []otelgrpc.Option{
@@ -167,6 +209,7 @@ func newETCD3Client(c storagebackend.TransportConfig) (*clientv3.Client, error)
}
dialOptions = append(dialOptions, grpc.WithContextDialer(dialer))
}
+
cfg := clientv3.Config{
DialTimeout: dialTimeout,
DialKeepAliveTime: keepaliveTime,
@@ -174,6 +217,7 @@ func newETCD3Client(c storagebackend.TransportConfig) (*clientv3.Client, error)
DialOptions: dialOptions,
Endpoints: c.ServerList,
TLS: tlsConfig,
+ Logger: etcd3ClientLogger,
}
return clientv3.New(cfg)
diff --git a/vendor/k8s.io/client-go/plugin/pkg/client/auth/exec/exec.go b/vendor/k8s.io/client-go/plugin/pkg/client/auth/exec/exec.go
index 7a0984626..710aad894 100644
--- a/vendor/k8s.io/client-go/plugin/pkg/client/auth/exec/exec.go
+++ b/vendor/k8s.io/client-go/plugin/pkg/client/auth/exec/exec.go
@@ -39,6 +39,7 @@ import (
"k8s.io/apimachinery/pkg/runtime/schema"
"k8s.io/apimachinery/pkg/runtime/serializer"
"k8s.io/apimachinery/pkg/util/clock"
+ utilnet "k8s.io/apimachinery/pkg/util/net"
"k8s.io/client-go/pkg/apis/clientauthentication"
"k8s.io/client-go/pkg/apis/clientauthentication/install"
clientauthenticationv1 "k8s.io/client-go/pkg/apis/clientauthentication/v1"
@@ -200,14 +201,18 @@ func newAuthenticator(c *cache, isTerminalFunc func(int) bool, config *api.ExecC
now: time.Now,
environ: os.Environ,
- defaultDialer: defaultDialer,
- connTracker: connTracker,
+ connTracker: connTracker,
}
for _, env := range config.Env {
a.env = append(a.env, env.Name+"="+env.Value)
}
+ // these functions are made comparable and stored in the cache so that repeated clientset
+ // construction with the same rest.Config results in a single TLS cache and Authenticator
+ a.getCert = &transport.GetCertHolder{GetCert: a.cert}
+ a.dial = &transport.DialHolder{Dial: defaultDialer.DialContext}
+
return c.put(key, a), nil
}
@@ -262,8 +267,6 @@ type Authenticator struct {
now func() time.Time
environ func() []string
- // defaultDialer is used for clients which don't specify a custom dialer
- defaultDialer *connrotation.Dialer
// connTracker tracks all connections opened that we need to close when rotating a client certificate
connTracker *connrotation.ConnectionTracker
@@ -274,6 +277,12 @@ type Authenticator struct {
mu sync.Mutex
cachedCreds *credentials
exp time.Time
+
+ // getCert makes Authenticator.cert comparable to support TLS config caching
+ getCert *transport.GetCertHolder
+ // dial is used for clients which do not specify a custom dialer
+ // it is comparable to support TLS config caching
+ dial *transport.DialHolder
}
type credentials struct {
@@ -301,26 +310,34 @@ func (a *Authenticator) UpdateTransportConfig(c *transport.Config) error {
if c.TLS.GetCert != nil {
return errors.New("can't add TLS certificate callback: transport.Config.TLS.GetCert already set")
}
- c.TLS.GetCert = a.cert
+ c.TLS.GetCert = a.getCert.GetCert
+ c.TLS.GetCertHolder = a.getCert // comparable for TLS config caching
- var d *connrotation.Dialer
if c.Dial != nil {
// if c has a custom dialer, we have to wrap it
- d = connrotation.NewDialerWithTracker(c.Dial, a.connTracker)
+ // TLS config caching is not supported for this config
+ d := connrotation.NewDialerWithTracker(c.Dial, a.connTracker)
+ c.Dial = d.DialContext
+ c.DialHolder = nil
} else {
- d = a.defaultDialer
+ c.Dial = a.dial.Dial
+ c.DialHolder = a.dial // comparable for TLS config caching
}
- c.Dial = d.DialContext
-
return nil
}
+var _ utilnet.RoundTripperWrapper = &roundTripper{}
+
type roundTripper struct {
a *Authenticator
base http.RoundTripper
}
+func (r *roundTripper) WrappedRoundTripper() http.RoundTripper {
+ return r.base
+}
+
func (r *roundTripper) RoundTrip(req *http.Request) (*http.Response, error) {
// If a user has already set credentials, use that. This makes commands like
// "kubectl get --token (token) pods" work.
diff --git a/vendor/k8s.io/client-go/rest/request.go b/vendor/k8s.io/client-go/rest/request.go
index e5a8100be..3deb0e5cd 100644
--- a/vendor/k8s.io/client-go/rest/request.go
+++ b/vendor/k8s.io/client-go/rest/request.go
@@ -82,6 +82,12 @@ func (r *RequestConstructionError) Error() string {
var noBackoff = &NoBackoff{}
+type requestRetryFunc func(maxRetries int) WithRetry
+
+func defaultRequestRetryFn(maxRetries int) WithRetry {
+ return &withRetry{maxRetries: maxRetries}
+}
+
// Request allows for building up a request to a server in a chained fashion.
// Any errors are stored until the end of your call, so you only have to
// check once.
@@ -93,6 +99,7 @@ type Request struct {
rateLimiter flowcontrol.RateLimiter
backoff BackoffManager
timeout time.Duration
+ maxRetries int
// generic components accessible via method setters
verb string
@@ -109,9 +116,10 @@ type Request struct {
subresource string
// output
- err error
- body io.Reader
- retry WithRetry
+ err error
+ body io.Reader
+
+ retryFn requestRetryFunc
}
// NewRequest creates a new request helper object for accessing runtime.Objects on a server.
@@ -142,7 +150,8 @@ func NewRequest(c *RESTClient) *Request {
backoff: backoff,
timeout: timeout,
pathPrefix: pathPrefix,
- retry: &withRetry{maxRetries: 10},
+ maxRetries: 10,
+ retryFn: defaultRequestRetryFn,
warningHandler: c.warningHandler,
}
@@ -408,7 +417,10 @@ func (r *Request) Timeout(d time.Duration) *Request {
// function is specifically called with a different value.
// A zero maxRetries prevent it from doing retires and return an error immediately.
func (r *Request) MaxRetries(maxRetries int) *Request {
- r.retry.SetMaxRetries(maxRetries)
+ if maxRetries < 0 {
+ maxRetries = 0
+ }
+ r.maxRetries = maxRetries
return r
}
@@ -688,8 +700,10 @@ func (r *Request) Watch(ctx context.Context) (watch.Interface, error) {
}
return false
}
+
var retryAfter *RetryAfter
url := r.URL().String()
+ withRetry := r.retryFn(r.maxRetries)
for {
req, err := r.newHTTPRequest(ctx)
if err != nil {
@@ -724,9 +738,9 @@ func (r *Request) Watch(ctx context.Context) (watch.Interface, error) {
defer readAndCloseResponseBody(resp)
var retry bool
- retryAfter, retry = r.retry.NextRetry(req, resp, err, isErrRetryableFunc)
+ retryAfter, retry = withRetry.NextRetry(req, resp, err, isErrRetryableFunc)
if retry {
- err := r.retry.BeforeNextRetry(ctx, r.backoff, retryAfter, url, r.body)
+ err := withRetry.BeforeNextRetry(ctx, r.backoff, retryAfter, url, r.body)
if err == nil {
return false, nil
}
@@ -817,6 +831,7 @@ func (r *Request) Stream(ctx context.Context) (io.ReadCloser, error) {
}
var retryAfter *RetryAfter
+ withRetry := r.retryFn(r.maxRetries)
url := r.URL().String()
for {
req, err := r.newHTTPRequest(ctx)
@@ -862,9 +877,9 @@ func (r *Request) Stream(ctx context.Context) (io.ReadCloser, error) {
defer resp.Body.Close()
var retry bool
- retryAfter, retry = r.retry.NextRetry(req, resp, err, neverRetryError)
+ retryAfter, retry = withRetry.NextRetry(req, resp, err, neverRetryError)
if retry {
- err := r.retry.BeforeNextRetry(ctx, r.backoff, retryAfter, url, r.body)
+ err := withRetry.BeforeNextRetry(ctx, r.backoff, retryAfter, url, r.body)
if err == nil {
return false, nil
}
@@ -961,6 +976,7 @@ func (r *Request) request(ctx context.Context, fn func(*http.Request, *http.Resp
// Right now we make about ten retry attempts if we get a Retry-After response.
var retryAfter *RetryAfter
+ withRetry := r.retryFn(r.maxRetries)
for {
req, err := r.newHTTPRequest(ctx)
if err != nil {
@@ -997,7 +1013,7 @@ func (r *Request) request(ctx context.Context, fn func(*http.Request, *http.Resp
}
var retry bool
- retryAfter, retry = r.retry.NextRetry(req, resp, err, func(req *http.Request, err error) bool {
+ retryAfter, retry = withRetry.NextRetry(req, resp, err, func(req *http.Request, err error) bool {
// "Connection reset by peer" or "apiserver is shutting down" are usually a transient errors.
// Thus in case of "GET" operations, we simply retry it.
// We are not automatically retrying "write" operations, as they are not idempotent.
@@ -1011,7 +1027,7 @@ func (r *Request) request(ctx context.Context, fn func(*http.Request, *http.Resp
return false
})
if retry {
- err := r.retry.BeforeNextRetry(ctx, r.backoff, retryAfter, req.URL.String(), r.body)
+ err := withRetry.BeforeNextRetry(ctx, r.backoff, retryAfter, req.URL.String(), r.body)
if err == nil {
return false
}
diff --git a/vendor/k8s.io/client-go/tools/clientcmd/auth_loaders.go b/vendor/k8s.io/client-go/tools/clientcmd/auth_loaders.go
index 0e4127762..5153a95a2 100644
--- a/vendor/k8s.io/client-go/tools/clientcmd/auth_loaders.go
+++ b/vendor/k8s.io/client-go/tools/clientcmd/auth_loaders.go
@@ -51,10 +51,10 @@ func (a *PromptingAuthLoader) LoadAuth(path string) (*clientauth.Info, error) {
// Prompt for user/pass and write a file if none exists.
if _, err := os.Stat(path); os.IsNotExist(err) {
authPtr, err := a.Prompt()
- auth := *authPtr
if err != nil {
return nil, err
}
+ auth := *authPtr
data, err := json.Marshal(auth)
if err != nil {
return &auth, err
diff --git a/vendor/k8s.io/client-go/transport/cache.go b/vendor/k8s.io/client-go/transport/cache.go
index 5fe768ed5..f4a864d05 100644
--- a/vendor/k8s.io/client-go/transport/cache.go
+++ b/vendor/k8s.io/client-go/transport/cache.go
@@ -17,6 +17,7 @@ limitations under the License.
package transport
import (
+ "context"
"fmt"
"net"
"net/http"
@@ -50,6 +51,9 @@ type tlsCacheKey struct {
serverName string
nextProtos string
disableCompression bool
+ // these functions are wrapped to allow them to be used as map keys
+ getCert *GetCertHolder
+ dial *DialHolder
}
func (t tlsCacheKey) String() string {
@@ -57,7 +61,8 @@ func (t tlsCacheKey) String() string {
if len(t.keyData) > 0 {
keyText = "<redacted>"
}
- return fmt.Sprintf("insecure:%v, caData:%#v, certData:%#v, keyData:%s, serverName:%s, disableCompression:%t", t.insecure, t.caData, t.certData, keyText, t.serverName, t.disableCompression)
+ return fmt.Sprintf("insecure:%v, caData:%#v, certData:%#v, keyData:%s, serverName:%s, disableCompression:%t, getCert:%p, dial:%p",
+ t.insecure, t.caData, t.certData, keyText, t.serverName, t.disableCompression, t.getCert, t.dial)
}
func (c *tlsTransportCache) get(config *Config) (http.RoundTripper, error) {
@@ -87,8 +92,10 @@ func (c *tlsTransportCache) get(config *Config) (http.RoundTripper, error) {
return http.DefaultTransport, nil
}
- dial := config.Dial
- if dial == nil {
+ var dial func(ctx context.Context, network, address string) (net.Conn, error)
+ if config.Dial != nil {
+ dial = config.Dial
+ } else {
dial = (&net.Dialer{
Timeout: 30 * time.Second,
KeepAlive: 30 * time.Second,
@@ -133,10 +140,18 @@ func tlsConfigKey(c *Config) (tlsCacheKey, bool, error) {
return tlsCacheKey{}, false, err
}
- if c.TLS.GetCert != nil || c.Dial != nil || c.Proxy != nil {
+ if c.Proxy != nil {
// cannot determine equality for functions
return tlsCacheKey{}, false, nil
}
+ if c.Dial != nil && c.DialHolder == nil {
+ // cannot determine equality for dial function that doesn't have non-nil DialHolder set as well
+ return tlsCacheKey{}, false, nil
+ }
+ if c.TLS.GetCert != nil && c.TLS.GetCertHolder == nil {
+ // cannot determine equality for getCert function that doesn't have non-nil GetCertHolder set as well
+ return tlsCacheKey{}, false, nil
+ }
k := tlsCacheKey{
insecure: c.TLS.Insecure,
@@ -144,6 +159,8 @@ func tlsConfigKey(c *Config) (tlsCacheKey, bool, error) {
serverName: c.TLS.ServerName,
nextProtos: strings.Join(c.TLS.NextProtos, ","),
disableCompression: c.DisableCompression,
+ getCert: c.TLS.GetCertHolder,
+ dial: c.DialHolder,
}
if c.TLS.ReloadTLSFiles {
diff --git a/vendor/k8s.io/client-go/transport/config.go b/vendor/k8s.io/client-go/transport/config.go
index 070474831..ab18911b1 100644
--- a/vendor/k8s.io/client-go/transport/config.go
+++ b/vendor/k8s.io/client-go/transport/config.go
@@ -68,7 +68,11 @@ type Config struct {
WrapTransport WrapperFunc
// Dial specifies the dial function for creating unencrypted TCP connections.
+ // If specified, this transport will be non-cacheable unless DialHolder is also set.
Dial func(ctx context.Context, network, address string) (net.Conn, error)
+ // DialHolder can be populated to make transport configs cacheable.
+ // If specified, DialHolder.Dial must be equal to Dial.
+ DialHolder *DialHolder
// Proxy is the proxy func to be used for all requests made by this
// transport. If Proxy is nil, http.ProxyFromEnvironment is used. If Proxy
@@ -78,6 +82,11 @@ type Config struct {
Proxy func(*http.Request) (*url.URL, error)
}
+// DialHolder is used to make the wrapped function comparable so that it can be used as a map key.
+type DialHolder struct {
+ Dial func(ctx context.Context, network, address string) (net.Conn, error)
+}
+
// ImpersonationConfig has all the available impersonation options
type ImpersonationConfig struct {
// UserName matches user.Info.GetName()
@@ -141,5 +150,15 @@ type TLSConfig struct {
// To use only http/1.1, set to ["http/1.1"].
NextProtos []string
- GetCert func() (*tls.Certificate, error) // Callback that returns a TLS client certificate. CertData, CertFile, KeyData and KeyFile supercede this field.
+ // Callback that returns a TLS client certificate. CertData, CertFile, KeyData and KeyFile supercede this field.
+ // If specified, this transport is non-cacheable unless CertHolder is populated.
+ GetCert func() (*tls.Certificate, error)
+ // CertHolder can be populated to make transport configs that set GetCert cacheable.
+ // If set, CertHolder.GetCert must be equal to GetCert.
+ GetCertHolder *GetCertHolder
+}
+
+// GetCertHolder is used to make the wrapped function comparable so that it can be used as a map key.
+type GetCertHolder struct {
+ GetCert func() (*tls.Certificate, error)
}
diff --git a/vendor/k8s.io/client-go/transport/transport.go b/vendor/k8s.io/client-go/transport/transport.go
index b4a7bfa67..eabfce72d 100644
--- a/vendor/k8s.io/client-go/transport/transport.go
+++ b/vendor/k8s.io/client-go/transport/transport.go
@@ -24,6 +24,7 @@ import (
"fmt"
"io/ioutil"
"net/http"
+ "reflect"
"sync"
"time"
@@ -39,6 +40,10 @@ func New(config *Config) (http.RoundTripper, error) {
return nil, fmt.Errorf("using a custom transport with TLS certificate options or the insecure flag is not allowed")
}
+ if !isValidHolders(config) {
+ return nil, fmt.Errorf("misconfigured holder for dialer or cert callback")
+ }
+
var (
rt http.RoundTripper
err error
@@ -56,6 +61,26 @@ func New(config *Config) (http.RoundTripper, error) {
return HTTPWrappersForConfig(config, rt)
}
+func isValidHolders(config *Config) bool {
+ if config.TLS.GetCertHolder != nil {
+ if config.TLS.GetCertHolder.GetCert == nil ||
+ config.TLS.GetCert == nil ||
+ reflect.ValueOf(config.TLS.GetCertHolder.GetCert).Pointer() != reflect.ValueOf(config.TLS.GetCert).Pointer() {
+ return false
+ }
+ }
+
+ if config.DialHolder != nil {
+ if config.DialHolder.Dial == nil ||
+ config.Dial == nil ||
+ reflect.ValueOf(config.DialHolder.Dial).Pointer() != reflect.ValueOf(config.Dial).Pointer() {
+ return false
+ }
+ }
+
+ return true
+}
+
// TLSConfigFor returns a tls.Config that will provide the transport level security defined
// by the provided Config. Will return nil if no transport level security is requested.
func TLSConfigFor(c *Config) (*tls.Config, error) {
diff --git a/vendor/k8s.io/cri-api/pkg/errors/doc.go b/vendor/k8s.io/cri-api/pkg/errors/doc.go
new file mode 100644
index 000000000..f3413ee98
--- /dev/null
+++ b/vendor/k8s.io/cri-api/pkg/errors/doc.go
@@ -0,0 +1,19 @@
+/*
+Copyright 2020 The Kubernetes Authors.
+
+Licensed 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.
+*/
+
+// Package errors provides helper functions for use by the kubelet
+// to deal with CRI errors.
+package errors // import "k8s.io/cri-api/pkg/errors"
diff --git a/vendor/k8s.io/cri-api/pkg/errors/errors.go b/vendor/k8s.io/cri-api/pkg/errors/errors.go
new file mode 100644
index 000000000..41d7b9246
--- /dev/null
+++ b/vendor/k8s.io/cri-api/pkg/errors/errors.go
@@ -0,0 +1,38 @@
+/*
+Copyright 2020 The Kubernetes Authors.
+
+Licensed 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.
+*/
+
+package errors
+
+import (
+ "google.golang.org/grpc/codes"
+ "google.golang.org/grpc/status"
+)
+
+// IsNotFound returns a boolean indicating whether the error
+// is grpc not found error.
+// See https://github.com/grpc/grpc/blob/master/doc/statuscodes.md
+// for a list of grpc status codes.
+func IsNotFound(err error) bool {
+ s, ok := status.FromError(err)
+ if !ok {
+ return ok
+ }
+ if s.Code() == codes.NotFound {
+ return true
+ }
+
+ return false
+}
diff --git a/vendor/k8s.io/csi-translation-lib/plugins/azure_disk.go b/vendor/k8s.io/csi-translation-lib/plugins/azure_disk.go
index f7d9da43e..6d327df4d 100644
--- a/vendor/k8s.io/csi-translation-lib/plugins/azure_disk.go
+++ b/vendor/k8s.io/csi-translation-lib/plugins/azure_disk.go
@@ -107,7 +107,7 @@ func (t *azureDiskCSITranslator) TranslateInTreeInlineVolumeToCSI(volume *v1.Vol
ObjectMeta: metav1.ObjectMeta{
// Must be unique per disk as it is used as the unique part of the
// staging path
- Name: fmt.Sprintf("%s-%s", AzureDiskDriverName, azureSource.DiskName),
+ Name: azureSource.DataDiskURI,
},
Spec: v1.PersistentVolumeSpec{
PersistentVolumeSource: v1.PersistentVolumeSource{
diff --git a/vendor/k8s.io/csi-translation-lib/plugins/azure_file.go b/vendor/k8s.io/csi-translation-lib/plugins/azure_file.go
index df8251325..5df63c716 100644
--- a/vendor/k8s.io/csi-translation-lib/plugins/azure_file.go
+++ b/vendor/k8s.io/csi-translation-lib/plugins/azure_file.go
@@ -34,7 +34,7 @@ const (
AzureFileInTreePluginName = "kubernetes.io/azure-file"
separator = "#"
- volumeIDTemplate = "%s#%s#%s#%s"
+ volumeIDTemplate = "%s#%s#%s#%s#%s"
// Parameter names defined in azure file CSI driver, refer to
// https://github.com/kubernetes-sigs/azurefile-csi-driver/blob/master/docs/driver-parameters.md
shareNameField = "sharename"
@@ -81,19 +81,19 @@ func (t *azureFileCSITranslator) TranslateInTreeInlineVolumeToCSI(volume *v1.Vol
if podNamespace != "" {
secretNamespace = podNamespace
}
+ volumeID := fmt.Sprintf(volumeIDTemplate, "", accountName, azureSource.ShareName, volume.Name, secretNamespace)
var (
pv = &v1.PersistentVolume{
ObjectMeta: metav1.ObjectMeta{
- // Must be unique per disk as it is used as the unique part of the
- // staging path
- Name: fmt.Sprintf("%s-%s", AzureFileDriverName, azureSource.ShareName),
+ // Must be unique as it is used as the unique part of the staging path
+ Name: volumeID,
},
Spec: v1.PersistentVolumeSpec{
PersistentVolumeSource: v1.PersistentVolumeSource{
CSI: &v1.CSIPersistentVolumeSource{
Driver: AzureFileDriverName,
- VolumeHandle: fmt.Sprintf(volumeIDTemplate, "", accountName, azureSource.ShareName, ""),
+ VolumeHandle: volumeID,
ReadOnly: azureSource.ReadOnly,
VolumeAttributes: map[string]string{shareNameField: azureSource.ShareName},
NodeStageSecretRef: &v1.SecretReference{
@@ -129,7 +129,24 @@ func (t *azureFileCSITranslator) TranslateInTreePVToCSI(pv *v1.PersistentVolume)
resourceGroup = v
}
}
- volumeID := fmt.Sprintf(volumeIDTemplate, resourceGroup, accountName, azureSource.ShareName, "")
+
+ // Secret is required when mounting a volume but pod presence cannot be assumed - we should not try to read pod now.
+ namespace := ""
+ // Try to read SecretNamespace from source pv.
+ if azureSource.SecretNamespace != nil {
+ namespace = *azureSource.SecretNamespace
+ } else {
+ // Try to read namespace from ClaimRef which should be always present.
+ if pv.Spec.ClaimRef != nil {
+ namespace = pv.Spec.ClaimRef.Namespace
+ }
+ }
+
+ if len(namespace) == 0 {
+ return nil, fmt.Errorf("could not find a secret namespace in PersistentVolumeSource or ClaimRef")
+ }
+
+ volumeID := fmt.Sprintf(volumeIDTemplate, resourceGroup, accountName, azureSource.ShareName, pv.ObjectMeta.Name, namespace)
var (
// refer to https://github.com/kubernetes-sigs/azurefile-csi-driver/blob/master/docs/driver-parameters.md
@@ -137,7 +154,7 @@ func (t *azureFileCSITranslator) TranslateInTreePVToCSI(pv *v1.PersistentVolume)
Driver: AzureFileDriverName,
NodeStageSecretRef: &v1.SecretReference{
Name: azureSource.SecretName,
- Namespace: defaultSecretNamespace,
+ Namespace: namespace,
},
ReadOnly: azureSource.ReadOnly,
VolumeAttributes: map[string]string{shareNameField: azureSource.ShareName},
@@ -145,10 +162,6 @@ func (t *azureFileCSITranslator) TranslateInTreePVToCSI(pv *v1.PersistentVolume)
}
)
- if azureSource.SecretNamespace != nil {
- csiSource.NodeStageSecretRef.Namespace = *azureSource.SecretNamespace
- }
-
pv.Spec.PersistentVolumeSource.AzureFile = nil
pv.Spec.PersistentVolumeSource.CSI = csiSource
diff --git a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/types.go b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/types.go
index ef21f2c0c..516bf0232 100644
--- a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/types.go
+++ b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/types.go
@@ -84,9 +84,15 @@ type ClusterConfiguration struct {
// Networking holds configuration for the networking topology of the cluster.
Networking Networking
+
// KubernetesVersion is the target version of the control plane.
KubernetesVersion string
+ // CIKubernetesVersion is the target CI version of the control plane.
+ // Useful for running kubeadm with CI Kubernetes version.
+ // +k8s:conversion-gen=false
+ CIKubernetesVersion string
+
// ControlPlaneEndpoint sets a stable IP address or DNS name for the control plane; it
// can be a valid IP address or a RFC-1123 DNS subdomain, both with optional TCP port.
// In case the ControlPlaneEndpoint is not specified, the AdvertiseAddress + BindPort
@@ -116,8 +122,8 @@ type ClusterConfiguration struct {
CertificatesDir string
// ImageRepository sets the container registry to pull images from.
- // If empty, `k8s.gcr.io` will be used by default; in case of kubernetes version is a CI build (kubernetes version starts with `ci/` or `ci-cross/`)
- // `gcr.io/k8s-staging-ci-images` will be used as a default for control plane components and for kube-proxy, while `k8s.gcr.io`
+ // If empty, `registry.k8s.io` will be used by default; in case of kubernetes version is a CI build (kubernetes version starts with `ci/`)
+ // `gcr.io/k8s-staging-ci-images` will be used as a default for control plane components and for kube-proxy, while `registry.k8s.io`
// will be used for all the other images.
ImageRepository string
diff --git a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/defaults.go b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/defaults.go
index 2611da7be..5760f8fea 100644
--- a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/defaults.go
+++ b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/defaults.go
@@ -40,7 +40,8 @@ const (
// DefaultCertificatesDir defines default certificate directory
DefaultCertificatesDir = "/etc/kubernetes/pki"
// DefaultImageRepository defines default image registry
- DefaultImageRepository = "k8s.gcr.io"
+ // (previously this defaulted to k8s.gcr.io)
+ DefaultImageRepository = "registry.k8s.io"
// DefaultManifestsDir defines default manifests directory
DefaultManifestsDir = "/etc/kubernetes/manifests"
// DefaultClusterName defines the default cluster name
diff --git a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/doc.go b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/doc.go
index 4d86d91db..e2e4509c0 100644
--- a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/doc.go
+++ b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/doc.go
@@ -186,7 +186,7 @@ limitations under the License.
// etcd:
// # one of local or external
// local:
-// imageRepository: "k8s.gcr.io"
+// imageRepository: "registry.k8s.io"
// imageTag: "3.2.24"
// dataDir: "/var/lib/etcd"
// extraArgs:
@@ -240,7 +240,7 @@ limitations under the License.
// readOnly: false
// pathType: File
// certificatesDir: "/etc/kubernetes/pki"
-// imageRepository: "k8s.gcr.io"
+// imageRepository: "registry.k8s.io"
// useHyperKubeImage: false
// clusterName: "example-cluster"
// ---
diff --git a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/types.go b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/types.go
index 15c535f49..ee6d5f073 100644
--- a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/types.go
+++ b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/types.go
@@ -96,8 +96,8 @@ type ClusterConfiguration struct {
CertificatesDir string `json:"certificatesDir,omitempty"`
// ImageRepository sets the container registry to pull images from.
- // If empty, `k8s.gcr.io` will be used by default; in case of kubernetes version is a CI build (kubernetes version starts with `ci/` or `ci-cross/`)
- // `gcr.io/k8s-staging-ci-images` will be used as a default for control plane components and for kube-proxy, while `k8s.gcr.io`
+ // If empty, `registry.k8s.io` will be used by default; in case of kubernetes version is a CI build (kubernetes version starts with `ci/`)
+ // `gcr.io/k8s-staging-ci-images` will be used as a default for control plane components and for kube-proxy, while `registry.k8s.io`
// will be used for all the other images.
ImageRepository string `json:"imageRepository,omitempty"`
diff --git a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/zz_generated.conversion.go b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/zz_generated.conversion.go
index edb19d114..262515d19 100644
--- a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/zz_generated.conversion.go
+++ b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta2/zz_generated.conversion.go
@@ -335,6 +335,7 @@ func autoConvert_kubeadm_ClusterConfiguration_To_v1beta2_ClusterConfiguration(in
return err
}
out.KubernetesVersion = in.KubernetesVersion
+ // INFO: in.CIKubernetesVersion opted out of conversion generation
out.ControlPlaneEndpoint = in.ControlPlaneEndpoint
if err := Convert_kubeadm_APIServer_To_v1beta2_APIServer(&in.APIServer, &out.APIServer, s); err != nil {
return err
diff --git a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/defaults.go b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/defaults.go
index feb72d780..5d9da2ecf 100644
--- a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/defaults.go
+++ b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/defaults.go
@@ -42,7 +42,8 @@ const (
// DefaultCertificatesDir defines default certificate directory
DefaultCertificatesDir = "/etc/kubernetes/pki"
// DefaultImageRepository defines default image registry
- DefaultImageRepository = "k8s.gcr.io"
+ // (previously this defaulted to k8s.gcr.io)
+ DefaultImageRepository = "registry.k8s.io"
// DefaultManifestsDir defines default manifests directory
DefaultManifestsDir = "/etc/kubernetes/manifests"
// DefaultClusterName defines the default cluster name
diff --git a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/doc.go b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/doc.go
index 7234f2e4c..1484cd0c2 100644
--- a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/doc.go
+++ b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/doc.go
@@ -195,7 +195,7 @@ limitations under the License.
// etcd:
// # one of local or external
// local:
-// imageRepository: "k8s.gcr.io"
+// imageRepository: "registry.k8s.io"
// imageTag: "3.2.24"
// dataDir: "/var/lib/etcd"
// extraArgs:
@@ -249,7 +249,7 @@ limitations under the License.
// readOnly: false
// pathType: File
// certificatesDir: "/etc/kubernetes/pki"
-// imageRepository: "k8s.gcr.io"
+// imageRepository: "registry.k8s.io"
// clusterName: "example-cluster"
// ---
// apiVersion: kubelet.config.k8s.io/v1beta1
diff --git a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/types.go b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/types.go
index c3177617e..431d83488 100644
--- a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/types.go
+++ b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/types.go
@@ -121,8 +121,8 @@ type ClusterConfiguration struct {
CertificatesDir string `json:"certificatesDir,omitempty"`
// ImageRepository sets the container registry to pull images from.
- // If empty, `k8s.gcr.io` will be used by default; in case of kubernetes version is a CI build (kubernetes version starts with `ci/` or `ci-cross/`)
- // `gcr.io/k8s-staging-ci-images` will be used as a default for control plane components and for kube-proxy, while `k8s.gcr.io`
+ // If empty, `registry.k8s.io` will be used by default; in case of kubernetes version is a CI build (kubernetes version starts with `ci/`)
+ // `gcr.io/k8s-staging-ci-images` will be used as a default for control plane components and for kube-proxy, while `registry.k8s.io`
// will be used for all the other images.
// +optional
ImageRepository string `json:"imageRepository,omitempty"`
diff --git a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/zz_generated.conversion.go b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/zz_generated.conversion.go
index 69d939c57..9a04e377d 100644
--- a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/zz_generated.conversion.go
+++ b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm/v1beta3/zz_generated.conversion.go
@@ -344,6 +344,7 @@ func autoConvert_kubeadm_ClusterConfiguration_To_v1beta3_ClusterConfiguration(in
return err
}
out.KubernetesVersion = in.KubernetesVersion
+ // INFO: in.CIKubernetesVersion opted out of conversion generation
out.ControlPlaneEndpoint = in.ControlPlaneEndpoint
if err := Convert_kubeadm_APIServer_To_v1beta3_APIServer(&in.APIServer, &out.APIServer, s); err != nil {
return err
diff --git a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/constants/constants.go b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/constants/constants.go
index b3559734a..d914f3ae1 100644
--- a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/constants/constants.go
+++ b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/constants/constants.go
@@ -286,7 +286,7 @@ const (
MinExternalEtcdVersion = "3.2.18"
// DefaultEtcdVersion indicates the default etcd version that kubeadm uses
- DefaultEtcdVersion = "3.5.0-0"
+ DefaultEtcdVersion = "3.5.6-0"
// Etcd defines variable used internally when referring to etcd component
Etcd = "etcd"
@@ -340,6 +340,9 @@ const (
// TODO: Find a better place for this constant
YAMLDocumentSeparator = "---\n"
+ // CIKubernetesVersionPrefix is the prefix for CI Kubernetes version
+ CIKubernetesVersionPrefix = "ci/"
+
// DefaultAPIServerBindAddress is the default bind address for the API Server
DefaultAPIServerBindAddress = "0.0.0.0"
@@ -467,8 +470,8 @@ var (
19: "3.4.13-0",
20: "3.4.13-0",
21: "3.4.13-0",
- 22: "3.5.0-0",
- 23: "3.5.0-0",
+ 22: "3.5.6-0",
+ 23: "3.5.6-0",
}
// KubeadmCertsClusterRoleName sets the name for the ClusterRole that allows
diff --git a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/images/images.go b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/images/images.go
index 7e97dbc94..38e65d1ed 100644
--- a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/images/images.go
+++ b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/images/images.go
@@ -48,7 +48,7 @@ func GetDNSImage(cfg *kubeadmapi.ClusterConfiguration) string {
if cfg.DNS.ImageRepository != "" {
dnsImageRepository = cfg.DNS.ImageRepository
}
- // Handle the renaming of the official image from "k8s.gcr.io/coredns" to "k8s.gcr.io/coredns/coredns
+ // Handle the renaming of the official image from "registry.k8s.io/coredns" to "registry.k8s.io/coredns/coredns
if dnsImageRepository == kubeadmapiv1beta2.DefaultImageRepository {
dnsImageRepository = fmt.Sprintf("%s/coredns", dnsImageRepository)
}
diff --git a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/phases/etcd/local.go b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/phases/etcd/local.go
index 144cda6cd..35f607b26 100644
--- a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/phases/etcd/local.go
+++ b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/phases/etcd/local.go
@@ -235,22 +235,25 @@ func getEtcdCommand(cfg *kubeadmapi.ClusterConfiguration, endpoint *kubeadmapi.A
etcdLocalhostAddress = "::1"
}
defaultArguments := map[string]string{
- "name": nodeName,
- "listen-client-urls": fmt.Sprintf("%s,%s", etcdutil.GetClientURLByIP(etcdLocalhostAddress), etcdutil.GetClientURL(endpoint)),
- "advertise-client-urls": etcdutil.GetClientURL(endpoint),
- "listen-peer-urls": etcdutil.GetPeerURL(endpoint),
- "initial-advertise-peer-urls": etcdutil.GetPeerURL(endpoint),
- "data-dir": cfg.Etcd.Local.DataDir,
- "cert-file": filepath.Join(cfg.CertificatesDir, kubeadmconstants.EtcdServerCertName),
- "key-file": filepath.Join(cfg.CertificatesDir, kubeadmconstants.EtcdServerKeyName),
- "trusted-ca-file": filepath.Join(cfg.CertificatesDir, kubeadmconstants.EtcdCACertName),
- "client-cert-auth": "true",
- "peer-cert-file": filepath.Join(cfg.CertificatesDir, kubeadmconstants.EtcdPeerCertName),
- "peer-key-file": filepath.Join(cfg.CertificatesDir, kubeadmconstants.EtcdPeerKeyName),
- "peer-trusted-ca-file": filepath.Join(cfg.CertificatesDir, kubeadmconstants.EtcdCACertName),
- "peer-client-cert-auth": "true",
- "snapshot-count": "10000",
- "listen-metrics-urls": fmt.Sprintf("http://%s", net.JoinHostPort(etcdLocalhostAddress, strconv.Itoa(kubeadmconstants.EtcdMetricsPort))),
+ "name": nodeName,
+ // TODO: start using --initial-corrupt-check once the graduated flag is available:
+ // https://github.com/kubernetes/kubeadm/issues/2676
+ "experimental-initial-corrupt-check": "true",
+ "listen-client-urls": fmt.Sprintf("%s,%s", etcdutil.GetClientURLByIP(etcdLocalhostAddress), etcdutil.GetClientURL(endpoint)),
+ "advertise-client-urls": etcdutil.GetClientURL(endpoint),
+ "listen-peer-urls": etcdutil.GetPeerURL(endpoint),
+ "initial-advertise-peer-urls": etcdutil.GetPeerURL(endpoint),
+ "data-dir": cfg.Etcd.Local.DataDir,
+ "cert-file": filepath.Join(cfg.CertificatesDir, kubeadmconstants.EtcdServerCertName),
+ "key-file": filepath.Join(cfg.CertificatesDir, kubeadmconstants.EtcdServerKeyName),
+ "trusted-ca-file": filepath.Join(cfg.CertificatesDir, kubeadmconstants.EtcdCACertName),
+ "client-cert-auth": "true",
+ "peer-cert-file": filepath.Join(cfg.CertificatesDir, kubeadmconstants.EtcdPeerCertName),
+ "peer-key-file": filepath.Join(cfg.CertificatesDir, kubeadmconstants.EtcdPeerKeyName),
+ "peer-trusted-ca-file": filepath.Join(cfg.CertificatesDir, kubeadmconstants.EtcdCACertName),
+ "peer-client-cert-auth": "true",
+ "snapshot-count": "10000",
+ "listen-metrics-urls": fmt.Sprintf("http://%s", net.JoinHostPort(etcdLocalhostAddress, strconv.Itoa(kubeadmconstants.EtcdMetricsPort))),
}
if len(initialCluster) == 0 {
diff --git a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/util/pkiutil/pki_helpers.go b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/util/pkiutil/pki_helpers.go
index f24bd975d..8cff24daf 100644
--- a/vendor/k8s.io/kubernetes/cmd/kubeadm/app/util/pkiutil/pki_helpers.go
+++ b/vendor/k8s.io/kubernetes/cmd/kubeadm/app/util/pkiutil/pki_helpers.go
@@ -353,7 +353,7 @@ func TryLoadCSRAndKeyFromDisk(pkiPath, name string) (*x509.CertificateRequest, c
}
// TryLoadPrivatePublicKeyFromDisk tries to load the key from the disk and validates that it is valid
-func TryLoadPrivatePublicKeyFromDisk(pkiPath, name string) (*rsa.PrivateKey, *rsa.PublicKey, error) {
+func TryLoadPrivatePublicKeyFromDisk(pkiPath, name string) (crypto.PrivateKey, crypto.PublicKey, error) {
privateKeyPath := pathForKey(pkiPath, name)
// Parse the private key from a file
@@ -370,15 +370,15 @@ func TryLoadPrivatePublicKeyFromDisk(pkiPath, name string) (*rsa.PrivateKey, *rs
return nil, nil, errors.Wrapf(err, "couldn't load the public key file %s", publicKeyPath)
}
- // Allow RSA format only
- k, ok := privKey.(*rsa.PrivateKey)
- if !ok {
- return nil, nil, errors.Errorf("the private key file %s isn't in RSA format", privateKeyPath)
+ // Allow RSA and ECDSA formats only
+ switch k := privKey.(type) {
+ case *rsa.PrivateKey:
+ return k, pubKeys[0].(*rsa.PublicKey), nil
+ case *ecdsa.PrivateKey:
+ return k, pubKeys[0].(*ecdsa.PublicKey), nil
+ default:
+ return nil, nil, errors.Errorf("the private key file %s is neither in RSA nor ECDSA format", privateKeyPath)
}
-
- p := pubKeys[0].(*rsa.PublicKey)
-
- return k, p, nil
}
// TryLoadCSRFromDisk tries to load the CSR from the disk
diff --git a/vendor/k8s.io/kubernetes/cmd/kubelet/app/options/container_runtime.go b/vendor/k8s.io/kubernetes/cmd/kubelet/app/options/container_runtime.go
index e557c113e..f3257ca0f 100644
--- a/vendor/k8s.io/kubernetes/cmd/kubelet/app/options/container_runtime.go
+++ b/vendor/k8s.io/kubernetes/cmd/kubelet/app/options/container_runtime.go
@@ -27,7 +27,7 @@ import (
const (
// When these values are updated, also update test/utils/image/manifest.go
- defaultPodSandboxImageName = "k8s.gcr.io/pause"
+ defaultPodSandboxImageName = "registry.k8s.io/pause"
defaultPodSandboxImageVersion = "3.5"
)
diff --git a/vendor/k8s.io/kubernetes/cmd/kubelet/app/server.go b/vendor/k8s.io/kubernetes/cmd/kubelet/app/server.go
index 592001984..8d2e2ea9b 100644
--- a/vendor/k8s.io/kubernetes/cmd/kubelet/app/server.go
+++ b/vendor/k8s.io/kubernetes/cmd/kubelet/app/server.go
@@ -245,7 +245,7 @@ HTTP server: The kubelet can also listen for HTTP and respond to a simple API
config.StaticPodURLHeader[k] = []string{"<masked>"}
}
// log the kubelet's config for inspection
- klog.V(5).InfoS("KubeletConfiguration", "configuration", kubeletServer.KubeletConfiguration)
+ klog.V(5).InfoS("KubeletConfiguration", "configuration", config)
// run the kubelet
if err := Run(ctx, kubeletServer, kubeletDeps, utilfeature.DefaultFeatureGate); err != nil {
diff --git a/vendor/k8s.io/kubernetes/cmd/kubelet/app/server_windows.go b/vendor/k8s.io/kubernetes/cmd/kubelet/app/server_windows.go
index b986f102b..3906dd7c3 100644
--- a/vendor/k8s.io/kubernetes/cmd/kubelet/app/server_windows.go
+++ b/vendor/k8s.io/kubernetes/cmd/kubelet/app/server_windows.go
@@ -19,66 +19,27 @@ limitations under the License.
package app
import (
- "fmt"
+ "errors"
"os/user"
"golang.org/x/sys/windows"
+ "k8s.io/klog/v2"
)
-func isAdmin() (bool, error) {
- // Get current user
- u, err := user.Current()
- if err != nil {
- return false, fmt.Errorf("error retrieving current user: %s", err)
- }
- // Get IDs of group user is a member of
- ids, err := u.GroupIds()
- if err != nil {
- return false, fmt.Errorf("error retrieving group ids: %s", err)
- }
-
- // Check for existence of BUILTIN\ADMINISTRATORS group id
- for i := range ids {
- // BUILTIN\ADMINISTRATORS
- if "S-1-5-32-544" == ids[i] {
- return true, nil
- }
- }
- return false, nil
-}
-
func checkPermissions() error {
- //https://github.com/golang/go/issues/28804#issuecomment-505326268
- var sid *windows.SID
- var userIsAdmin bool
- // https://docs.microsoft.com/en-us/windows/desktop/api/securitybaseapi/nf-securitybaseapi-checktokenmembership
- err := windows.AllocateAndInitializeSid(
- &windows.SECURITY_NT_AUTHORITY,
- 2,
- windows.SECURITY_BUILTIN_DOMAIN_RID,
- windows.DOMAIN_ALIAS_RID_ADMINS,
- 0, 0, 0, 0, 0, 0,
- &sid)
+ u, err := user.Current()
if err != nil {
- return fmt.Errorf("error while checking for elevated permissions: %s", err)
+ klog.ErrorS(err, "Unable to get current user")
+ return err
}
- //We must free the sid to prevent security token leaks
- defer windows.FreeSid(sid)
- token := windows.Token(0)
+ // For Windows user.UserName contains the login name and user.Name contains
+ // the user's display name - https://pkg.go.dev/os/user#User
+ klog.InfoS("Kubelet is running as", "login name", u.Username, "dispaly name", u.Name)
- userIsAdmin, err = isAdmin()
- if err != nil {
- return fmt.Errorf("error while checking admin group membership: %s", err)
- }
-
- member, err := token.IsMember(sid)
- if err != nil {
- return fmt.Errorf("error while checking for elevated permissions: %s", err)
- }
- if !member {
- return fmt.Errorf("kubelet needs to run with administrator permissions. Run as admin is: %t, User in admin group: %t", member, userIsAdmin)
+ if !windows.GetCurrentProcessToken().IsElevated() {
+ return errors.New("kubelet needs to run with elevated permissions!")
}
return nil
diff --git a/vendor/k8s.io/kubernetes/pkg/api/v1/pod/util.go b/vendor/k8s.io/kubernetes/pkg/api/v1/pod/util.go
index 9560121bb..8bfc21a67 100644
--- a/vendor/k8s.io/kubernetes/pkg/api/v1/pod/util.go
+++ b/vendor/k8s.io/kubernetes/pkg/api/v1/pod/util.go
@@ -301,12 +301,28 @@ func IsPodReady(pod *v1.Pod) bool {
return IsPodReadyConditionTrue(pod.Status)
}
+// IsPodTerminal returns true if a pod is terminal, all containers are stopped and cannot ever regress.
+func IsPodTerminal(pod *v1.Pod) bool {
+ return IsPodPhaseTerminal(pod.Status.Phase)
+}
+
+// IsPhaseTerminal returns true if the pod's phase is terminal.
+func IsPodPhaseTerminal(phase v1.PodPhase) bool {
+ return phase == v1.PodFailed || phase == v1.PodSucceeded
+}
+
// IsPodReadyConditionTrue returns true if a pod is ready; false otherwise.
func IsPodReadyConditionTrue(status v1.PodStatus) bool {
condition := GetPodReadyCondition(status)
return condition != nil && condition.Status == v1.ConditionTrue
}
+// IsContainersReadyConditionTrue returns true if a pod is ready; false otherwise.
+func IsContainersReadyConditionTrue(status v1.PodStatus) bool {
+ condition := GetContainersReadyCondition(status)
+ return condition != nil && condition.Status == v1.ConditionTrue
+}
+
// GetPodReadyCondition extracts the pod ready condition from the given status and returns that.
// Returns nil if the condition is not present.
func GetPodReadyCondition(status v1.PodStatus) *v1.PodCondition {
@@ -314,6 +330,13 @@ func GetPodReadyCondition(status v1.PodStatus) *v1.PodCondition {
return condition
}
+// GetContainersReadyCondition extracts the containers ready condition from the given status and returns that.
+// Returns nil if the condition is not present.
+func GetContainersReadyCondition(status v1.PodStatus) *v1.PodCondition {
+ _, condition := GetPodCondition(&status, v1.ContainersReady)
+ return condition
+}
+
// GetPodCondition extracts the provided condition from the given status and returns that.
// Returns nil and -1 if the condition is not present, and the index of the located condition.
func GetPodCondition(status *v1.PodStatus, conditionType v1.PodConditionType) (int, *v1.PodCondition) {
diff --git a/vendor/k8s.io/kubernetes/pkg/apis/core/validation/events.go b/vendor/k8s.io/kubernetes/pkg/apis/core/validation/events.go
index adb0177b7..a01d825d7 100644
--- a/vendor/k8s.io/kubernetes/pkg/apis/core/validation/events.go
+++ b/vendor/k8s.io/kubernetes/pkg/apis/core/validation/events.go
@@ -95,7 +95,18 @@ func ValidateEventUpdate(newEvent, oldEvent *core.Event, requestVersion schema.G
allErrs = append(allErrs, ValidateImmutableField(newEvent.Count, oldEvent.Count, field.NewPath("count"))...)
allErrs = append(allErrs, ValidateImmutableField(newEvent.Reason, oldEvent.Reason, field.NewPath("reason"))...)
allErrs = append(allErrs, ValidateImmutableField(newEvent.Type, oldEvent.Type, field.NewPath("type"))...)
- allErrs = append(allErrs, ValidateImmutableField(newEvent.EventTime, oldEvent.EventTime, field.NewPath("eventTime"))...)
+
+ // Disallow changes to eventTime greater than microsecond-level precision.
+ // Tolerating sub-microsecond changes is required to tolerate updates
+ // from clients that correctly truncate to microsecond-precision when serializing,
+ // or from clients built with incorrect nanosecond-precision protobuf serialization.
+ // See https://github.com/kubernetes/kubernetes/issues/111928
+ newTruncated := newEvent.EventTime.Truncate(time.Microsecond).UTC()
+ oldTruncated := oldEvent.EventTime.Truncate(time.Microsecond).UTC()
+ if newTruncated != oldTruncated {
+ allErrs = append(allErrs, ValidateImmutableField(newEvent.EventTime, oldEvent.EventTime, field.NewPath("eventTime"))...)
+ }
+
allErrs = append(allErrs, ValidateImmutableField(newEvent.Action, oldEvent.Action, field.NewPath("action"))...)
allErrs = append(allErrs, ValidateImmutableField(newEvent.Related, oldEvent.Related, field.NewPath("related"))...)
allErrs = append(allErrs, ValidateImmutableField(newEvent.ReportingController, oldEvent.ReportingController, field.NewPath("reportingController"))...)
diff --git a/vendor/k8s.io/kubernetes/pkg/credentialprovider/azure/azure_acr_helper.go b/vendor/k8s.io/kubernetes/pkg/credentialprovider/azure/azure_acr_helper.go
index ec4af71a2..aa75e98ee 100644
--- a/vendor/k8s.io/kubernetes/pkg/credentialprovider/azure/azure_acr_helper.go
+++ b/vendor/k8s.io/kubernetes/pkg/credentialprovider/azure/azure_acr_helper.go
@@ -80,7 +80,7 @@ const dockerTokenLoginUsernameGUID = "00000000-0000-0000-0000-000000000000"
var client = &http.Client{
Transport: utilnet.SetTransportDefaults(&http.Transport{}),
- Timeout: time.Second * 10,
+ Timeout: time.Second * 60,
}
func receiveChallengeFromLoginServer(serverAddress string) (*authDirective, error) {
diff --git a/vendor/k8s.io/kubernetes/pkg/credentialprovider/azure/azure_credentials.go b/vendor/k8s.io/kubernetes/pkg/credentialprovider/azure/azure_credentials.go
index 2d9281482..1ad75ac6d 100644
--- a/vendor/k8s.io/kubernetes/pkg/credentialprovider/azure/azure_credentials.go
+++ b/vendor/k8s.io/kubernetes/pkg/credentialprovider/azure/azure_credentials.go
@@ -29,7 +29,6 @@ import (
"time"
"github.com/Azure/azure-sdk-for-go/services/containerregistry/mgmt/2019-05-01/containerregistry"
- "github.com/Azure/go-autorest/autorest"
"github.com/Azure/go-autorest/autorest/adal"
"github.com/Azure/go-autorest/autorest/azure"
"github.com/spf13/pflag"
@@ -88,39 +87,6 @@ type RegistriesClient interface {
List(ctx context.Context) ([]containerregistry.Registry, error)
}
-// azRegistriesClient implements RegistriesClient.
-type azRegistriesClient struct {
- client containerregistry.RegistriesClient
-}
-
-func newAzRegistriesClient(subscriptionID, endpoint string, token *adal.ServicePrincipalToken) *azRegistriesClient {
- registryClient := containerregistry.NewRegistriesClient(subscriptionID)
- registryClient.BaseURI = endpoint
- registryClient.Authorizer = autorest.NewBearerAuthorizer(token)
-
- return &azRegistriesClient{
- client: registryClient,
- }
-}
-
-func (az *azRegistriesClient) List(ctx context.Context) ([]containerregistry.Registry, error) {
- iterator, err := az.client.ListComplete(ctx)
- if err != nil {
- return nil, err
- }
-
- result := make([]containerregistry.Registry, 0)
- for ; iterator.NotDone(); err = iterator.Next() {
- if err != nil {
- return nil, err
- }
-
- result = append(result, iterator.Value())
- }
-
- return result, nil
-}
-
// NewACRProvider parses the specified configFile and returns a DockerConfigProvider
func NewACRProvider(configFile *string) credentialprovider.DockerConfigProvider {
return &acrProvider{
@@ -133,7 +99,6 @@ type acrProvider struct {
file *string
config *auth.AzureAuthConfig
environment *azure.Environment
- registryClient RegistriesClient
servicePrincipalToken *adal.ServicePrincipalToken
cache cache.Store
}
@@ -199,10 +164,7 @@ func (a *acrProvider) Enabled() bool {
a.servicePrincipalToken, err = auth.GetServicePrincipalToken(a.config, a.environment)
if err != nil {
klog.Errorf("Failed to create service principal token: %v", err)
- return false
}
-
- a.registryClient = newAzRegistriesClient(a.config.SubscriptionID, a.environment.ResourceManagerEndpoint, a.servicePrincipalToken)
return true
}
@@ -313,11 +275,21 @@ func getLoginServer(registry containerregistry.Registry) string {
}
func getACRDockerEntryFromARMToken(a *acrProvider, loginServer string) (*credentialprovider.DockerConfigEntry, error) {
- // Run EnsureFresh to make sure the token is valid and does not expire
- if err := a.servicePrincipalToken.EnsureFresh(); err != nil {
- klog.Errorf("Failed to ensure fresh service principal token: %v", err)
- return nil, err
+ if a.servicePrincipalToken == nil {
+ token, err := auth.GetServicePrincipalToken(a.config, a.environment)
+ if err != nil {
+ klog.Errorf("Failed to create service principal token: %v", err)
+ return nil, err
+ }
+ a.servicePrincipalToken = token
+ } else {
+ // Run EnsureFresh to make sure the token is valid and does not expire
+ if err := a.servicePrincipalToken.EnsureFresh(); err != nil {
+ klog.Errorf("Failed to ensure fresh service principal token: %v", err)
+ return nil, err
+ }
}
+
armAccessToken := a.servicePrincipalToken.OAuthToken()
klog.V(4).Infof("discovering auth redirects for: %s", loginServer)
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/cm/containermap/container_map.go b/vendor/k8s.io/kubernetes/pkg/kubelet/cm/containermap/container_map.go
index c324e028b..7dc4689af 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/cm/containermap/container_map.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/cm/containermap/container_map.go
@@ -69,3 +69,10 @@ func (cm ContainerMap) GetContainerRef(containerID string) (string, string, erro
}
return cm[containerID].podUID, cm[containerID].containerName, nil
}
+
+// Visit invoke visitor function to walks all of the entries in the container map
+func (cm ContainerMap) Visit(visitor func(podUID, containerName, containerID string)) {
+ for k, v := range cm {
+ visitor(v.podUID, v.containerName, k)
+ }
+}
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/cm/cpumanager/cpu_manager.go b/vendor/k8s.io/kubernetes/pkg/kubelet/cm/cpumanager/cpu_manager.go
index 4777c132e..6ed4a72bb 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/cm/cpumanager/cpu_manager.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/cm/cpumanager/cpu_manager.go
@@ -384,6 +384,16 @@ func (m *manager) removeStaleState() {
}
}
}
+
+ m.containerMap.Visit(func(podUID, containerName, containerID string) {
+ if _, ok := activeContainers[podUID][containerName]; !ok {
+ klog.ErrorS(nil, "RemoveStaleState: removing container", "podUID", podUID, "containerName", containerName)
+ err := m.policyRemoveContainerByRef(podUID, containerName)
+ if err != nil {
+ klog.ErrorS(err, "RemoveStaleState: failed to remove container", "podUID", podUID, "containerName", containerName)
+ }
+ }
+ })
}
func (m *manager) reconcileState() (success []reconciledContainer, failure []reconciledContainer) {
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/cm/memorymanager/memory_manager.go b/vendor/k8s.io/kubernetes/pkg/kubelet/cm/memorymanager/memory_manager.go
index 3bf4dc0bf..41fdac903 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/cm/memorymanager/memory_manager.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/cm/memorymanager/memory_manager.go
@@ -339,6 +339,13 @@ func (m *manager) removeStaleState() {
}
}
}
+
+ m.containerMap.Visit(func(podUID, containerName, containerID string) {
+ if _, ok := activeContainers[podUID][containerName]; !ok {
+ klog.InfoS("RemoveStaleState removing state", "podUID", podUID, "containerName", containerName)
+ m.policyRemoveContainerByRef(podUID, containerName)
+ }
+ })
}
func (m *manager) policyRemoveContainerByRef(podUID string, containerName string) {
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/dockershim/docker_sandbox.go b/vendor/k8s.io/kubernetes/pkg/kubelet/dockershim/docker_sandbox.go
index c9f0d5e0c..c21c4b749 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/dockershim/docker_sandbox.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/dockershim/docker_sandbox.go
@@ -40,7 +40,7 @@ import (
)
const (
- defaultSandboxImage = "k8s.gcr.io/pause:3.5"
+ defaultSandboxImage = "registry.k8s.io/pause:3.5"
// Various default sandbox resources requests/limits.
defaultSandboxCPUshares int64 = 2
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/kubelet.go b/vendor/k8s.io/kubernetes/pkg/kubelet/kubelet.go
index 278ee5bbc..b8ffc1ed8 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/kubelet.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/kubelet.go
@@ -1367,15 +1367,26 @@ func (kl *Kubelet) Run(updates <-chan kubetypes.PodUpdate) {
}
// syncPod is the transaction script for the sync of a single pod (setting up)
-// a pod. The reverse (teardown) is handled in syncTerminatingPod and
-// syncTerminatedPod. If syncPod exits without error, then the pod runtime
-// state is in sync with the desired configuration state (pod is running).
-// If syncPod exits with a transient error, the next invocation of syncPod
-// is expected to make progress towards reaching the runtime state.
+// a pod. This method is reentrant and expected to converge a pod towards the
+// desired state of the spec. The reverse (teardown) is handled in
+// syncTerminatingPod and syncTerminatedPod. If syncPod exits without error,
+// then the pod runtime state is in sync with the desired configuration state
+// (pod is running). If syncPod exits with a transient error, the next
+// invocation of syncPod is expected to make progress towards reaching the
+// runtime state. syncPod exits with isTerminal when the pod was detected to
+// have reached a terminal lifecycle phase due to container exits (for
+// RestartNever or RestartOnFailure) and the next method invoked will by
+// syncTerminatingPod.
//
// Arguments:
//
-// o - the SyncPodOptions for this invocation
+// updateType - whether this is a create (first time) or an update, should
+// only be used for metrics since this method must be reentrant
+// pod - the pod that is being set up
+// mirrorPod - the mirror pod known to the kubelet for this pod, if any
+// podStatus - the most recent pod status observed for this pod which can
+// be used to determine the set of actions that should be taken during
+// this loop of syncPod
//
// The workflow is:
// * If the pod is being created, record pod worker start latency
@@ -1383,7 +1394,9 @@ func (kl *Kubelet) Run(updates <-chan kubetypes.PodUpdate) {
// * If the pod is being seen as running for the first time, record pod
// start latency
// * Update the status of the pod in the status manager
-// * Kill the pod if it should not be running due to soft admission
+// * Stop the pod's containers if it should not be running due to soft
+// admission
+// * Ensure any background tracking for a runnable pod is started
// * Create a mirror pod if the pod is a static pod, and does not
// already have a mirror pod
// * Create the data directories for the pod if they do not exist
@@ -1397,10 +1410,12 @@ func (kl *Kubelet) Run(updates <-chan kubetypes.PodUpdate) {
//
// This operation writes all events that are dispatched in order to provide
// the most accurate information possible about an error situation to aid debugging.
-// Callers should not throw an event if this operation returns an error.
-func (kl *Kubelet) syncPod(ctx context.Context, updateType kubetypes.SyncPodType, pod *v1.Pod, podStatus *kubecontainer.PodStatus) error {
+// Callers should not write an event if this operation returns an error.
+func (kl *Kubelet) syncPod(ctx context.Context, updateType kubetypes.SyncPodType, pod *v1.Pod, podStatus *kubecontainer.PodStatus) (isTerminal bool, err error) {
klog.V(4).InfoS("syncPod enter", "pod", klog.KObj(pod), "podUID", pod.UID)
- defer klog.V(4).InfoS("syncPod exit", "pod", klog.KObj(pod), "podUID", pod.UID)
+ defer func() {
+ klog.V(4).InfoS("syncPod exit", "pod", klog.KObj(pod), "podUID", pod.UID, "isTerminal", isTerminal)
+ }()
// Latency measurements for the main workflow are relative to the
// first time the pod was seen by the API server.
@@ -1432,11 +1447,17 @@ func (kl *Kubelet) syncPod(ctx context.Context, updateType kubetypes.SyncPodType
for _, ipInfo := range apiPodStatus.PodIPs {
podStatus.IPs = append(podStatus.IPs, ipInfo.IP)
}
-
if len(podStatus.IPs) == 0 && len(apiPodStatus.PodIP) > 0 {
podStatus.IPs = []string{apiPodStatus.PodIP}
}
+ // If the pod is terminal, we don't need to continue to setup the pod
+ if apiPodStatus.Phase == v1.PodSucceeded || apiPodStatus.Phase == v1.PodFailed {
+ kl.statusManager.SetPodStatus(pod, apiPodStatus)
+ isTerminal = true
+ return isTerminal, nil
+ }
+
// If the pod should not be running, we request the pod's containers be stopped. This is not the same
// as termination (we want to stop the pod, but potentially restart it later if soft admission allows
// it later). Set the status and phase appropriately
@@ -1485,13 +1506,23 @@ func (kl *Kubelet) syncPod(ctx context.Context, updateType kubetypes.SyncPodType
// Return an error to signal that the sync loop should back off.
syncErr = fmt.Errorf("pod cannot be run: %s", runnable.Message)
}
- return syncErr
+ return false, syncErr
}
// If the network plugin is not ready, only start the pod if it uses the host network
if err := kl.runtimeState.networkErrors(); err != nil && !kubecontainer.IsHostNetworkPod(pod) {
kl.recorder.Eventf(pod, v1.EventTypeWarning, events.NetworkNotReady, "%s: %v", NetworkNotReadyErrorMsg, err)
- return fmt.Errorf("%s: %v", NetworkNotReadyErrorMsg, err)
+ return false, fmt.Errorf("%s: %v", NetworkNotReadyErrorMsg, err)
+ }
+
+ // ensure the kubelet knows about referenced secrets or configmaps used by the pod
+ if !kl.podWorkers.IsPodTerminationRequested(pod.UID) {
+ if kl.secretManager != nil {
+ kl.secretManager.RegisterPod(pod)
+ }
+ if kl.configMapManager != nil {
+ kl.configMapManager.RegisterPod(pod)
+ }
}
// Create Cgroups for the pod and apply resource parameters
@@ -1538,7 +1569,7 @@ func (kl *Kubelet) syncPod(ctx context.Context, updateType kubetypes.SyncPodType
}
if err := pcm.EnsureExists(pod); err != nil {
kl.recorder.Eventf(pod, v1.EventTypeWarning, events.FailedToCreatePodContainer, "unable to ensure pod container exists: %v", err)
- return fmt.Errorf("failed to ensure that the pod: %v cgroups exist and are correctly applied: %v", pod.UID, err)
+ return false, fmt.Errorf("failed to ensure that the pod: %v cgroups exist and are correctly applied: %v", pod.UID, err)
}
}
}
@@ -1548,7 +1579,7 @@ func (kl *Kubelet) syncPod(ctx context.Context, updateType kubetypes.SyncPodType
if err := kl.makePodDataDirs(pod); err != nil {
kl.recorder.Eventf(pod, v1.EventTypeWarning, events.FailedToMakePodDataDirectories, "error making pod data directories: %v", err)
klog.ErrorS(err, "Unable to make pod data directories for pod", "pod", klog.KObj(pod))
- return err
+ return false, err
}
// Volume manager will not mount volumes for terminating pods
@@ -1558,13 +1589,16 @@ func (kl *Kubelet) syncPod(ctx context.Context, updateType kubetypes.SyncPodType
if err := kl.volumeManager.WaitForAttachAndMount(pod); err != nil {
kl.recorder.Eventf(pod, v1.EventTypeWarning, events.FailedMountVolume, "Unable to attach or mount volumes: %v", err)
klog.ErrorS(err, "Unable to attach or mount volumes for pod; skipping pod", "pod", klog.KObj(pod))
- return err
+ return false, err
}
}
// Fetch the pull secrets for the pod
pullSecrets := kl.getPullSecretsForPod(pod)
+ // Ensure the pod is being probed
+ kl.probeManager.AddPod(pod)
+
// Call the container runtime's SyncPod callback
result := kl.containerRuntime.SyncPod(pod, podStatus, pullSecrets, kl.backOff)
kl.reasonCache.Update(pod.UID, result)
@@ -1573,15 +1607,15 @@ func (kl *Kubelet) syncPod(ctx context.Context, updateType kubetypes.SyncPodType
for _, r := range result.SyncResults {
if r.Error != kubecontainer.ErrCrashLoopBackOff && r.Error != images.ErrImagePullBackOff {
// Do not record an event here, as we keep all event logging for sync pod failures
- // local to container runtime so we get better errors
- return err
+ // local to container runtime, so we get better errors.
+ return false, err
}
}
- return nil
+ return false, nil
}
- return nil
+ return false, nil
}
// syncTerminatingPod is expected to terminate all running containers in a pod. Once this method
@@ -1622,6 +1656,9 @@ func (kl *Kubelet) syncTerminatingPod(ctx context.Context, pod *v1.Pod, podStatu
} else {
klog.V(4).InfoS("Pod terminating with grace period", "pod", klog.KObj(pod), "podUID", pod.UID, "gracePeriod", nil)
}
+
+ kl.probeManager.StopLivenessAndStartup(pod)
+
p := kubecontainer.ConvertPodStatusToRunningPod(kl.getRuntime().Type(), podStatus)
if err := kl.killPod(pod, p, gracePeriod); err != nil {
kl.recorder.Eventf(pod, v1.EventTypeWarning, events.FailedToKillPod, "error killing pod: %v", err)
@@ -1630,6 +1667,12 @@ func (kl *Kubelet) syncTerminatingPod(ctx context.Context, pod *v1.Pod, podStatu
return err
}
+ // Once the containers are stopped, we can stop probing for liveness and readiness.
+ // TODO: once a pod is terminal, certain probes (liveness exec) could be stopped immediately after
+ // the detection of a container shutdown or (for readiness) after the first failure. Tracked as
+ // https://github.com/kubernetes/kubernetes/issues/107894 although may not be worth optimizing.
+ kl.probeManager.RemovePod(pod)
+
// Guard against consistency issues in KillPod implementations by checking that there are no
// running containers. This method is invoked infrequently so this is effectively free and can
// catch race conditions introduced by callers updating pod status out of order.
@@ -1682,6 +1725,14 @@ func (kl *Kubelet) syncTerminatedPod(ctx context.Context, pod *v1.Pod, podStatus
}
klog.V(4).InfoS("Pod termination unmounted volumes", "pod", klog.KObj(pod), "podUID", pod.UID)
+ // After volume unmount is complete, let the secret and configmap managers know we're done with this pod
+ if kl.secretManager != nil {
+ kl.secretManager.UnregisterPod(pod)
+ }
+ if kl.configMapManager != nil {
+ kl.configMapManager.UnregisterPod(pod)
+ }
+
// Note: we leave pod containers to be reclaimed in the background since dockershim requires the
// container for retrieving logs and we want to make sure logs are available until the pod is
// physically deleted.
@@ -2032,7 +2083,6 @@ func (kl *Kubelet) HandlePodAdditions(pods []*v1.Pod) {
// the apiserver and no action (other than cleanup) is required.
kl.podManager.AddPod(pod)
-
// Only go through the admission process if the pod is not requested
// for termination by another part of the kubelet. If the pod is already
// using resources (previously admitted), the pod worker is going to be
@@ -2051,9 +2101,6 @@ func (kl *Kubelet) HandlePodAdditions(pods []*v1.Pod) {
}
}
kl.dispatchWork(pod, kubetypes.SyncPodCreate, nil, start)
- // TODO: move inside syncPod and make reentrant
- // https://github.com/kubernetes/kubernetes/issues/105014
- kl.probeManager.AddPod(pod)
}
}
@@ -2077,10 +2124,6 @@ func (kl *Kubelet) HandlePodRemoves(pods []*v1.Pod) {
if err := kl.deletePod(pod); err != nil {
klog.V(2).InfoS("Failed to delete pod", "pod", klog.KObj(pod), "err", err)
}
- // TODO: move inside syncTerminatingPod|syncTerminatedPod (we should stop probing
- // once the pod kill is acknowledged and during eviction)
- // https://github.com/kubernetes/kubernetes/issues/105014
- kl.probeManager.RemovePod(pod)
}
}
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/kubelet_pods.go b/vendor/k8s.io/kubernetes/pkg/kubelet/kubelet_pods.go
index fcb854073..2a39f4867 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/kubelet_pods.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/kubelet_pods.go
@@ -845,6 +845,12 @@ func countRunningContainerStatus(status v1.PodStatus) int {
return runningContainers
}
+// PodCouldHaveRunningContainers returns true if the pod with the given UID could still have running
+// containers. This returns false if the pod has not yet been started or the pod is unknown.
+func (kl *Kubelet) PodCouldHaveRunningContainers(pod *v1.Pod) bool {
+ return kl.podWorkers.CouldHaveRunningContainers(pod.UID)
+}
+
// PodResourcesAreReclaimed returns true if all required node-level resources that a pod was consuming have
// been reclaimed by the kubelet. Reclaiming resources is a prerequisite to deleting a pod from the API server.
func (kl *Kubelet) PodResourcesAreReclaimed(pod *v1.Pod, status v1.PodStatus) bool {
@@ -1008,8 +1014,8 @@ func (kl *Kubelet) HandlePodCleanups() error {
}
// Stop probing pods that are not running
- klog.V(3).InfoS("Clean up probes for terminating and terminated pods")
- kl.probeManager.CleanupPods(runningPods)
+ klog.V(3).InfoS("Clean up probes for terminated pods")
+ kl.probeManager.CleanupPods(possiblyRunningPods)
// Terminate any pods that are observed in the runtime but not
// present in the list of known running pods from config.
@@ -1102,9 +1108,6 @@ func (kl *Kubelet) HandlePodCleanups() error {
start := kl.clock.Now()
klog.V(3).InfoS("Pod is restartable after termination due to UID reuse", "pod", klog.KObj(pod), "podUID", pod.UID)
kl.dispatchWork(pod, kubetypes.SyncPodCreate, nil, start)
- // TODO: move inside syncPod and make reentrant
- // https://github.com/kubernetes/kubernetes/issues/105014
- kl.probeManager.AddPod(pod)
}
return nil
@@ -1330,7 +1333,7 @@ func getPhase(spec *v1.PodSpec, info []v1.ContainerStatus) v1.PodPhase {
}
// generateAPIPodStatus creates the final API pod status for a pod, given the
-// internal pod status.
+// internal pod status. This method should only be called from within sync*Pod methods.
func (kl *Kubelet) generateAPIPodStatus(pod *v1.Pod, podStatus *kubecontainer.PodStatus) v1.PodStatus {
klog.V(3).InfoS("Generating pod status", "pod", klog.KObj(pod))
@@ -1535,12 +1538,8 @@ func (kl *Kubelet) convertToAPIContainerStatuses(pod *v1.Pod, podStatus *kubecon
case cs.State == kubecontainer.ContainerStateRunning:
status.State.Running = &v1.ContainerStateRunning{StartedAt: metav1.NewTime(cs.StartedAt)}
case cs.State == kubecontainer.ContainerStateCreated:
- // Treat containers in the "created" state as if they are exited.
- // The pod workers are supposed start all containers it creates in
- // one sync (syncPod) iteration. There should not be any normal
- // "created" containers when the pod worker generates the status at
- // the beginning of a sync iteration.
- fallthrough
+ // containers that are created but not running are "waiting to be running"
+ status.State.Waiting = &v1.ContainerStateWaiting{}
case cs.State == kubecontainer.ContainerStateExited:
status.State.Terminated = &v1.ContainerStateTerminated{
ExitCode: int32(cs.ExitCode),
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/kuberuntime/kuberuntime_container.go b/vendor/k8s.io/kubernetes/pkg/kubelet/kuberuntime/kuberuntime_container.go
index 44fbfee46..3cf0c4a7c 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/kuberuntime/kuberuntime_container.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/kuberuntime/kuberuntime_container.go
@@ -34,6 +34,8 @@ import (
"sync"
"time"
+ crierror "k8s.io/cri-api/pkg/errors"
+
grpcstatus "google.golang.org/grpc/status"
"github.com/armon/circbuf"
@@ -502,10 +504,17 @@ func (m *kubeGenericRuntimeManager) getPodContainerStatuses(uid kubetypes.UID, n
return nil, err
}
- statuses := make([]*kubecontainer.Status, len(containers))
+ statuses := []*kubecontainer.Status{}
// TODO: optimization: set maximum number of containers per container name to examine.
- for i, c := range containers {
+ for _, c := range containers {
status, err := m.runtimeService.ContainerStatus(c.Id)
+ // Between List (ListContainers) and check (ContainerStatus) another thread might remove a container, and that is normal.
+ // The previous call (ListContainers) never fails due to a pod container not existing.
+ // Therefore, this method should not either, but instead act as if the previous call failed,
+ // which means the error should be ignored.
+ if crierror.IsNotFound(err) {
+ continue
+ }
if err != nil {
// Merely log this here; GetPodStatus will actually report the error out.
klog.V(4).InfoS("ContainerStatus return error", "containerID", c.Id, "err", err)
@@ -538,7 +547,7 @@ func (m *kubeGenericRuntimeManager) getPodContainerStatuses(uid kubetypes.UID, n
cStatus.Message += tMessage
}
}
- statuses[i] = cStatus
+ statuses = append(statuses, cStatus)
}
sort.Sort(containerStatusByCreated(statuses))
@@ -715,15 +724,15 @@ func (m *kubeGenericRuntimeManager) killContainer(pod *v1.Pod, containerID kubec
"containerName", containerName, "containerID", containerID.String(), "gracePeriod", gracePeriod)
err := m.runtimeService.StopContainer(containerID.ID, gracePeriod)
- if err != nil {
+ if err != nil && !crierror.IsNotFound(err) {
klog.ErrorS(err, "Container termination failed with gracePeriod", "pod", klog.KObj(pod), "podUID", pod.UID,
"containerName", containerName, "containerID", containerID.String(), "gracePeriod", gracePeriod)
- } else {
- klog.V(3).InfoS("Container exited normally", "pod", klog.KObj(pod), "podUID", pod.UID,
- "containerName", containerName, "containerID", containerID.String())
+ return err
}
+ klog.V(3).InfoS("Container exited normally", "pod", klog.KObj(pod), "podUID", pod.UID,
+ "containerName", containerName, "containerID", containerID.String())
- return err
+ return nil
}
// killContainersWithSyncResult kills all pod's containers with sync results.
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/kuberuntime/kuberuntime_manager.go b/vendor/k8s.io/kubernetes/pkg/kubelet/kuberuntime/kuberuntime_manager.go
index 96343832a..0dcb31fec 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/kuberuntime/kuberuntime_manager.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/kuberuntime/kuberuntime_manager.go
@@ -24,6 +24,7 @@ import (
"time"
cadvisorapi "github.com/google/cadvisor/info/v1"
+ crierror "k8s.io/cri-api/pkg/errors"
"k8s.io/klog/v2"
v1 "k8s.io/api/core/v1"
@@ -987,7 +988,7 @@ func (m *kubeGenericRuntimeManager) killPodWithSyncResult(pod *v1.Pod, runningPo
result.AddSyncResult(killSandboxResult)
// Stop all sandboxes belongs to same pod
for _, podSandbox := range runningPod.Sandboxes {
- if err := m.runtimeService.StopPodSandbox(podSandbox.ID.ID); err != nil {
+ if err := m.runtimeService.StopPodSandbox(podSandbox.ID.ID); err != nil && !crierror.IsNotFound(err) {
killSandboxResult.Fail(kubecontainer.ErrKillPodSandbox, err.Error())
klog.ErrorS(nil, "Failed to stop sandbox", "podSandboxID", podSandbox.ID)
}
@@ -1029,16 +1030,22 @@ func (m *kubeGenericRuntimeManager) GetPodStatus(uid kubetypes.UID, name, namesp
klog.V(4).InfoS("getSandboxIDByPodUID got sandbox IDs for pod", "podSandboxID", podSandboxIDs, "pod", klog.KObj(pod))
- sandboxStatuses := make([]*runtimeapi.PodSandboxStatus, len(podSandboxIDs))
+ sandboxStatuses := []*runtimeapi.PodSandboxStatus{}
podIPs := []string{}
for idx, podSandboxID := range podSandboxIDs {
podSandboxStatus, err := m.runtimeService.PodSandboxStatus(podSandboxID)
+ // Between List (getSandboxIDByPodUID) and check (PodSandboxStatus) another thread might remove a container, and that is normal.
+ // The previous call (getSandboxIDByPodUID) never fails due to a pod sandbox not existing.
+ // Therefore, this method should not either, but instead act as if the previous call failed,
+ // which means the error should be ignored.
+ if crierror.IsNotFound(err) {
+ continue
+ }
if err != nil {
klog.ErrorS(err, "PodSandboxStatus of sandbox for pod", "podSandboxID", podSandboxID, "pod", klog.KObj(pod))
return nil, err
}
- sandboxStatuses[idx] = podSandboxStatus
-
+ sandboxStatuses = append(sandboxStatuses, podSandboxStatus)
// Only get pod IP from latest sandbox
if idx == 0 && podSandboxStatus.State == runtimeapi.PodSandboxState_SANDBOX_READY {
podIPs = m.determinePodSandboxIPs(namespace, name, podSandboxStatus)
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/kuberuntime/security_context_windows.go b/vendor/k8s.io/kubernetes/pkg/kubelet/kuberuntime/security_context_windows.go
index 343757c00..0ffcfe1af 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/kuberuntime/security_context_windows.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/kuberuntime/security_context_windows.go
@@ -24,6 +24,7 @@ import (
"k8s.io/klog/v2"
"k8s.io/kubernetes/pkg/kubelet/util/format"
"k8s.io/kubernetes/pkg/securitycontext"
+ "strings"
)
var (
@@ -36,6 +37,7 @@ var (
// According to the discussion of sig-windows, at present, we assume that ContainerAdministrator is the windows container root user,
// and then optimize this logic according to the best time.
// https://docs.google.com/document/d/1Tjxzjjuy4SQsFSUVXZbvqVb64hjNAG5CQX8bK7Yda9w
+// note: usernames on Windows are NOT case sensitive!
func verifyRunAsNonRoot(pod *v1.Pod, container *v1.Container, uid *int64, username string) error {
effectiveSc := securitycontext.DetermineEffectiveSecurityContext(pod, container)
// If the option is not set, or if running as root is allowed, return nil.
@@ -53,15 +55,17 @@ func verifyRunAsNonRoot(pod *v1.Pod, container *v1.Container, uid *int64, userna
if effectiveSc.RunAsGroup != nil {
klog.InfoS("Windows container does not support SecurityContext.RunAsGroup", "pod", klog.KObj(pod), "containerName", container.Name)
}
+ // Verify that if runAsUserName is set for the pod and/or container that it is not set to 'ContainerAdministrator'
if effectiveSc.WindowsOptions != nil {
if effectiveSc.WindowsOptions.RunAsUserName != nil {
- if *effectiveSc.WindowsOptions.RunAsUserName == windowsRootUserName {
- return fmt.Errorf("container's runAsUser (%s) which will be regarded as root identity and will break non-root policy (pod: %q, container: %s)", username, format.Pod(pod), container.Name)
+ if strings.EqualFold(*effectiveSc.WindowsOptions.RunAsUserName, windowsRootUserName) {
+ return fmt.Errorf("container's runAsUserName (%s) which will be regarded as root identity and will break non-root policy (pod: %q, container: %s)", *effectiveSc.WindowsOptions.RunAsUserName, format.Pod(pod), container.Name)
}
return nil
}
}
- if len(username) > 0 && username == windowsRootUserName {
+ // Verify that if runAsUserName is NOT set for the pod and/or container that the default user for the container image is not set to 'ContainerAdministrator'
+ if len(username) > 0 && strings.EqualFold(username, windowsRootUserName) {
return fmt.Errorf("container's runAsUser (%s) which will be regarded as root identity and will break non-root policy (pod: %q, container: %s)", username, format.Pod(pod), container.Name)
}
return nil
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/metrics/collectors/resource_metrics.go b/vendor/k8s.io/kubernetes/pkg/kubelet/metrics/collectors/resource_metrics.go
index 74d2ea7af..0f0662a2c 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/metrics/collectors/resource_metrics.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/metrics/collectors/resource_metrics.go
@@ -142,7 +142,7 @@ func (rc *resourceMetricsCollector) CollectWithStability(ch chan<- metrics.Metri
}
func (rc *resourceMetricsCollector) collectNodeCPUMetrics(ch chan<- metrics.Metric, s summary.NodeStats) {
- if s.CPU == nil {
+ if s.CPU == nil || s.CPU.UsageCoreNanoSeconds == nil {
return
}
@@ -151,7 +151,7 @@ func (rc *resourceMetricsCollector) collectNodeCPUMetrics(ch chan<- metrics.Metr
}
func (rc *resourceMetricsCollector) collectNodeMemoryMetrics(ch chan<- metrics.Metric, s summary.NodeStats) {
- if s.Memory == nil {
+ if s.Memory == nil || s.Memory.WorkingSetBytes == nil {
return
}
@@ -169,7 +169,7 @@ func (rc *resourceMetricsCollector) collectContainerStartTime(ch chan<- metrics.
}
func (rc *resourceMetricsCollector) collectContainerCPUMetrics(ch chan<- metrics.Metric, pod summary.PodStats, s summary.ContainerStats) {
- if s.CPU == nil {
+ if s.CPU == nil || s.CPU.UsageCoreNanoSeconds == nil {
return
}
@@ -179,7 +179,7 @@ func (rc *resourceMetricsCollector) collectContainerCPUMetrics(ch chan<- metrics
}
func (rc *resourceMetricsCollector) collectContainerMemoryMetrics(ch chan<- metrics.Metric, pod summary.PodStats, s summary.ContainerStats) {
- if s.Memory == nil {
+ if s.Memory == nil || s.Memory.WorkingSetBytes == nil {
return
}
@@ -189,7 +189,7 @@ func (rc *resourceMetricsCollector) collectContainerMemoryMetrics(ch chan<- metr
}
func (rc *resourceMetricsCollector) collectPodCPUMetrics(ch chan<- metrics.Metric, pod summary.PodStats) {
- if pod.CPU == nil {
+ if pod.CPU == nil || pod.CPU.UsageCoreNanoSeconds == nil {
return
}
@@ -199,7 +199,7 @@ func (rc *resourceMetricsCollector) collectPodCPUMetrics(ch chan<- metrics.Metri
}
func (rc *resourceMetricsCollector) collectPodMemoryMetrics(ch chan<- metrics.Metric, pod summary.PodStats) {
- if pod.Memory == nil {
+ if pod.Memory == nil || pod.Memory.WorkingSetBytes == nil {
return
}
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/pod/pod_manager.go b/vendor/k8s.io/kubernetes/pkg/kubelet/pod/pod_manager.go
index 42a559cdb..6d5553798 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/pod/pod_manager.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/pod/pod_manager.go
@@ -139,10 +139,6 @@ func (pm *basicManager) UpdatePod(pod *v1.Pod) {
pm.updatePodsInternal(pod)
}
-func isPodInTerminatedState(pod *v1.Pod) bool {
- return pod.Status.Phase == v1.PodFailed || pod.Status.Phase == v1.PodSucceeded
-}
-
// updateMetrics updates the metrics surfaced by the pod manager.
// oldPod or newPod may be nil to signify creation or deletion.
func updateMetrics(oldPod, newPod *v1.Pod) {
@@ -167,32 +163,6 @@ func updateMetrics(oldPod, newPod *v1.Pod) {
// lock.
func (pm *basicManager) updatePodsInternal(pods ...*v1.Pod) {
for _, pod := range pods {
- if pm.secretManager != nil {
- if isPodInTerminatedState(pod) {
- // Pods that are in terminated state and no longer running can be
- // ignored as they no longer require access to secrets.
- // It is especially important in watch-based manager, to avoid
- // unnecessary watches for terminated pods waiting for GC.
- pm.secretManager.UnregisterPod(pod)
- } else {
- // TODO: Consider detecting only status update and in such case do
- // not register pod, as it doesn't really matter.
- pm.secretManager.RegisterPod(pod)
- }
- }
- if pm.configMapManager != nil {
- if isPodInTerminatedState(pod) {
- // Pods that are in terminated state and no longer running can be
- // ignored as they no longer require access to configmaps.
- // It is especially important in watch-based manager, to avoid
- // unnecessary watches for terminated pods waiting for GC.
- pm.configMapManager.UnregisterPod(pod)
- } else {
- // TODO: Consider detecting only status update and in such case do
- // not register pod, as it doesn't really matter.
- pm.configMapManager.RegisterPod(pod)
- }
- }
podFullName := kubecontainer.GetPodFullName(pod)
resolvedPodUID := kubetypes.ResolvedPodUID(pod.UID)
pm.podByUID[resolvedPodUID] = pod
@@ -204,12 +174,6 @@ func (pm *basicManager) DeletePod(pod *v1.Pod) {
updateMetrics(pod, nil)
pm.lock.Lock()
defer pm.lock.Unlock()
- if pm.secretManager != nil {
- pm.secretManager.UnregisterPod(pod)
- }
- if pm.configMapManager != nil {
- pm.configMapManager.UnregisterPod(pod)
- }
podFullName := kubecontainer.GetPodFullName(pod)
delete(pm.podByUID, kubetypes.ResolvedPodUID(pod.UID))
delete(pm.podByFullName, podFullName)
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/pod_workers.go b/vendor/k8s.io/kubernetes/pkg/kubelet/pod_workers.go
index 185de3dcc..f2c5db9ec 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/pod_workers.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/pod_workers.go
@@ -215,7 +215,7 @@ type PodWorkers interface {
}
// the function to invoke to perform a sync (reconcile the kubelet state to the desired shape of the pod)
-type syncPodFnType func(ctx context.Context, updateType kubetypes.SyncPodType, pod *v1.Pod, podStatus *kubecontainer.PodStatus) error
+type syncPodFnType func(ctx context.Context, updateType kubetypes.SyncPodType, pod *v1.Pod, podStatus *kubecontainer.PodStatus) (bool, error)
// the function to invoke to terminate a pod (ensure no running processes are present)
type syncTerminatingPodFnType func(ctx context.Context, pod *v1.Pod, podStatus *kubecontainer.PodStatus, runningPod *kubecontainer.Pod, gracePeriod *int64, podStatusFn func(*v1.PodStatus)) error
@@ -389,6 +389,11 @@ type podWorkers struct {
syncTerminatingPodFn syncTerminatingPodFnType
syncTerminatedPodFn syncTerminatedPodFnType
+ // workerChannelFn is exposed for testing to allow unit tests to impose delays
+ // in channel communication. The function is invoked once each time a new worker
+ // goroutine starts.
+ workerChannelFn func(uid types.UID, in chan podWork) (out <-chan podWork)
+
// The EventRecorder to use
recorder record.EventRecorder
@@ -577,6 +582,7 @@ func (p *podWorkers) UpdatePod(options UpdatePodOptions) {
syncedAt: now,
startedTerminating: true,
finished: true,
+ fullname: kubecontainer.GetPodFullName(pod),
}
}
}
@@ -695,9 +701,8 @@ func (p *podWorkers) UpdatePod(options UpdatePodOptions) {
}
// start the pod worker goroutine if it doesn't exist
- var podUpdates chan podWork
- var exists bool
- if podUpdates, exists = p.podUpdates[uid]; !exists {
+ podUpdates, exists := p.podUpdates[uid]
+ if !exists {
// We need to have a buffer here, because checkForUpdates() method that
// puts an update into channel is called from the same goroutine where
// the channel is consumed. However, it is guaranteed that in such case
@@ -711,13 +716,21 @@ func (p *podWorkers) UpdatePod(options UpdatePodOptions) {
append(p.waitingToStartStaticPodsByFullname[status.fullname], uid)
}
+ // allow testing of delays in the pod update channel
+ var outCh <-chan podWork
+ if p.workerChannelFn != nil {
+ outCh = p.workerChannelFn(uid, podUpdates)
+ } else {
+ outCh = podUpdates
+ }
+
// Creating a new pod worker either means this is a new pod, or that the
// kubelet just restarted. In either case the kubelet is willing to believe
// the status of the pod for the first pod worker sync. See corresponding
// comment in syncPod.
go func() {
defer runtime.HandleCrash()
- p.managePodLoop(podUpdates)
+ p.managePodLoop(outCh)
}()
}
@@ -781,28 +794,31 @@ func calculateEffectiveGracePeriod(status *podSyncStatus, pod *v1.Pod, options *
}
// allowPodStart tries to start the pod and returns true if allowed, otherwise
-// it requeues the pod and returns false.
-func (p *podWorkers) allowPodStart(pod *v1.Pod) bool {
+// it requeues the pod and returns false. If the pod will never be able to start
+// because data is missing, or the pod was terminated before start, canEverStart
+// is false.
+func (p *podWorkers) allowPodStart(pod *v1.Pod) (canStart bool, canEverStart bool) {
if !kubetypes.IsStaticPod(pod) {
- // TBD: Do we want to allow non-static pods with the same full name?
+ // TODO: Do we want to allow non-static pods with the same full name?
// Note that it may disable the force deletion of pods.
- return true
+ return true, true
}
p.podLock.Lock()
defer p.podLock.Unlock()
status, ok := p.podSyncStatuses[pod.UID]
if !ok {
- klog.ErrorS(nil, "Failed to get a valid podSyncStatuses", "pod", klog.KObj(pod), "podUID", pod.UID)
- p.workQueue.Enqueue(pod.UID, wait.Jitter(p.backOffPeriod, workerBackOffPeriodJitterFactor))
- status.working = false
- return false
+ klog.ErrorS(nil, "Pod sync status does not exist, the worker should not be running", "pod", klog.KObj(pod), "podUID", pod.UID)
+ return false, false
+ }
+ if status.IsTerminationRequested() {
+ return false, false
}
if !p.allowStaticPodStart(status.fullname, pod.UID) {
p.workQueue.Enqueue(pod.UID, wait.Jitter(p.backOffPeriod, workerBackOffPeriodJitterFactor))
status.working = false
- return false
+ return false, true
}
- return true
+ return true, true
}
// allowStaticPodStart tries to start the static pod and returns true if
@@ -815,9 +831,12 @@ func (p *podWorkers) allowStaticPodStart(fullname string, uid types.UID) bool {
}
waitingPods := p.waitingToStartStaticPodsByFullname[fullname]
+ // TODO: This is O(N) with respect to the number of updates to static pods
+ // with overlapping full names, and ideally would be O(1).
for i, waitingUID := range waitingPods {
// has pod already terminated or been deleted?
- if _, ok := p.podSyncStatuses[waitingUID]; !ok {
+ status, ok := p.podSyncStatuses[waitingUID]
+ if !ok || status.IsTerminationRequested() || status.IsTerminated() {
continue
}
// another pod is next in line
@@ -843,8 +862,20 @@ func (p *podWorkers) managePodLoop(podUpdates <-chan podWork) {
var podStarted bool
for update := range podUpdates {
pod := update.Options.Pod
+
+ // Decide whether to start the pod. If the pod was terminated prior to the pod being allowed
+ // to start, we have to clean it up and then exit the pod worker loop.
if !podStarted {
- if !p.allowPodStart(pod) {
+ canStart, canEverStart := p.allowPodStart(pod)
+ if !canEverStart {
+ p.completeUnstartedTerminated(pod)
+ if start := update.Options.StartTime; !start.IsZero() {
+ metrics.PodWorkerDuration.WithLabelValues("terminated").Observe(metrics.SinceInSeconds(start))
+ }
+ klog.V(4).InfoS("Processing pod event done", "pod", klog.KObj(pod), "podUID", pod.UID, "updateType", update.WorkType)
+ return
+ }
+ if !canStart {
klog.V(4).InfoS("Pod cannot start yet", "pod", klog.KObj(pod), "podUID", pod.UID)
continue
}
@@ -852,6 +883,7 @@ func (p *podWorkers) managePodLoop(podUpdates <-chan podWork) {
}
klog.V(4).InfoS("Processing pod event", "pod", klog.KObj(pod), "podUID", pod.UID, "updateType", update.WorkType)
+ var isTerminal bool
err := func() error {
// The worker is responsible for ensuring the sync method sees the appropriate
// status updates on resyncs (the result of the last sync), transitions to
@@ -898,13 +930,14 @@ func (p *podWorkers) managePodLoop(podUpdates <-chan podWork) {
err = p.syncTerminatingPodFn(ctx, pod, status, update.Options.RunningPod, gracePeriod, podStatusFn)
default:
- err = p.syncPodFn(ctx, update.Options.UpdateType, pod, status)
+ isTerminal, err = p.syncPodFn(ctx, update.Options.UpdateType, pod, status)
}
lastSyncTime = time.Now()
return err
}()
+ var phaseTransition bool
switch {
case err == context.Canceled:
// when the context is cancelled we expect an update to already be queued
@@ -935,10 +968,17 @@ func (p *podWorkers) managePodLoop(podUpdates <-chan podWork) {
}
// otherwise we move to the terminating phase
p.completeTerminating(pod)
+ phaseTransition = true
+
+ case isTerminal:
+ // if syncPod indicated we are now terminal, set the appropriate pod status to move to terminating
+ klog.V(4).InfoS("Pod is terminal", "pod", klog.KObj(pod), "podUID", pod.UID, "updateType", update.WorkType)
+ p.completeSync(pod)
+ phaseTransition = true
}
- // queue a retry for errors if necessary, then put the next event in the channel if any
- p.completeWork(pod, err)
+ // queue a retry if necessary, then put the next event in the channel if any
+ p.completeWork(pod, phaseTransition, err)
if start := update.Options.StartTime; !start.IsZero() {
metrics.PodWorkerDuration.WithLabelValues(update.Options.UpdateType.String()).Observe(metrics.SinceInSeconds(start))
}
@@ -969,6 +1009,33 @@ func (p *podWorkers) acknowledgeTerminating(pod *v1.Pod) PodStatusFunc {
return nil
}
+// completeSync is invoked when syncPod completes successfully and indicates the pod is now terminal and should
+// be terminated. This happens when the natural pod lifecycle completes - any pod which is not RestartAlways
+// exits. Unnatural completions, such as evictions, API driven deletion or phase transition, are handled by
+// UpdatePod.
+func (p *podWorkers) completeSync(pod *v1.Pod) {
+ p.podLock.Lock()
+ defer p.podLock.Unlock()
+
+ klog.V(4).InfoS("Pod indicated lifecycle completed naturally and should now terminate", "pod", klog.KObj(pod), "podUID", pod.UID)
+
+ if status, ok := p.podSyncStatuses[pod.UID]; ok {
+ if status.terminatingAt.IsZero() {
+ status.terminatingAt = time.Now()
+ } else {
+ klog.V(4).InfoS("Pod worker attempted to set terminatingAt twice, likely programmer error", "pod", klog.KObj(pod), "podUID", pod.UID)
+ }
+ status.startedTerminating = true
+ }
+
+ p.lastUndeliveredWorkUpdate[pod.UID] = podWork{
+ WorkType: TerminatingPodWork,
+ Options: UpdatePodOptions{
+ Pod: pod,
+ },
+ }
+}
+
// completeTerminating is invoked when syncTerminatingPod completes successfully, which means
// no container is running, no container will be started in the future, and we are ready for
// cleanup. This updates the termination state which prevents future syncs and will ensure
@@ -1023,12 +1090,7 @@ func (p *podWorkers) completeTerminatingRuntimePod(pod *v1.Pod) {
}
}
- ch, ok := p.podUpdates[pod.UID]
- if ok {
- close(ch)
- }
- delete(p.podUpdates, pod.UID)
- delete(p.lastUndeliveredWorkUpdate, pod.UID)
+ p.cleanupPodUpdates(pod.UID)
}
// completeTerminated is invoked after syncTerminatedPod completes successfully and means we
@@ -1039,12 +1101,7 @@ func (p *podWorkers) completeTerminated(pod *v1.Pod) {
klog.V(4).InfoS("Pod is complete and the worker can now stop", "pod", klog.KObj(pod), "podUID", pod.UID)
- ch, ok := p.podUpdates[pod.UID]
- if ok {
- close(ch)
- }
- delete(p.podUpdates, pod.UID)
- delete(p.lastUndeliveredWorkUpdate, pod.UID)
+ p.cleanupPodUpdates(pod.UID)
if status, ok := p.podSyncStatuses[pod.UID]; ok {
if status.terminatingAt.IsZero() {
@@ -1062,11 +1119,40 @@ func (p *podWorkers) completeTerminated(pod *v1.Pod) {
}
}
+// completeUnstartedTerminated is invoked if a pod that has never been started receives a termination
+// signal before it can be started.
+func (p *podWorkers) completeUnstartedTerminated(pod *v1.Pod) {
+ p.podLock.Lock()
+ defer p.podLock.Unlock()
+
+ klog.V(4).InfoS("Pod never started and the worker can now stop", "pod", klog.KObj(pod), "podUID", pod.UID)
+
+ p.cleanupPodUpdates(pod.UID)
+
+ if status, ok := p.podSyncStatuses[pod.UID]; ok {
+ if status.terminatingAt.IsZero() {
+ klog.V(4).InfoS("Pod worker is complete but did not have terminatingAt set, likely programmer error", "pod", klog.KObj(pod), "podUID", pod.UID)
+ }
+ if !status.terminatedAt.IsZero() {
+ klog.V(4).InfoS("Pod worker is complete and had terminatedAt set, likely programmer error", "pod", klog.KObj(pod), "podUID", pod.UID)
+ }
+ status.finished = true
+ status.working = false
+ status.terminatedAt = time.Now()
+
+ if p.startedStaticPodsByFullname[status.fullname] == pod.UID {
+ delete(p.startedStaticPodsByFullname, status.fullname)
+ }
+ }
+}
+
// completeWork requeues on error or the next sync interval and then immediately executes any pending
// work.
-func (p *podWorkers) completeWork(pod *v1.Pod, syncErr error) {
+func (p *podWorkers) completeWork(pod *v1.Pod, phaseTransition bool, syncErr error) {
// Requeue the last update if the last sync returned error.
switch {
+ case phaseTransition:
+ p.workQueue.Enqueue(pod.UID, 0)
case syncErr == nil:
// No error; requeue at the regular resync interval.
p.workQueue.Enqueue(pod.UID, wait.Jitter(p.resyncInterval, workerResyncIntervalJitterFactor))
@@ -1146,10 +1232,10 @@ func (p *podWorkers) SyncKnownPods(desiredPods []*v1.Pod) map[types.UID]PodWorke
return workers
}
-// removeTerminatedWorker cleans up and removes the worker status for a worker that
-// has reached a terminal state of "finished" - has successfully exited
-// syncTerminatedPod. This "forgets" a pod by UID and allows another pod to be recreated
-// with the same UID.
+// removeTerminatedWorker cleans up and removes the worker status for a worker
+// that has reached a terminal state of "finished" - has successfully exited
+// syncTerminatedPod. This "forgets" a pod by UID and allows another pod to be
+// recreated with the same UID.
func (p *podWorkers) removeTerminatedWorker(uid types.UID) {
status, ok := p.podSyncStatuses[uid]
if !ok {
@@ -1158,11 +1244,6 @@ func (p *podWorkers) removeTerminatedWorker(uid types.UID) {
return
}
- if startedUID, started := p.startedStaticPodsByFullname[status.fullname]; started && startedUID != uid {
- klog.V(4).InfoS("Pod cannot start yet but is no longer known to the kubelet, finish it", "podUID", uid)
- status.finished = true
- }
-
if !status.finished {
klog.V(4).InfoS("Pod worker has been requested for removal but is still not fully terminated", "podUID", uid)
return
@@ -1174,8 +1255,7 @@ func (p *podWorkers) removeTerminatedWorker(uid types.UID) {
klog.V(4).InfoS("Pod has been terminated and is no longer known to the kubelet, remove all history", "podUID", uid)
}
delete(p.podSyncStatuses, uid)
- delete(p.podUpdates, uid)
- delete(p.lastUndeliveredWorkUpdate, uid)
+ p.cleanupPodUpdates(uid)
if p.startedStaticPodsByFullname[status.fullname] == uid {
delete(p.startedStaticPodsByFullname, status.fullname)
@@ -1226,3 +1306,15 @@ func killPodNow(podWorkers PodWorkers, recorder record.EventRecorder) eviction.K
}
}
}
+
+// cleanupPodUpdates closes the podUpdates channel and removes it from
+// podUpdates map so that the corresponding pod worker can stop. It also
+// removes any undelivered work. This method must be called holding the
+// pod lock.
+func (p *podWorkers) cleanupPodUpdates(uid types.UID) {
+ if ch, ok := p.podUpdates[uid]; ok {
+ close(ch)
+ }
+ delete(p.podUpdates, uid)
+ delete(p.lastUndeliveredWorkUpdate, uid)
+}
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/prober/prober_manager.go b/vendor/k8s.io/kubernetes/pkg/kubelet/prober/prober_manager.go
index 83532a313..9a6dc26a7 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/prober/prober_manager.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/prober/prober_manager.go
@@ -58,6 +58,9 @@ type Manager interface {
// pod created.
AddPod(pod *v1.Pod)
+ // StopLivenessAndStartup handles stopping liveness and startup probes during termination.
+ StopLivenessAndStartup(pod *v1.Pod)
+
// RemovePod handles cleaning up the removed pod state, including terminating probe workers and
// deleting cached results.
RemovePod(pod *v1.Pod)
@@ -161,7 +164,7 @@ func (m *manager) AddPod(pod *v1.Pod) {
if c.StartupProbe != nil {
key.probeType = startup
if _, ok := m.workers[key]; ok {
- klog.ErrorS(nil, "Startup probe already exists for container",
+ klog.V(8).ErrorS(nil, "Startup probe already exists for container",
"pod", klog.KObj(pod), "containerName", c.Name)
return
}
@@ -173,7 +176,7 @@ func (m *manager) AddPod(pod *v1.Pod) {
if c.ReadinessProbe != nil {
key.probeType = readiness
if _, ok := m.workers[key]; ok {
- klog.ErrorS(nil, "Readiness probe already exists for container",
+ klog.V(8).ErrorS(nil, "Readiness probe already exists for container",
"pod", klog.KObj(pod), "containerName", c.Name)
return
}
@@ -185,7 +188,7 @@ func (m *manager) AddPod(pod *v1.Pod) {
if c.LivenessProbe != nil {
key.probeType = liveness
if _, ok := m.workers[key]; ok {
- klog.ErrorS(nil, "Liveness probe already exists for container",
+ klog.V(8).ErrorS(nil, "Liveness probe already exists for container",
"pod", klog.KObj(pod), "containerName", c.Name)
return
}
@@ -196,6 +199,22 @@ func (m *manager) AddPod(pod *v1.Pod) {
}
}
+func (m *manager) StopLivenessAndStartup(pod *v1.Pod) {
+ m.workerLock.RLock()
+ defer m.workerLock.RUnlock()
+
+ key := probeKey{podUID: pod.UID}
+ for _, c := range pod.Spec.Containers {
+ key.containerName = c.Name
+ for _, probeType := range [...]probeType{liveness, startup} {
+ key.probeType = probeType
+ if worker, ok := m.workers[key]; ok {
+ worker.stop()
+ }
+ }
+ }
+}
+
func (m *manager) RemovePod(pod *v1.Pod) {
m.workerLock.RLock()
defer m.workerLock.RUnlock()
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/runonce.go b/vendor/k8s.io/kubernetes/pkg/kubelet/runonce.go
index ce7985021..4544cc992 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/runonce.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/runonce.go
@@ -112,9 +112,10 @@ func (kl *Kubelet) runOnce(pods []*v1.Pod, retryDelay time.Duration) (results []
// runPod runs a single pod and wait until all containers are running.
func (kl *Kubelet) runPod(pod *v1.Pod, retryDelay time.Duration) error {
+ var isTerminal bool
delay := retryDelay
retry := 0
- for {
+ for !isTerminal {
status, err := kl.containerRuntime.GetPodStatus(pod.UID, pod.Name, pod.Namespace)
if err != nil {
return fmt.Errorf("unable to get status for pod %q: %v", format.Pod(pod), err)
@@ -126,7 +127,7 @@ func (kl *Kubelet) runPod(pod *v1.Pod, retryDelay time.Duration) error {
}
klog.InfoS("Pod's containers not running: syncing", "pod", klog.KObj(pod))
- if err = kl.syncPod(context.Background(), kubetypes.SyncPodUpdate, pod, status); err != nil {
+ if isTerminal, err = kl.syncPod(context.Background(), kubetypes.SyncPodUpdate, pod, status); err != nil {
return fmt.Errorf("error syncing pod %q: %v", format.Pod(pod), err)
}
if retry >= runOnceMaxRetries {
@@ -138,6 +139,7 @@ func (kl *Kubelet) runPod(pod *v1.Pod, retryDelay time.Duration) error {
retry++
delay *= runOnceRetryDelayBackoff
}
+ return nil
}
// isPodRunning returns true if all containers of a manifest are running.
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/stats/cri_stats_provider.go b/vendor/k8s.io/kubernetes/pkg/kubelet/stats/cri_stats_provider.go
index e2f6ed41a..9c959c3c2 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/stats/cri_stats_provider.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/stats/cri_stats_provider.go
@@ -826,7 +826,24 @@ func getCRICadvisorStats(infos map[string]cadvisorapiv2.ContainerInfo) map[strin
if !isPodManagedContainer(&info) {
continue
}
- stats[path.Base(key)] = info
+ stats[extractIDFromCgroupPath(key)] = info
}
return stats
}
+
+func extractIDFromCgroupPath(cgroupPath string) string {
+ // case0 == cgroupfs: "/kubepods/burstable/pod2fc932ce-fdcc-454b-97bd-aadfdeb4c340/9be25294016e2dc0340dd605ce1f57b492039b267a6a618a7ad2a7a58a740f32"
+ id := path.Base(cgroupPath)
+
+ // case1 == systemd: "/kubepods.slice/kubepods-burstable.slice/kubepods-burstable-pod2fc932ce_fdcc_454b_97bd_aadfdeb4c340.slice/cri-containerd-aaefb9d8feed2d453b543f6d928cede7a4dbefa6a0ae7c9b990dd234c56e93b9.scope"
+ // trim anything before the final '-' and suffix .scope
+ systemdSuffix := ".scope"
+ if strings.HasSuffix(id, systemdSuffix) {
+ id = strings.TrimSuffix(id, systemdSuffix)
+ components := strings.Split(id, "-")
+ if len(components) > 1 {
+ id = components[len(components)-1]
+ }
+ }
+ return id
+}
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/status/generate.go b/vendor/k8s.io/kubernetes/pkg/kubelet/status/generate.go
index ef87faa99..024f1cf4b 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/status/generate.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/status/generate.go
@@ -20,7 +20,7 @@ import (
"fmt"
"strings"
- "k8s.io/api/core/v1"
+ v1 "k8s.io/api/core/v1"
podutil "k8s.io/kubernetes/pkg/api/v1/pod"
)
@@ -29,6 +29,8 @@ const (
UnknownContainerStatuses = "UnknownContainerStatuses"
// PodCompleted says that all related containers have succeeded.
PodCompleted = "PodCompleted"
+ // PodFailed says that the pod has failed and as such the containers have failed.
+ PodFailed = "PodFailed"
// ContainersNotReady says that one or more containers are not ready.
ContainersNotReady = "ContainersNotReady"
// ContainersNotInitialized says that one or more init containers have not succeeded.
@@ -62,11 +64,12 @@ func GenerateContainersReadyCondition(spec *v1.PodSpec, containerStatuses []v1.C
// If all containers are known and succeeded, just return PodCompleted.
if podPhase == v1.PodSucceeded && len(unknownContainers) == 0 {
- return v1.PodCondition{
- Type: v1.ContainersReady,
- Status: v1.ConditionFalse,
- Reason: PodCompleted,
- }
+ return generateContainersReadyConditionForTerminalPhase(podPhase)
+ }
+
+ // If the pod phase is failed, explicitly set the ready condition to false for containers since they may be in progress of terminating.
+ if podPhase == v1.PodFailed {
+ return generateContainersReadyConditionForTerminalPhase(podPhase)
}
// Generate message for containers in unknown condition.
@@ -191,3 +194,33 @@ func GeneratePodInitializedCondition(spec *v1.PodSpec, containerStatuses []v1.Co
Status: v1.ConditionTrue,
}
}
+
+func generateContainersReadyConditionForTerminalPhase(podPhase v1.PodPhase) v1.PodCondition {
+ condition := v1.PodCondition{
+ Type: v1.ContainersReady,
+ Status: v1.ConditionFalse,
+ }
+
+ if podPhase == v1.PodFailed {
+ condition.Reason = PodFailed
+ } else if podPhase == v1.PodSucceeded {
+ condition.Reason = PodCompleted
+ }
+
+ return condition
+}
+
+func generatePodReadyConditionForTerminalPhase(podPhase v1.PodPhase) v1.PodCondition {
+ condition := v1.PodCondition{
+ Type: v1.PodReady,
+ Status: v1.ConditionFalse,
+ }
+
+ if podPhase == v1.PodFailed {
+ condition.Reason = PodFailed
+ } else if podPhase == v1.PodSucceeded {
+ condition.Reason = PodCompleted
+ }
+
+ return condition
+}
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/status/status_manager.go b/vendor/k8s.io/kubernetes/pkg/kubelet/status/status_manager.go
index 88aac7fbc..8ebd5c79b 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/status/status_manager.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/status/status_manager.go
@@ -82,8 +82,10 @@ type PodStatusProvider interface {
// PodDeletionSafetyProvider provides guarantees that a pod can be safely deleted.
type PodDeletionSafetyProvider interface {
- // A function which returns true if the pod can safely be deleted
+ // PodResourcesAreReclaimed returns true if the pod can safely be deleted.
PodResourcesAreReclaimed(pod *v1.Pod, status v1.PodStatus) bool
+ // PodCouldHaveRunningContainers returns true if the pod could have running containers.
+ PodCouldHaveRunningContainers(pod *v1.Pod) bool
}
// Manager is the Source of truth for kubelet pod status, and should be kept up-to-date with
@@ -321,6 +323,14 @@ func findContainerStatus(status *v1.PodStatus, containerID string) (containerSta
}
+// TerminatePod ensures that the status of containers is properly defaulted at the end of the pod
+// lifecycle. As the Kubelet must reconcile with the container runtime to observe container status
+// there is always the possibility we are unable to retrieve one or more container statuses due to
+// garbage collection, admin action, or loss of temporary data on a restart. This method ensures
+// that any absent container status is treated as a failure so that we do not incorrectly describe
+// the pod as successful. If we have not yet initialized the pod in the presence of init containers,
+// the init container failure status is sufficient to describe the pod as failing, and we do not need
+// to override waiting containers (unless there is evidence the pod previously started those containers).
func (m *manager) TerminatePod(pod *v1.Pod) {
m.podStatusesLock.Lock()
defer m.podStatusesLock.Unlock()
@@ -332,19 +342,26 @@ func (m *manager) TerminatePod(pod *v1.Pod) {
oldStatus = &cachedStatus.status
}
status := *oldStatus.DeepCopy()
- for i := range status.ContainerStatuses {
- if status.ContainerStatuses[i].State.Terminated != nil {
- continue
- }
- status.ContainerStatuses[i].State = v1.ContainerState{
- Terminated: &v1.ContainerStateTerminated{
- Reason: "ContainerStatusUnknown",
- Message: "The container could not be located when the pod was terminated",
- ExitCode: 137,
- },
+
+ // once a pod has initialized, any missing status is treated as a failure
+ if hasPodInitialized(pod) {
+ for i := range status.ContainerStatuses {
+ if status.ContainerStatuses[i].State.Terminated != nil {
+ continue
+ }
+ status.ContainerStatuses[i].State = v1.ContainerState{
+ Terminated: &v1.ContainerStateTerminated{
+ Reason: "ContainerStatusUnknown",
+ Message: "The container could not be located when the pod was terminated",
+ ExitCode: 137,
+ },
+ }
}
}
- for i := range status.InitContainerStatuses {
+
+ // all but the final suffix of init containers which have no evidence of a container start are
+ // marked as failed containers
+ for i := range initializedContainers(status.InitContainerStatuses) {
if status.InitContainerStatuses[i].State.Terminated != nil {
continue
}
@@ -361,6 +378,49 @@ func (m *manager) TerminatePod(pod *v1.Pod) {
m.updateStatusInternal(pod, status, true)
}
+// hasPodInitialized returns true if the pod has no evidence of ever starting a regular container, which
+// implies those containers should not be transitioned to terminated status.
+func hasPodInitialized(pod *v1.Pod) bool {
+ // a pod without init containers is always initialized
+ if len(pod.Spec.InitContainers) == 0 {
+ return true
+ }
+ // if any container has ever moved out of waiting state, the pod has initialized
+ for _, status := range pod.Status.ContainerStatuses {
+ if status.LastTerminationState.Terminated != nil || status.State.Waiting == nil {
+ return true
+ }
+ }
+ // if the last init container has ever completed with a zero exit code, the pod is initialized
+ if l := len(pod.Status.InitContainerStatuses); l > 0 {
+ container := pod.Status.InitContainerStatuses[l-1]
+ if state := container.LastTerminationState; state.Terminated != nil && state.Terminated.ExitCode == 0 {
+ return true
+ }
+ if state := container.State; state.Terminated != nil && state.Terminated.ExitCode == 0 {
+ return true
+ }
+ }
+ // otherwise the pod has no record of being initialized
+ return false
+}
+
+// initializedContainers returns all status except for suffix of containers that are in Waiting
+// state, which is the set of containers that have attempted to start at least once. If all containers
+// are Watiing, the first container is always returned.
+func initializedContainers(containers []v1.ContainerStatus) []v1.ContainerStatus {
+ for i := len(containers) - 1; i >= 0; i-- {
+ if containers[i].State.Waiting == nil || containers[i].LastTerminationState.Terminated != nil {
+ return containers[0 : i+1]
+ }
+ }
+ // always return at least one container
+ if len(containers) > 0 {
+ return containers[0:1]
+ }
+ return nil
+}
+
// checkContainerStateTransition ensures that no container is trying to transition
// from a terminated to non-terminated state, which is illegal and indicates a
// logical error in the kubelet.
@@ -603,8 +663,9 @@ func (m *manager) syncPod(uid types.UID, status versionedPodStatus) {
return
}
- oldStatus := pod.Status.DeepCopy()
- newPod, patchBytes, unchanged, err := statusutil.PatchPodStatus(m.kubeClient, pod.Namespace, pod.Name, pod.UID, *oldStatus, mergePodStatus(*oldStatus, status.status))
+ mergedStatus := mergePodStatus(pod.Status, status.status, m.podDeletionSafety.PodCouldHaveRunningContainers(pod))
+
+ newPod, patchBytes, unchanged, err := statusutil.PatchPodStatus(m.kubeClient, pod.Namespace, pod.Name, pod.UID, pod.Status, mergedStatus)
klog.V(3).InfoS("Patch status for pod", "pod", klog.KObj(pod), "patch", string(patchBytes))
if err != nil {
@@ -614,7 +675,7 @@ func (m *manager) syncPod(uid types.UID, status versionedPodStatus) {
if unchanged {
klog.V(3).InfoS("Status for pod is up-to-date", "pod", klog.KObj(pod), "statusVersion", status.version)
} else {
- klog.V(3).InfoS("Status for pod updated successfully", "pod", klog.KObj(pod), "statusVersion", status.version, "status", status.status)
+ klog.V(3).InfoS("Status for pod updated successfully", "pod", klog.KObj(pod), "statusVersion", status.version, "status", mergedStatus)
pod = newPod
}
@@ -746,22 +807,55 @@ func normalizeStatus(pod *v1.Pod, status *v1.PodStatus) *v1.PodStatus {
return status
}
-// mergePodStatus merges oldPodStatus and newPodStatus where pod conditions
-// not owned by kubelet is preserved from oldPodStatus
-func mergePodStatus(oldPodStatus, newPodStatus v1.PodStatus) v1.PodStatus {
- podConditions := []v1.PodCondition{}
+// mergePodStatus merges oldPodStatus and newPodStatus to preserve where pod conditions
+// not owned by kubelet and to ensure terminal phase transition only happens after all
+// running containers have terminated. This method does not modify the old status.
+func mergePodStatus(oldPodStatus, newPodStatus v1.PodStatus, couldHaveRunningContainers bool) v1.PodStatus {
+ podConditions := make([]v1.PodCondition, 0, len(oldPodStatus.Conditions)+len(newPodStatus.Conditions))
+
for _, c := range oldPodStatus.Conditions {
if !kubetypes.PodConditionByKubelet(c.Type) {
podConditions = append(podConditions, c)
}
}
-
for _, c := range newPodStatus.Conditions {
if kubetypes.PodConditionByKubelet(c.Type) {
podConditions = append(podConditions, c)
}
}
newPodStatus.Conditions = podConditions
+
+ // Delay transitioning a pod to a terminal status unless the pod is actually terminal.
+ // The Kubelet should never transition a pod to terminal status that could have running
+ // containers and thus actively be leveraging exclusive resources. Note that resources
+ // like volumes are reconciled by a subsystem in the Kubelet and will converge if a new
+ // pod reuses an exclusive resource (unmount -> free -> mount), which means we do not
+ // need wait for those resources to be detached by the Kubelet. In general, resources
+ // the Kubelet exclusively owns must be released prior to a pod being reported terminal,
+ // while resources that have participanting components above the API use the pod's
+ // transition to a terminal phase (or full deletion) to release those resources.
+ if !podutil.IsPodPhaseTerminal(oldPodStatus.Phase) && podutil.IsPodPhaseTerminal(newPodStatus.Phase) {
+ if couldHaveRunningContainers {
+ newPodStatus.Phase = oldPodStatus.Phase
+ newPodStatus.Reason = oldPodStatus.Reason
+ newPodStatus.Message = oldPodStatus.Message
+ }
+ }
+
+ // If the new phase is terminal, explicitly set the ready condition to false for v1.PodReady and v1.ContainersReady.
+ // It may take some time for kubelet to reconcile the ready condition, so explicitly set ready conditions to false if the phase is terminal.
+ // This is done to ensure kubelet does not report a status update with terminal pod phase and ready=true.
+ // See https://issues.k8s.io/108594 for more details.
+ if podutil.IsPodPhaseTerminal(newPodStatus.Phase) {
+ if podutil.IsPodReadyConditionTrue(newPodStatus) || podutil.IsContainersReadyConditionTrue(newPodStatus) {
+ containersReadyCondition := generateContainersReadyConditionForTerminalPhase(newPodStatus.Phase)
+ podutil.UpdatePodCondition(&newPodStatus, &containersReadyCondition)
+
+ podReadyCondition := generatePodReadyConditionForTerminalPhase(newPodStatus.Phase)
+ podutil.UpdatePodCondition(&newPodStatus, &podReadyCondition)
+ }
+ }
+
return newPodStatus
}
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/token/token_manager.go b/vendor/k8s.io/kubernetes/pkg/kubelet/token/token_manager.go
index 7901c958f..f2f3a8053 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/token/token_manager.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/token/token_manager.go
@@ -74,6 +74,9 @@ func NewManager(c clientset.Interface) *Manager {
}
return tokenRequest, err
},
+ deleteToken: func(podUID types.UID) {
+ c.CoreV1().ServiceAccounts("default").Delete(context.TODO(), string(podUID), metav1.DeleteOptions{})
+ },
cache: make(map[string]*authenticationv1.TokenRequest),
clock: clock.RealClock{},
}
@@ -89,8 +92,9 @@ type Manager struct {
cache map[string]*authenticationv1.TokenRequest
// mocked for testing
- getToken func(name, namespace string, tr *authenticationv1.TokenRequest) (*authenticationv1.TokenRequest, error)
- clock clock.Clock
+ getToken func(name, namespace string, tr *authenticationv1.TokenRequest) (*authenticationv1.TokenRequest, error)
+ deleteToken func(podUID types.UID)
+ clock clock.Clock
}
// GetServiceAccountToken gets a service account token for a pod from cache or
@@ -137,6 +141,7 @@ func (m *Manager) DeleteServiceAccountToken(podUID types.UID) {
delete(m.cache, k)
}
}
+ m.deleteToken(podUID)
}
func (m *Manager) cleanup() {
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/util/manager/cache_based_manager.go b/vendor/k8s.io/kubernetes/pkg/kubelet/util/manager/cache_based_manager.go
index 591a57ab8..07ae0b39f 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/util/manager/cache_based_manager.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/util/manager/cache_based_manager.go
@@ -29,6 +29,7 @@ import (
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime"
+ "k8s.io/apimachinery/pkg/types"
"k8s.io/apimachinery/pkg/util/clock"
"k8s.io/apimachinery/pkg/util/sets"
)
@@ -42,6 +43,7 @@ type GetObjectFunc func(string, string, metav1.GetOptions) (runtime.Object, erro
type objectKey struct {
namespace string
name string
+ uid types.UID
}
// objectStoreItems is a single item stored in objectStore.
@@ -226,7 +228,7 @@ func (c *cacheBasedManager) RegisterPod(pod *v1.Pod) {
c.objectStore.AddReference(pod.Namespace, name)
}
var prev *v1.Pod
- key := objectKey{namespace: pod.Namespace, name: pod.Name}
+ key := objectKey{namespace: pod.Namespace, name: pod.Name, uid: pod.UID}
prev = c.registeredPods[key]
c.registeredPods[key] = pod
if prev != nil {
@@ -243,7 +245,7 @@ func (c *cacheBasedManager) RegisterPod(pod *v1.Pod) {
func (c *cacheBasedManager) UnregisterPod(pod *v1.Pod) {
var prev *v1.Pod
- key := objectKey{namespace: pod.Namespace, name: pod.Name}
+ key := objectKey{namespace: pod.Namespace, name: pod.Name, uid: pod.UID}
c.lock.Lock()
defer c.lock.Unlock()
prev = c.registeredPods[key]
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/volume_host.go b/vendor/k8s.io/kubernetes/pkg/kubelet/volume_host.go
index aac566eeb..c6f83437e 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/volume_host.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/volume_host.go
@@ -265,6 +265,20 @@ func (kvh *kubeletVolumeHost) GetNodeLabels() (map[string]string, error) {
return node.Labels, nil
}
+func (kvh *kubeletVolumeHost) GetAttachedVolumesFromNodeStatus() (map[v1.UniqueVolumeName]string, error) {
+ node, err := kvh.kubelet.GetNode()
+ if err != nil {
+ return nil, fmt.Errorf("error retrieving node: %v", err)
+ }
+ attachedVolumes := node.Status.VolumesAttached
+ result := map[v1.UniqueVolumeName]string{}
+ for i := range attachedVolumes {
+ attachedVolume := attachedVolumes[i]
+ result[attachedVolume.Name] = attachedVolume.DevicePath
+ }
+ return result, nil
+}
+
func (kvh *kubeletVolumeHost) GetNodeName() types.NodeName {
return kvh.kubelet.nodeName
}
diff --git a/vendor/k8s.io/kubernetes/pkg/kubelet/volumemanager/reconciler/reconciler.go b/vendor/k8s.io/kubernetes/pkg/kubelet/volumemanager/reconciler/reconciler.go
index 6e95e72a2..c228681de 100644
--- a/vendor/k8s.io/kubernetes/pkg/kubelet/volumemanager/reconciler/reconciler.go
+++ b/vendor/k8s.io/kubernetes/pkg/kubelet/volumemanager/reconciler/reconciler.go
@@ -185,9 +185,7 @@ func (rc *reconciler) unmountVolumes() {
klog.V(5).InfoS(mountedVolume.GenerateMsgDetailed("Starting operationExecutor.UnmountVolume", ""))
err := rc.operationExecutor.UnmountVolume(
mountedVolume.MountedVolume, rc.actualStateOfWorld, rc.kubeletPodsDir)
- if err != nil &&
- !nestedpendingoperations.IsAlreadyExists(err) &&
- !exponentialbackoff.IsExponentialBackoff(err) {
+ if err != nil && !isExpectedError(err) {
// Ignore nestedpendingoperations.IsAlreadyExists and exponentialbackoff.IsExponentialBackoff errors, they are expected.
// Log all other errors.
klog.ErrorS(err, mountedVolume.GenerateErrorDetailed(fmt.Sprintf("operationExecutor.UnmountVolume failed (controllerAttachDetachEnabled %v)", rc.controllerAttachDetachEnabled), err).Error())
@@ -206,6 +204,11 @@ func (rc *reconciler) mountAttachVolumes() {
volumeToMount.DevicePath = devicePath
if cache.IsVolumeNotAttachedError(err) {
if rc.controllerAttachDetachEnabled || !volumeToMount.PluginIsAttachable {
+ //// lets not spin a goroutine and unnecessarily trigger exponential backoff if this happens
+ if volumeToMount.PluginIsAttachable && !volumeToMount.ReportedInUse {
+ klog.V(5).InfoS(volumeToMount.GenerateMsgDetailed("operationExecutor.VerifyControllerAttachedVolume failed", " volume not marked in-use"), "pod", klog.KObj(volumeToMount.Pod))
+ continue
+ }
// Volume is not attached (or doesn't implement attacher), kubelet attach is disabled, wait
// for controller to finish attaching volume.
klog.V(5).InfoS(volumeToMount.GenerateMsgDetailed("Starting operationExecutor.VerifyControllerAttachedVolume", ""))
@@ -213,9 +216,7 @@ func (rc *reconciler) mountAttachVolumes() {
volumeToMount.VolumeToMount,
rc.nodeName,
rc.actualStateOfWorld)
- if err != nil &&
- !nestedpendingoperations.IsAlreadyExists(err) &&
- !exponentialbackoff.IsExponentialBackoff(err) {
+ if err != nil && !isExpectedError(err) {
// Ignore nestedpendingoperations.IsAlreadyExists and exponentialbackoff.IsExponentialBackoff errors, they are expected.
// Log all other errors.
klog.ErrorS(err, volumeToMount.GenerateErrorDetailed(fmt.Sprintf("operationExecutor.VerifyControllerAttachedVolume failed (controllerAttachDetachEnabled %v)", rc.controllerAttachDetachEnabled), err).Error())
@@ -233,9 +234,7 @@ func (rc *reconciler) mountAttachVolumes() {
}
klog.V(5).InfoS(volumeToAttach.GenerateMsgDetailed("Starting operationExecutor.AttachVolume", ""))
err := rc.operationExecutor.AttachVolume(volumeToAttach, rc.actualStateOfWorld)
- if err != nil &&
- !nestedpendingoperations.IsAlreadyExists(err) &&
- !exponentialbackoff.IsExponentialBackoff(err) {
+ if err != nil && !isExpectedError(err) {
// Ignore nestedpendingoperations.IsAlreadyExists and exponentialbackoff.IsExponentialBackoff errors, they are expected.
// Log all other errors.
klog.ErrorS(err, volumeToMount.GenerateErrorDetailed(fmt.Sprintf("operationExecutor.AttachVolume failed (controllerAttachDetachEnabled %v)", rc.controllerAttachDetachEnabled), err).Error())
@@ -257,9 +256,7 @@ func (rc *reconciler) mountAttachVolumes() {
volumeToMount.VolumeToMount,
rc.actualStateOfWorld,
isRemount)
- if err != nil &&
- !nestedpendingoperations.IsAlreadyExists(err) &&
- !exponentialbackoff.IsExponentialBackoff(err) {
+ if err != nil && !isExpectedError(err) {
// Ignore nestedpendingoperations.IsAlreadyExists and exponentialbackoff.IsExponentialBackoff errors, they are expected.
// Log all other errors.
klog.ErrorS(err, volumeToMount.GenerateErrorDetailed(fmt.Sprintf("operationExecutor.MountVolume failed (controllerAttachDetachEnabled %v)", rc.controllerAttachDetachEnabled), err).Error())
@@ -277,9 +274,7 @@ func (rc *reconciler) mountAttachVolumes() {
err := rc.operationExecutor.ExpandInUseVolume(
volumeToMount.VolumeToMount,
rc.actualStateOfWorld)
- if err != nil &&
- !nestedpendingoperations.IsAlreadyExists(err) &&
- !exponentialbackoff.IsExponentialBackoff(err) {
+ if err != nil && !isExpectedError(err) {
// Ignore nestedpendingoperations.IsAlreadyExists and exponentialbackoff.IsExponentialBackoff errors, they are expected.
// Log all other errors.
klog.ErrorS(err, volumeToMount.GenerateErrorDetailed("operationExecutor.ExpandInUseVolume failed", err).Error())
@@ -301,9 +296,7 @@ func (rc *reconciler) unmountDetachDevices() {
klog.V(5).InfoS(attachedVolume.GenerateMsgDetailed("Starting operationExecutor.UnmountDevice", ""))
err := rc.operationExecutor.UnmountDevice(
attachedVolume.AttachedVolume, rc.actualStateOfWorld, rc.hostutil)
- if err != nil &&
- !nestedpendingoperations.IsAlreadyExists(err) &&
- !exponentialbackoff.IsExponentialBackoff(err) {
+ if err != nil && !isExpectedError(err) {
// Ignore nestedpendingoperations.IsAlreadyExists and exponentialbackoff.IsExponentialBackoff errors, they are expected.
// Log all other errors.
klog.ErrorS(err, attachedVolume.GenerateErrorDetailed(fmt.Sprintf("operationExecutor.UnmountDevice failed (controllerAttachDetachEnabled %v)", rc.controllerAttachDetachEnabled), err).Error())
@@ -322,9 +315,7 @@ func (rc *reconciler) unmountDetachDevices() {
klog.V(5).InfoS(attachedVolume.GenerateMsgDetailed("Starting operationExecutor.DetachVolume", ""))
err := rc.operationExecutor.DetachVolume(
attachedVolume.AttachedVolume, false /* verifySafeToDetach */, rc.actualStateOfWorld)
- if err != nil &&
- !nestedpendingoperations.IsAlreadyExists(err) &&
- !exponentialbackoff.IsExponentialBackoff(err) {
+ if err != nil && !isExpectedError(err) {
// Ignore nestedpendingoperations.IsAlreadyExists && exponentialbackoff.IsExponentialBackoff errors, they are expected.
// Log all other errors.
klog.ErrorS(err, attachedVolume.GenerateErrorDetailed(fmt.Sprintf("operationExecutor.DetachVolume failed (controllerAttachDetachEnabled %v)", rc.controllerAttachDetachEnabled), err).Error())
@@ -729,3 +720,8 @@ func getVolumesFromPodDir(podDir string) ([]podVolume, error) {
klog.V(4).InfoS("Get volumes from pod directory", "path", podDir, "volumes", volumes)
return volumes, nil
}
+
+// ignore nestedpendingoperations.IsAlreadyExists and exponentialbackoff.IsExponentialBackoff errors, they are expected.
+func isExpectedError(err error) bool {
+ return nestedpendingoperations.IsAlreadyExists(err) || exponentialbackoff.IsExponentialBackoff(err) || operationexecutor.IsMountFailedPreconditionError(err)
+}
diff --git a/vendor/k8s.io/kubernetes/pkg/printers/internalversion/printers.go b/vendor/k8s.io/kubernetes/pkg/printers/internalversion/printers.go
index 7d7693c8a..296c5ce8b 100644
--- a/vendor/k8s.io/kubernetes/pkg/printers/internalversion/printers.go
+++ b/vendor/k8s.io/kubernetes/pkg/printers/internalversion/printers.go
@@ -1375,7 +1375,7 @@ func printCSINode(obj *storage.CSINode, options printers.GenerateOptions) ([]met
row := metav1.TableRow{
Object: runtime.RawExtension{Object: obj},
}
- row.Cells = append(row.Cells, obj.Name, len(obj.Spec.Drivers), translateTimestampSince(obj.CreationTimestamp))
+ row.Cells = append(row.Cells, obj.Name, int64(len(obj.Spec.Drivers)), translateTimestampSince(obj.CreationTimestamp))
return []metav1.TableRow{row}, nil
}
@@ -1481,7 +1481,7 @@ func printMutatingWebhook(obj *admissionregistration.MutatingWebhookConfiguratio
row := metav1.TableRow{
Object: runtime.RawExtension{Object: obj},
}
- row.Cells = append(row.Cells, obj.Name, len(obj.Webhooks), translateTimestampSince(obj.CreationTimestamp))
+ row.Cells = append(row.Cells, obj.Name, int64(len(obj.Webhooks)), translateTimestampSince(obj.CreationTimestamp))
return []metav1.TableRow{row}, nil
}
@@ -1501,7 +1501,7 @@ func printValidatingWebhook(obj *admissionregistration.ValidatingWebhookConfigur
row := metav1.TableRow{
Object: runtime.RawExtension{Object: obj},
}
- row.Cells = append(row.Cells, obj.Name, len(obj.Webhooks), translateTimestampSince(obj.CreationTimestamp))
+ row.Cells = append(row.Cells, obj.Name, int64(len(obj.Webhooks)), translateTimestampSince(obj.CreationTimestamp))
return []metav1.TableRow{row}, nil
}
@@ -2301,7 +2301,7 @@ func printStatus(obj *metav1.Status, options printers.GenerateOptions) ([]metav1
row := metav1.TableRow{
Object: runtime.RawExtension{Object: obj},
}
- row.Cells = append(row.Cells, obj.Status, obj.Reason, obj.Message)
+ row.Cells = append(row.Cells, obj.Status, string(obj.Reason), obj.Message)
return []metav1.TableRow{row}, nil
}
@@ -2530,7 +2530,7 @@ func printFlowSchema(obj *flowcontrol.FlowSchema, options printers.GenerateOptio
break
}
}
- row.Cells = append(row.Cells, name, plName, obj.Spec.MatchingPrecedence, distinguisherMethod, translateTimestampSince(obj.CreationTimestamp), badPLRef)
+ row.Cells = append(row.Cells, name, plName, int64(obj.Spec.MatchingPrecedence), distinguisherMethod, translateTimestampSince(obj.CreationTimestamp), badPLRef)
return []metav1.TableRow{row}, nil
}
diff --git a/vendor/k8s.io/kubernetes/pkg/proxy/endpointslicecache.go b/vendor/k8s.io/kubernetes/pkg/proxy/endpointslicecache.go
index de2fe43a8..bef55fa1c 100644
--- a/vendor/k8s.io/kubernetes/pkg/proxy/endpointslicecache.go
+++ b/vendor/k8s.io/kubernetes/pkg/proxy/endpointslicecache.go
@@ -80,6 +80,7 @@ type endpointSliceInfo struct {
type endpointInfo struct {
Addresses []string
NodeName *string
+ Topology map[string]string
Zone *string
ZoneHints sets.String
@@ -132,6 +133,7 @@ func newEndpointSliceInfo(endpointSlice *discovery.EndpointSlice, remove bool) *
Addresses: endpoint.Addresses,
Zone: endpoint.Zone,
NodeName: endpoint.NodeName,
+ Topology: endpoint.DeprecatedTopology,
// conditions
Ready: endpoint.Conditions.Ready == nil || *endpoint.Conditions.Ready,
@@ -284,6 +286,12 @@ func (cache *EndpointSliceCache) addEndpoints(serviceNN types.NamespacedName, po
if endpoint.NodeName != nil {
isLocal = cache.isLocal(*endpoint.NodeName)
nodeName = *endpoint.NodeName
+ } else {
+ deprecatedHostname, ok := endpoint.Topology[v1.LabelHostname]
+ if ok {
+ isLocal = cache.isLocal(deprecatedHostname)
+ nodeName = deprecatedHostname
+ }
}
zone := ""
diff --git a/vendor/k8s.io/kubernetes/pkg/proxy/iptables/proxier.go b/vendor/k8s.io/kubernetes/pkg/proxy/iptables/proxier.go
index 44021ee79..072274de5 100644
--- a/vendor/k8s.io/kubernetes/pkg/proxy/iptables/proxier.go
+++ b/vendor/k8s.io/kubernetes/pkg/proxy/iptables/proxier.go
@@ -190,7 +190,6 @@ type Proxier struct {
mu sync.Mutex // protects the following fields
serviceMap proxy.ServiceMap
endpointsMap proxy.EndpointsMap
- portsMap map[utilnet.LocalPort]utilnet.Closeable
nodeLabels map[string]string
// endpointSlicesSynced, and servicesSynced are set to true
// when corresponding objects are synced after startup. This is used to avoid
@@ -209,7 +208,6 @@ type Proxier struct {
localDetector proxyutiliptables.LocalTrafficDetector
hostname string
nodeIP net.IP
- portMapper utilnet.PortOpener
recorder events.EventRecorder
serviceHealthServer healthcheck.ServiceHealthServer
@@ -296,7 +294,6 @@ func NewProxier(ipt utiliptables.Interface,
}
proxier := &Proxier{
- portsMap: make(map[utilnet.LocalPort]utilnet.Closeable),
serviceMap: make(proxy.ServiceMap),
serviceChanges: proxy.NewServiceChangeTracker(newServiceInfo, ipFamily, recorder, nil),
endpointsMap: make(proxy.EndpointsMap),
@@ -309,7 +306,6 @@ func NewProxier(ipt utiliptables.Interface,
localDetector: localDetector,
hostname: hostname,
nodeIP: nodeIP,
- portMapper: &utilnet.ListenPortOpener,
recorder: recorder,
serviceHealthServer: serviceHealthServer,
healthzServer: healthzServer,
@@ -966,9 +962,6 @@ func (proxier *Proxier) syncProxyRules() {
// Accumulate NAT chains to keep.
activeNATChains := map[utiliptables.Chain]bool{} // use a map as a set
- // Accumulate the set of local ports that we will be holding open once this update is complete
- replacementPortsMap := map[utilnet.LocalPort]utilnet.Closeable{}
-
// We are creating those slices ones here to avoid memory reallocations
// in every loop. Note that reuse the memory, instead of doing:
// slice = <some new slice>
@@ -990,7 +983,6 @@ func (proxier *Proxier) syncProxyRules() {
proxier.endpointChainsNumber += len(proxier.endpointsMap[svcName])
}
- localAddrSet := utilproxy.GetLocalAddrSet()
nodeAddresses, err := utilproxy.GetNodeAddresses(proxier.nodePortAddresses, proxier.networkInterfacer)
if err != nil {
klog.ErrorS(err, "Failed to get node ip address matching nodeport cidrs, services with nodeport may not work as intended", "CIDRs", proxier.nodePortAddresses)
@@ -1004,10 +996,6 @@ func (proxier *Proxier) syncProxyRules() {
continue
}
isIPv6 := utilnet.IsIPv6(svcInfo.ClusterIP())
- localPortIPFamily := utilnet.IPv4
- if isIPv6 {
- localPortIPFamily = utilnet.IPv6
- }
protocol := strings.ToLower(string(svcInfo.Protocol()))
svcNameString := svcInfo.serviceNameString
@@ -1151,40 +1139,6 @@ func (proxier *Proxier) syncProxyRules() {
// Capture externalIPs.
for _, externalIP := range svcInfo.ExternalIPStrings() {
- // If the "external" IP happens to be an IP that is local to this
- // machine, hold the local port open so no other process can open it
- // (because the socket might open but it would never work).
- if (svcInfo.Protocol() != v1.ProtocolSCTP) && localAddrSet.Has(net.ParseIP(externalIP)) {
- lp := utilnet.LocalPort{
- Description: "externalIP for " + svcNameString,
- IP: externalIP,
- IPFamily: localPortIPFamily,
- Port: svcInfo.Port(),
- Protocol: utilnet.Protocol(svcInfo.Protocol()),
- }
- if proxier.portsMap[lp] != nil {
- klog.V(4).InfoS("Port was open before and is still needed", "port", lp.String())
- replacementPortsMap[lp] = proxier.portsMap[lp]
- } else {
- socket, err := proxier.portMapper.OpenLocalPort(&lp)
- if err != nil {
- msg := fmt.Sprintf("can't open port %s, skipping it", lp.String())
-
- proxier.recorder.Eventf(
- &v1.ObjectReference{
- Kind: "Node",
- Name: proxier.hostname,
- UID: types.UID(proxier.hostname),
- Namespace: "",
- }, nil, v1.EventTypeWarning, err.Error(), "SyncProxyRules", msg)
- klog.ErrorS(err, "can't open port, skipping it", "port", lp.String())
- continue
- }
- klog.V(2).InfoS("Opened local port", "port", lp.String())
- replacementPortsMap[lp] = socket
- }
- }
-
if hasEndpoints {
args = append(args[:0],
"-m", "comment", "--comment", fmt.Sprintf(`"%s external IP"`, svcNameString),
@@ -1306,56 +1260,7 @@ func (proxier *Proxier) syncProxyRules() {
// Capture nodeports. If we had more than 2 rules it might be
// worthwhile to make a new per-service chain for nodeport rules, but
// with just 2 rules it ends up being a waste and a cognitive burden.
- if svcInfo.NodePort() != 0 {
- // Hold the local port open so no other process can open it
- // (because the socket might open but it would never work).
- if len(nodeAddresses) == 0 {
- continue
- }
-
- lps := make([]utilnet.LocalPort, 0)
- for address := range nodeAddresses {
- lp := utilnet.LocalPort{
- Description: "nodePort for " + svcNameString,
- IP: address,
- IPFamily: localPortIPFamily,
- Port: svcInfo.NodePort(),
- Protocol: utilnet.Protocol(svcInfo.Protocol()),
- }
- if utilproxy.IsZeroCIDR(address) {
- // Empty IP address means all
- lp.IP = ""
- lps = append(lps, lp)
- // If we encounter a zero CIDR, then there is no point in processing the rest of the addresses.
- break
- }
- lps = append(lps, lp)
- }
-
- // For ports on node IPs, open the actual port and hold it.
- for _, lp := range lps {
- if proxier.portsMap[lp] != nil {
- klog.V(4).InfoS("Port was open before and is still needed", "port", lp.String())
- replacementPortsMap[lp] = proxier.portsMap[lp]
- } else if svcInfo.Protocol() != v1.ProtocolSCTP {
- socket, err := proxier.portMapper.OpenLocalPort(&lp)
- if err != nil {
- msg := fmt.Sprintf("can't open port %s, skipping it", lp.String())
-
- proxier.recorder.Eventf(
- &v1.ObjectReference{
- Kind: "Node",
- Name: proxier.hostname,
- UID: types.UID(proxier.hostname),
- Namespace: "",
- }, nil, v1.EventTypeWarning, err.Error(), "SyncProxyRules", msg)
- klog.ErrorS(err, "can't open port, skipping it", "port", lp.String())
- continue
- }
- klog.V(2).InfoS("Opened local port", "port", lp.String())
- replacementPortsMap[lp] = socket
- }
- }
+ if svcInfo.NodePort() != 0 && len(nodeAddresses) != 0 {
if hasEndpoints {
args = append(args[:0],
@@ -1623,9 +1528,6 @@ func (proxier *Proxier) syncProxyRules() {
if err != nil {
klog.ErrorS(err, "Failed to execute iptables-restore")
metrics.IptablesRestoreFailuresTotal.Inc()
- // Revert new local ports.
- klog.V(2).InfoS("Closing local ports after iptables-restore failure")
- utilproxy.RevertPorts(replacementPortsMap, proxier.portsMap)
return
}
success = true
@@ -1638,14 +1540,6 @@ func (proxier *Proxier) syncProxyRules() {
}
}
- // Close old local ports and save new ones.
- for k, v := range proxier.portsMap {
- if replacementPortsMap[k] == nil {
- v.Close()
- }
- }
- proxier.portsMap = replacementPortsMap
-
if proxier.healthzServer != nil {
proxier.healthzServer.Updated()
}
diff --git a/vendor/k8s.io/kubernetes/pkg/proxy/util/utils.go b/vendor/k8s.io/kubernetes/pkg/proxy/util/utils.go
index c1286029a..3d7f75f74 100644
--- a/vendor/k8s.io/kubernetes/pkg/proxy/util/utils.go
+++ b/vendor/k8s.io/kubernetes/pkg/proxy/util/utils.go
@@ -96,7 +96,7 @@ func IsProxyableIP(ip string) error {
}
func isProxyableIP(ip net.IP) error {
- if ip.IsLoopback() || ip.IsLinkLocalUnicast() || ip.IsLinkLocalMulticast() || ip.IsInterfaceLocalMulticast() {
+ if !ip.IsGlobalUnicast() {
return ErrAddressNotAllowed
}
return nil
diff --git a/vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/nodeaffinity/node_affinity.go b/vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/nodeaffinity/node_affinity.go
index 11934a4a3..d0fafd1fa 100644
--- a/vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/nodeaffinity/node_affinity.go
+++ b/vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/nodeaffinity/node_affinity.go
@@ -78,7 +78,7 @@ func (s *preFilterState) Clone() framework.StateData {
// failed by this plugin schedulable.
func (pl *NodeAffinity) EventsToRegister() []framework.ClusterEvent {
return []framework.ClusterEvent{
- {Resource: framework.Node, ActionType: framework.Add | framework.UpdateNodeLabel},
+ {Resource: framework.Node, ActionType: framework.Add | framework.Update},
}
}
diff --git a/vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/nodename/node_name.go b/vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/nodename/node_name.go
index 356fc5c29..e02134f79 100644
--- a/vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/nodename/node_name.go
+++ b/vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/nodename/node_name.go
@@ -43,7 +43,7 @@ const (
// failed by this plugin schedulable.
func (pl *NodeName) EventsToRegister() []framework.ClusterEvent {
return []framework.ClusterEvent{
- {Resource: framework.Node, ActionType: framework.Add},
+ {Resource: framework.Node, ActionType: framework.Add | framework.Update},
}
}
diff --git a/vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/nodeports/node_ports.go b/vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/nodeports/node_ports.go
index 12ffc87a7..e77f8411a 100644
--- a/vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/nodeports/node_ports.go
+++ b/vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/nodeports/node_ports.go
@@ -105,7 +105,7 @@ func (pl *NodePorts) EventsToRegister() []framework.ClusterEvent {
return []framework.ClusterEvent{
// Due to immutable fields `spec.containers[*].ports`, pod update events are ignored.
{Resource: framework.Pod, ActionType: framework.Delete},
- {Resource: framework.Node, ActionType: framework.Add},
+ {Resource: framework.Node, ActionType: framework.Add | framework.Update},
}
}
diff --git a/vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/noderesources/fit.go b/vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/noderesources/fit.go
index 268783a43..26c81e515 100644
--- a/vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/noderesources/fit.go
+++ b/vendor/k8s.io/kubernetes/pkg/scheduler/framework/plugins/noderesources/fit.go
@@ -211,7 +211,7 @@ func getPreFilterState(cycleState *framework.CycleState) (*preFilterState, error
func (f *Fit) EventsToRegister() []framework.ClusterEvent {
return []framework.ClusterEvent{
{Resource: framework.Pod, ActionType: framework.Delete},
- {Resource: framework.Node, ActionType: framework.Add | framework.UpdateNodeAllocatable},
+ {Resource: framework.Node, ActionType: framework.Add | framework.Update},
}
}
diff --git a/vendor/k8s.io/kubernetes/pkg/volume/csi/csi_attacher.go b/vendor/k8s.io/kubernetes/pkg/volume/csi/csi_attacher.go
index 3528dc3b8..7024c4a83 100644
--- a/vendor/k8s.io/kubernetes/pkg/volume/csi/csi_attacher.go
+++ b/vendor/k8s.io/kubernetes/pkg/volume/csi/csi_attacher.go
@@ -413,30 +413,17 @@ func (c *csiAttacher) Detach(volumeName string, nodeName types.NodeName) error {
return errors.New("missing expected parameter volumeName")
}
- if isAttachmentName(volumeName) {
- // Detach can also be called with the attach ID as the `volumeName`. This codepath is
- // hit only when we have migrated an in-tree volume to CSI and the A/D Controller is shut
- // down, the pod with the volume is deleted, and the A/D Controller starts back up in that
- // order.
- attachID = volumeName
-
- // Vol ID should be the volume handle, except that is not available here.
- // It is only used in log messages so in the event that this happens log messages will be
- // printing out the attachID instead of the volume handle.
- volID = volumeName
- } else {
- // volumeName in format driverName<SEP>volumeHandle generated by plugin.GetVolumeName()
- parts := strings.Split(volumeName, volNameSep)
- if len(parts) != 2 {
- klog.Error(log("detacher.Detach insufficient info encoded in volumeName"))
- return errors.New("volumeName missing expected data")
- }
-
- driverName := parts[0]
- volID = parts[1]
- attachID = getAttachmentName(volID, driverName, string(nodeName))
+ // volumeName in format driverName<SEP>volumeHandle generated by plugin.GetVolumeName()
+ parts := strings.Split(volumeName, volNameSep)
+ if len(parts) != 2 {
+ klog.Error(log("detacher.Detach insufficient info encoded in volumeName"))
+ return errors.New("volumeName missing expected data")
}
+ driverName := parts[0]
+ volID = parts[1]
+ attachID = getAttachmentName(volID, driverName, string(nodeName))
+
if err := c.k8s.StorageV1().VolumeAttachments().Delete(context.TODO(), attachID, metav1.DeleteOptions{}); err != nil {
if apierrors.IsNotFound(err) {
// object deleted or never existed, done
@@ -640,13 +627,6 @@ func getAttachmentName(volName, csiDriverName, nodeName string) string {
return fmt.Sprintf("csi-%x", result)
}
-// isAttachmentName returns true if the string given is of the form of an Attach ID
-// and false otherwise
-func isAttachmentName(unknownString string) bool {
- // 68 == "csi-" + len(sha256hash)
- return strings.HasPrefix(unknownString, "csi-") && len(unknownString) == 68
-}
-
func makeDeviceMountPath(plugin *csiPlugin, spec *volume.Spec) (string, error) {
if spec == nil {
return "", errors.New(log("makeDeviceMountPath failed, spec is nil"))
diff --git a/vendor/k8s.io/kubernetes/pkg/volume/plugins.go b/vendor/k8s.io/kubernetes/pkg/volume/plugins.go
index 17f149804..1d07b8f80 100644
--- a/vendor/k8s.io/kubernetes/pkg/volume/plugins.go
+++ b/vendor/k8s.io/kubernetes/pkg/volume/plugins.go
@@ -445,6 +445,8 @@ type VolumeHost interface {
// Returns the name of the node
GetNodeName() types.NodeName
+ GetAttachedVolumesFromNodeStatus() (map[v1.UniqueVolumeName]string, error)
+
// Returns the event recorder of kubelet.
GetEventRecorder() record.EventRecorder
diff --git a/vendor/k8s.io/kubernetes/pkg/volume/util/operationexecutor/operation_executor.go b/vendor/k8s.io/kubernetes/pkg/volume/util/operationexecutor/operation_executor.go
index ca4f6ca0c..d7ec7989e 100644
--- a/vendor/k8s.io/kubernetes/pkg/volume/util/operationexecutor/operation_executor.go
+++ b/vendor/k8s.io/kubernetes/pkg/volume/util/operationexecutor/operation_executor.go
@@ -21,6 +21,7 @@ limitations under the License.
package operationexecutor
import (
+ "errors"
"fmt"
"time"
@@ -418,6 +419,23 @@ const (
VolumeNotMounted VolumeMountState = "VolumeNotMounted"
)
+type MountPreConditionFailed struct {
+ msg string
+}
+
+func (err *MountPreConditionFailed) Error() string {
+ return err.msg
+}
+
+func NewMountPreConditionFailedError(msg string) *MountPreConditionFailed {
+ return &MountPreConditionFailed{msg: msg}
+}
+
+func IsMountFailedPreconditionError(err error) bool {
+ var failedPreconditionError *MountPreConditionFailed
+ return errors.As(err, &failedPreconditionError)
+}
+
// GenerateMsgDetailed returns detailed msgs for volumes to mount
func (volume *VolumeToMount) GenerateMsgDetailed(prefixMsg, suffixMsg string) (detailedMsg string) {
detailedStr := fmt.Sprintf("(UniqueName: %q) pod %q (UID: %q)", volume.VolumeName, volume.Pod.Name, volume.Pod.UID)
diff --git a/vendor/k8s.io/kubernetes/pkg/volume/util/operationexecutor/operation_generator.go b/vendor/k8s.io/kubernetes/pkg/volume/util/operationexecutor/operation_generator.go
index 96d29c2e7..ab42fc7df 100644
--- a/vendor/k8s.io/kubernetes/pkg/volume/util/operationexecutor/operation_generator.go
+++ b/vendor/k8s.io/kubernetes/pkg/volume/util/operationexecutor/operation_generator.go
@@ -46,9 +46,10 @@ import (
)
const (
- unknownVolumePlugin string = "UnknownVolumePlugin"
- unknownAttachableVolumePlugin string = "UnknownAttachableVolumePlugin"
- DetachOperationName string = "volume_detach"
+ unknownVolumePlugin string = "UnknownVolumePlugin"
+ unknownAttachableVolumePlugin string = "UnknownAttachableVolumePlugin"
+ DetachOperationName string = "volume_detach"
+ VerifyControllerAttachedVolumeOpName string = "verify_controller_attached_volume"
)
// InTreeToCSITranslator contains methods required to check migratable status
@@ -966,6 +967,12 @@ func (og *operationGenerator) GenerateUnmountDeviceFunc(
}
// The device is still in use elsewhere. Caller will log and retry.
if deviceOpened {
+ // Mark the device as uncertain, so MountDevice is called for new pods.
+ markDeviceUncertainErr := actualStateOfWorld.MarkDeviceAsUncertain(deviceToDetach.VolumeName, deviceToDetach.DevicePath, deviceMountPath)
+ if markDeviceUncertainErr != nil {
+ // There is nothing else we can do. Hope that UnmountDevice will be re-tried shortly.
+ klog.Errorf(deviceToDetach.GenerateErrorDetailed("UnmountDevice.MarkDeviceAsUncertain failed", markDeviceUncertainErr).Error())
+ }
eventErr, detailedErr := deviceToDetach.GenerateError(
"UnmountDevice failed",
goerrors.New("the device is in use when it was no longer expected to be in use"))
@@ -1472,6 +1479,20 @@ func (og *operationGenerator) GenerateVerifyControllerAttachedVolumeFunc(
return volumetypes.GeneratedOperations{}, volumeToMount.GenerateErrorDetailed("VerifyControllerAttachedVolume.FindPluginBySpec failed", err)
}
+ // For attachable volume types, lets check if volume is attached by reading from node lister.
+ // This would avoid exponential back-off and creation of goroutine unnecessarily. We still
+ // verify status of attached volume by directly reading from API server later on.This is necessarily
+ // to ensure any race conditions because of cached state in the informer.
+ if volumeToMount.PluginIsAttachable {
+ cachedAttachedVolumes, _ := og.volumePluginMgr.Host.GetAttachedVolumesFromNodeStatus()
+ if cachedAttachedVolumes != nil {
+ _, volumeFound := cachedAttachedVolumes[volumeToMount.VolumeName]
+ if !volumeFound {
+ return volumetypes.GeneratedOperations{}, NewMountPreConditionFailedError(fmt.Sprintf("volume %s is not yet in node's status", volumeToMount.VolumeName))
+ }
+ }
+ }
+
verifyControllerAttachedVolumeFunc := func() volumetypes.OperationContext {
migrated := getMigratedStatusBySpec(volumeToMount.VolumeSpec)
if !volumeToMount.PluginIsAttachable {
@@ -1537,7 +1558,7 @@ func (og *operationGenerator) GenerateVerifyControllerAttachedVolumeFunc(
}
return volumetypes.GeneratedOperations{
- OperationName: "verify_controller_attached_volume",
+ OperationName: VerifyControllerAttachedVolumeOpName,
OperationFunc: verifyControllerAttachedVolumeFunc,
CompleteFunc: util.OperationCompleteHook(util.GetFullQualifiedPluginNameForVolume(volumePlugin.GetPluginName(), volumeToMount.VolumeSpec), "verify_controller_attached_volume"),
EventRecorderFunc: nil, // nil because we do not want to generate event on error
diff --git a/vendor/k8s.io/mount-utils/go.mod b/vendor/k8s.io/mount-utils/go.mod
index d61ed417e..13985ce9a 100644
--- a/vendor/k8s.io/mount-utils/go.mod
+++ b/vendor/k8s.io/mount-utils/go.mod
@@ -11,7 +11,7 @@ require (
gopkg.in/check.v1 v1.0.0-20200227125254-8fa46927fb4f // indirect
gopkg.in/yaml.v3 v3.0.0-20210107192922-496545a6307b // indirect
k8s.io/klog/v2 v2.9.0
- k8s.io/utils v0.0.0-20210819203725-bdf08cb9a70a
+ k8s.io/utils v0.0.0-20211116205334-6203023598ed
)
replace k8s.io/mount-utils => ../mount-utils
diff --git a/vendor/k8s.io/mount-utils/go.sum b/vendor/k8s.io/mount-utils/go.sum
index 2c7cfce98..dc0d2e4db 100644
--- a/vendor/k8s.io/mount-utils/go.sum
+++ b/vendor/k8s.io/mount-utils/go.sum
@@ -28,5 +28,5 @@ gopkg.in/yaml.v3 v3.0.0-20210107192922-496545a6307b/go.mod h1:K4uyk7z7BCEPqu6E+C
k8s.io/klog/v2 v2.0.0/go.mod h1:PBfzABfn139FHAV07az/IF9Wp1bkk3vpT2XSJ76fSDE=
k8s.io/klog/v2 v2.9.0 h1:D7HV+n1V57XeZ0m6tdRkfknthUaM06VFbWldOFh8kzM=
k8s.io/klog/v2 v2.9.0/go.mod h1:hy9LJ/NvuK+iVyP4Ehqva4HxZG/oXyIS3n3Jmire4Ec=
-k8s.io/utils v0.0.0-20210819203725-bdf08cb9a70a h1:8dYfu/Fc9Gz2rNJKB9IQRGgQOh2clmRzNIPPY1xLY5g=
-k8s.io/utils v0.0.0-20210819203725-bdf08cb9a70a/go.mod h1:jPW/WVKK9YHAvNhRxK0md/EJ228hCsBRufyofKtW8HA=
+k8s.io/utils v0.0.0-20211116205334-6203023598ed h1:ck1fRPWPJWsMd8ZRFsWc6mh/zHp5fZ/shhbrgPUxDAE=
+k8s.io/utils v0.0.0-20211116205334-6203023598ed/go.mod h1:jPW/WVKK9YHAvNhRxK0md/EJ228hCsBRufyofKtW8HA=
diff --git a/vendor/k8s.io/utils/clock/clock.go b/vendor/k8s.io/utils/clock/clock.go
index a3efed96a..dd181ce8d 100644
--- a/vendor/k8s.io/utils/clock/clock.go
+++ b/vendor/k8s.io/utils/clock/clock.go
@@ -53,6 +53,16 @@ type WithTicker interface {
NewTicker(time.Duration) Ticker
}
+// WithDelayedExecution allows for injecting fake or real clocks into
+// code that needs to make use of AfterFunc functionality.
+type WithDelayedExecution interface {
+ Clock
+ // AfterFunc executes f in its own goroutine after waiting
+ // for d duration and returns a Timer whose channel can be
+ // closed by calling Stop() on the Timer.
+ AfterFunc(d time.Duration, f func()) Timer
+}
+
// Ticker defines the Ticker interface.
type Ticker interface {
C() <-chan time.Time
@@ -88,6 +98,13 @@ func (RealClock) NewTimer(d time.Duration) Timer {
}
}
+// AfterFunc is the same as time.AfterFunc(d, f).
+func (RealClock) AfterFunc(d time.Duration, f func()) Timer {
+ return &realTimer{
+ timer: time.AfterFunc(d, f),
+ }
+}
+
// Tick is the same as time.Tick(d)
// This method does not allow to free/GC the backing ticker. Use
// NewTicker instead.
diff --git a/vendor/k8s.io/utils/inotify/inotify_linux.go b/vendor/k8s.io/utils/inotify/inotify_linux.go
index 6258277c9..2963042e4 100644
--- a/vendor/k8s.io/utils/inotify/inotify_linux.go
+++ b/vendor/k8s.io/utils/inotify/inotify_linux.go
@@ -120,7 +120,11 @@ func (w *Watcher) RemoveWatch(path string) error {
}
success, errno := syscall.InotifyRmWatch(w.fd, watch.wd)
if success == -1 {
- return os.NewSyscallError("inotify_rm_watch", errno)
+ // when file descriptor or watch descriptor not found, InotifyRmWatch syscall return EINVAL error
+ // if return error, it may lead this path remain in watches and paths map, and no other event can trigger remove action.
+ if errno != syscall.EINVAL {
+ return os.NewSyscallError("inotify_rm_watch", errno)
+ }
}
delete(w.watches, path)
// Locking here to protect the read from paths in readEvents.
diff --git a/vendor/modules.txt b/vendor/modules.txt
index f0649f926..f4251002c 100644
--- a/vendor/modules.txt
+++ b/vendor/modules.txt
@@ -344,7 +344,7 @@ github.com/golang/protobuf/ptypes/timestamp
github.com/golang/protobuf/ptypes/wrappers
# github.com/google/btree v1.0.1
github.com/google/btree
-# github.com/google/cadvisor v0.39.3
+# github.com/google/cadvisor v0.39.4
github.com/google/cadvisor/accelerators
github.com/google/cadvisor/cache/memory
github.com/google/cadvisor/collector
@@ -1066,7 +1066,7 @@ helm.sh/helm/v3/pkg/storage
helm.sh/helm/v3/pkg/storage/driver
helm.sh/helm/v3/pkg/strvals
helm.sh/helm/v3/pkg/time
-# k8s.io/api v0.22.6 => github.com/kubeedge/kubernetes/staging/src/k8s.io/api v1.22.6-kubeedge4
+# k8s.io/api v0.22.17 => github.com/kubeedge/kubernetes/staging/src/k8s.io/api v1.22.17-kubeedge1
## explicit
k8s.io/api/admission/v1
k8s.io/api/admission/v1beta1
@@ -1114,13 +1114,13 @@ k8s.io/api/scheduling/v1beta1
k8s.io/api/storage/v1
k8s.io/api/storage/v1alpha1
k8s.io/api/storage/v1beta1
-# k8s.io/apiextensions-apiserver v0.22.6 => github.com/kubeedge/kubernetes/staging/src/k8s.io/apiextensions-apiserver v1.22.6-kubeedge4
+# k8s.io/apiextensions-apiserver v0.22.17 => github.com/kubeedge/kubernetes/staging/src/k8s.io/apiextensions-apiserver v1.22.17-kubeedge1
## explicit
k8s.io/apiextensions-apiserver/pkg/apis/apiextensions
k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1
k8s.io/apiextensions-apiserver/pkg/apis/apiextensions/v1beta1
k8s.io/apiextensions-apiserver/pkg/crdserverscheme
-# k8s.io/apimachinery v0.22.6 => github.com/kubeedge/kubernetes/staging/src/k8s.io/apimachinery v1.22.6-kubeedge4
+# k8s.io/apimachinery v0.22.17 => github.com/kubeedge/kubernetes/staging/src/k8s.io/apimachinery v1.22.17-kubeedge1
## explicit
k8s.io/apimachinery/pkg/api/equality
k8s.io/apimachinery/pkg/api/errors
@@ -1184,7 +1184,7 @@ k8s.io/apimachinery/pkg/watch
k8s.io/apimachinery/third_party/forked/golang/json
k8s.io/apimachinery/third_party/forked/golang/netutil
k8s.io/apimachinery/third_party/forked/golang/reflect
-# k8s.io/apiserver v0.22.6 => github.com/kubeedge/kubernetes/staging/src/k8s.io/apiserver v1.22.6-kubeedge4
+# k8s.io/apiserver v0.22.17 => github.com/kubeedge/kubernetes/staging/src/k8s.io/apiserver v1.22.17-kubeedge1
## explicit
k8s.io/apiserver/pkg/admission
k8s.io/apiserver/pkg/admission/configuration
@@ -1313,12 +1313,12 @@ k8s.io/apiserver/plugin/pkg/audit/truncate
k8s.io/apiserver/plugin/pkg/audit/webhook
k8s.io/apiserver/plugin/pkg/authenticator/token/webhook
k8s.io/apiserver/plugin/pkg/authorizer/webhook
-# k8s.io/cli-runtime v0.22.6 => github.com/kubeedge/kubernetes/staging/src/k8s.io/cli-runtime v1.22.6-kubeedge4
+# k8s.io/cli-runtime v0.22.17 => github.com/kubeedge/kubernetes/staging/src/k8s.io/cli-runtime v1.22.17-kubeedge1
## explicit
k8s.io/cli-runtime/pkg/genericclioptions
k8s.io/cli-runtime/pkg/printers
k8s.io/cli-runtime/pkg/resource
-# k8s.io/client-go v0.22.6 => github.com/kubeedge/kubernetes/staging/src/k8s.io/client-go v1.22.6-kubeedge4
+# k8s.io/client-go v0.22.17 => github.com/kubeedge/kubernetes/staging/src/k8s.io/client-go v1.22.17-kubeedge1
## explicit
k8s.io/client-go/applyconfigurations/admissionregistration/v1
k8s.io/client-go/applyconfigurations/admissionregistration/v1beta1
@@ -1603,16 +1603,16 @@ k8s.io/client-go/util/homedir
k8s.io/client-go/util/jsonpath
k8s.io/client-go/util/keyutil
k8s.io/client-go/util/workqueue
-# k8s.io/cloud-provider v0.22.6 => github.com/kubeedge/kubernetes/staging/src/k8s.io/cloud-provider v1.22.6-kubeedge4
+# k8s.io/cloud-provider v0.22.17 => github.com/kubeedge/kubernetes/staging/src/k8s.io/cloud-provider v1.22.17-kubeedge1
## explicit
k8s.io/cloud-provider/credentialconfig
k8s.io/cloud-provider/volume/errors
-# k8s.io/cluster-bootstrap v0.22.6 => github.com/kubeedge/kubernetes/staging/src/k8s.io/cluster-bootstrap v1.22.6-kubeedge4
+# k8s.io/cluster-bootstrap v0.22.17 => github.com/kubeedge/kubernetes/staging/src/k8s.io/cluster-bootstrap v1.22.17-kubeedge1
## explicit
k8s.io/cluster-bootstrap/token/api
k8s.io/cluster-bootstrap/token/util
k8s.io/cluster-bootstrap/util/secrets
-# k8s.io/code-generator v0.22.6 => github.com/kubeedge/kubernetes/staging/src/k8s.io/code-generator v1.22.6-kubeedge4
+# k8s.io/code-generator v0.22.17 => github.com/kubeedge/kubernetes/staging/src/k8s.io/code-generator v1.22.17-kubeedge1
## explicit
k8s.io/code-generator/cmd/client-gen
k8s.io/code-generator/cmd/client-gen/args
@@ -1634,7 +1634,7 @@ k8s.io/code-generator/cmd/lister-gen/args
k8s.io/code-generator/cmd/lister-gen/generators
k8s.io/code-generator/pkg/namer
k8s.io/code-generator/pkg/util
-# k8s.io/component-base v0.22.6 => github.com/kubeedge/kubernetes/staging/src/k8s.io/component-base v1.22.6-kubeedge4
+# k8s.io/component-base v0.22.17 => github.com/kubeedge/kubernetes/staging/src/k8s.io/component-base v1.22.17-kubeedge1
## explicit
k8s.io/component-base/cli/flag
k8s.io/component-base/cli/globalflag
@@ -1656,17 +1656,18 @@ k8s.io/component-base/term
k8s.io/component-base/traces
k8s.io/component-base/version
k8s.io/component-base/version/verflag
-# k8s.io/component-helpers v0.0.0 => github.com/kubeedge/kubernetes/staging/src/k8s.io/component-helpers v1.22.6-kubeedge4
+# k8s.io/component-helpers v0.0.0 => github.com/kubeedge/kubernetes/staging/src/k8s.io/component-helpers v1.22.17-kubeedge1
k8s.io/component-helpers/apimachinery/lease
k8s.io/component-helpers/scheduling/corev1
k8s.io/component-helpers/scheduling/corev1/nodeaffinity
k8s.io/component-helpers/storage/volume
-# k8s.io/cri-api v0.22.6 => github.com/kubeedge/kubernetes/staging/src/k8s.io/cri-api v1.22.6-kubeedge4
+# k8s.io/cri-api v0.22.17 => github.com/kubeedge/kubernetes/staging/src/k8s.io/cri-api v1.22.17-kubeedge1
## explicit
k8s.io/cri-api/pkg/apis
k8s.io/cri-api/pkg/apis/runtime/v1alpha2
k8s.io/cri-api/pkg/apis/testing
-# k8s.io/csi-translation-lib v0.22.6 => github.com/kubeedge/kubernetes/staging/src/k8s.io/csi-translation-lib v1.22.6-kubeedge4
+k8s.io/cri-api/pkg/errors
+# k8s.io/csi-translation-lib v0.22.17 => github.com/kubeedge/kubernetes/staging/src/k8s.io/csi-translation-lib v1.22.17-kubeedge1
## explicit
k8s.io/csi-translation-lib
k8s.io/csi-translation-lib/plugins
@@ -1697,13 +1698,13 @@ k8s.io/kube-openapi/pkg/util/proto
k8s.io/kube-openapi/pkg/util/proto/validation
k8s.io/kube-openapi/pkg/util/sets
k8s.io/kube-openapi/pkg/validation/spec
-# k8s.io/kube-scheduler v0.22.6 => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-scheduler v1.22.6-kubeedge4
+# k8s.io/kube-scheduler v0.22.17 => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-scheduler v1.22.17-kubeedge1
## explicit
k8s.io/kube-scheduler/config/v1
k8s.io/kube-scheduler/config/v1beta1
k8s.io/kube-scheduler/config/v1beta2
k8s.io/kube-scheduler/extender/v1
-# k8s.io/kubectl v0.22.4 => github.com/kubeedge/kubernetes/staging/src/k8s.io/kubectl v1.22.6-kubeedge4
+# k8s.io/kubectl v0.22.4 => github.com/kubeedge/kubernetes/staging/src/k8s.io/kubectl v1.22.17-kubeedge1
k8s.io/kubectl/pkg/cmd/util
k8s.io/kubectl/pkg/scheme
k8s.io/kubectl/pkg/util/interrupt
@@ -1712,7 +1713,7 @@ k8s.io/kubectl/pkg/util/openapi/validation
k8s.io/kubectl/pkg/util/templates
k8s.io/kubectl/pkg/util/term
k8s.io/kubectl/pkg/validation
-# k8s.io/kubelet v0.22.6 => github.com/kubeedge/kubernetes/staging/src/k8s.io/kubelet v1.22.6-kubeedge4
+# k8s.io/kubelet v0.22.17 => github.com/kubeedge/kubernetes/staging/src/k8s.io/kubelet v1.22.17-kubeedge1
## explicit
k8s.io/kubelet/config/v1alpha1
k8s.io/kubelet/config/v1beta1
@@ -1724,7 +1725,7 @@ k8s.io/kubelet/pkg/apis/deviceplugin/v1beta1
k8s.io/kubelet/pkg/apis/pluginregistration/v1
k8s.io/kubelet/pkg/apis/podresources/v1
k8s.io/kubelet/pkg/apis/stats/v1alpha1
-# k8s.io/kubernetes v1.22.6 => github.com/kubeedge/kubernetes v1.22.6-kubeedge4
+# k8s.io/kubernetes v1.22.17 => github.com/kubeedge/kubernetes v1.22.17-kubeedge1
## explicit
k8s.io/kubernetes/cmd/kubeadm/app/apis/bootstraptoken/v1
k8s.io/kubernetes/cmd/kubeadm/app/apis/kubeadm
@@ -2035,15 +2036,15 @@ k8s.io/kubernetes/pkg/volume/util/volumepathhandler
k8s.io/kubernetes/pkg/volume/validation
k8s.io/kubernetes/pkg/windows/service
k8s.io/kubernetes/third_party/forked/golang/expansion
-# k8s.io/legacy-cloud-providers v0.0.0 => github.com/kubeedge/kubernetes/staging/src/k8s.io/legacy-cloud-providers v1.22.6-kubeedge4
+# k8s.io/legacy-cloud-providers v0.0.0 => github.com/kubeedge/kubernetes/staging/src/k8s.io/legacy-cloud-providers v1.22.17-kubeedge1
k8s.io/legacy-cloud-providers/azure/auth
k8s.io/legacy-cloud-providers/gce/gcpcredential
-# k8s.io/mount-utils v0.22.6 => github.com/kubeedge/kubernetes/staging/src/k8s.io/mount-utils v1.22.6-kubeedge4
+# k8s.io/mount-utils v0.22.17 => github.com/kubeedge/kubernetes/staging/src/k8s.io/mount-utils v1.22.17-kubeedge1
## explicit
k8s.io/mount-utils
# k8s.io/system-validators v1.5.0
k8s.io/system-validators/validators
-# k8s.io/utils v0.0.0-20210819203725-bdf08cb9a70a
+# k8s.io/utils v0.0.0-20211116205334-6203023598ed
## explicit
k8s.io/utils/buffer
k8s.io/utils/clock
@@ -2081,7 +2082,7 @@ sigs.k8s.io/apiserver-network-proxy/pkg/server/metrics
sigs.k8s.io/apiserver-network-proxy/pkg/util
sigs.k8s.io/apiserver-network-proxy/proto/agent
sigs.k8s.io/apiserver-network-proxy/proto/header
-# sigs.k8s.io/apiserver-network-proxy/konnectivity-client v0.0.27 => sigs.k8s.io/apiserver-network-proxy/konnectivity-client v0.0.27
+# sigs.k8s.io/apiserver-network-proxy/konnectivity-client v0.0.30 => sigs.k8s.io/apiserver-network-proxy/konnectivity-client v0.0.27
## explicit
sigs.k8s.io/apiserver-network-proxy/konnectivity-client/pkg/client
sigs.k8s.io/apiserver-network-proxy/konnectivity-client/proto/client
@@ -2227,37 +2228,37 @@ sigs.k8s.io/yaml
# go.etcd.io/etcd/pkg/v3 => go.etcd.io/etcd/pkg/v3 v3.5.0
# go.etcd.io/etcd/raft/v3 => go.etcd.io/etcd/raft/v3 v3.5.0
# go.etcd.io/etcd/server/v3 => go.etcd.io/etcd/server/v3 v3.5.0
-# k8s.io/api => github.com/kubeedge/kubernetes/staging/src/k8s.io/api v1.22.6-kubeedge4
-# k8s.io/apiextensions-apiserver => github.com/kubeedge/kubernetes/staging/src/k8s.io/apiextensions-apiserver v1.22.6-kubeedge4
-# k8s.io/apimachinery => github.com/kubeedge/kubernetes/staging/src/k8s.io/apimachinery v1.22.6-kubeedge4
-# k8s.io/apiserver => github.com/kubeedge/kubernetes/staging/src/k8s.io/apiserver v1.22.6-kubeedge4
-# k8s.io/cli-runtime => github.com/kubeedge/kubernetes/staging/src/k8s.io/cli-runtime v1.22.6-kubeedge4
-# k8s.io/client-go => github.com/kubeedge/kubernetes/staging/src/k8s.io/client-go v1.22.6-kubeedge4
-# k8s.io/cloud-provider => github.com/kubeedge/kubernetes/staging/src/k8s.io/cloud-provider v1.22.6-kubeedge4
-# k8s.io/cluster-bootstrap => github.com/kubeedge/kubernetes/staging/src/k8s.io/cluster-bootstrap v1.22.6-kubeedge4
-# k8s.io/code-generator => github.com/kubeedge/kubernetes/staging/src/k8s.io/code-generator v1.22.6-kubeedge4
-# k8s.io/component-base => github.com/kubeedge/kubernetes/staging/src/k8s.io/component-base v1.22.6-kubeedge4
-# k8s.io/component-helpers => github.com/kubeedge/kubernetes/staging/src/k8s.io/component-helpers v1.22.6-kubeedge4
-# k8s.io/controller-manager => github.com/kubeedge/kubernetes/staging/src/k8s.io/controller-manager v1.22.6-kubeedge4
-# k8s.io/cri-api => github.com/kubeedge/kubernetes/staging/src/k8s.io/cri-api v1.22.6-kubeedge4
-# k8s.io/csi-api => github.com/kubeedge/kubernetes/staging/src/k8s.io/csi-api v1.22.6-kubeedge4
-# k8s.io/csi-translation-lib => github.com/kubeedge/kubernetes/staging/src/k8s.io/csi-translation-lib v1.22.6-kubeedge4
-# k8s.io/gengo v0.0.0 => k8s.io/gengo v0.22.6
+# k8s.io/api => github.com/kubeedge/kubernetes/staging/src/k8s.io/api v1.22.17-kubeedge1
+# k8s.io/apiextensions-apiserver => github.com/kubeedge/kubernetes/staging/src/k8s.io/apiextensions-apiserver v1.22.17-kubeedge1
+# k8s.io/apimachinery => github.com/kubeedge/kubernetes/staging/src/k8s.io/apimachinery v1.22.17-kubeedge1
+# k8s.io/apiserver => github.com/kubeedge/kubernetes/staging/src/k8s.io/apiserver v1.22.17-kubeedge1
+# k8s.io/cli-runtime => github.com/kubeedge/kubernetes/staging/src/k8s.io/cli-runtime v1.22.17-kubeedge1
+# k8s.io/client-go => github.com/kubeedge/kubernetes/staging/src/k8s.io/client-go v1.22.17-kubeedge1
+# k8s.io/cloud-provider => github.com/kubeedge/kubernetes/staging/src/k8s.io/cloud-provider v1.22.17-kubeedge1
+# k8s.io/cluster-bootstrap => github.com/kubeedge/kubernetes/staging/src/k8s.io/cluster-bootstrap v1.22.17-kubeedge1
+# k8s.io/code-generator => github.com/kubeedge/kubernetes/staging/src/k8s.io/code-generator v1.22.17-kubeedge1
+# k8s.io/component-base => github.com/kubeedge/kubernetes/staging/src/k8s.io/component-base v1.22.17-kubeedge1
+# k8s.io/component-helpers => github.com/kubeedge/kubernetes/staging/src/k8s.io/component-helpers v1.22.17-kubeedge1
+# k8s.io/controller-manager => github.com/kubeedge/kubernetes/staging/src/k8s.io/controller-manager v1.22.17-kubeedge1
+# k8s.io/cri-api => github.com/kubeedge/kubernetes/staging/src/k8s.io/cri-api v1.22.17-kubeedge1
+# k8s.io/csi-api => github.com/kubeedge/kubernetes/staging/src/k8s.io/csi-api v1.22.17-kubeedge1
+# k8s.io/csi-translation-lib => github.com/kubeedge/kubernetes/staging/src/k8s.io/csi-translation-lib v1.22.17-kubeedge1
+# k8s.io/gengo v0.0.0 => k8s.io/gengo v0.22.17
# k8s.io/heapster => k8s.io/heapster v1.2.0-beta.1
# k8s.io/klog/v2 => k8s.io/klog/v2 v2.8.0
-# k8s.io/kube-aggregator => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-aggregator v1.22.6-kubeedge4
-# k8s.io/kube-controller-manager => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-controller-manager v1.22.6-kubeedge4
-# k8s.io/kube-proxy => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-proxy v1.22.6-kubeedge4
-# k8s.io/kube-scheduler => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-scheduler v1.22.6-kubeedge4
-# k8s.io/kubectl => github.com/kubeedge/kubernetes/staging/src/k8s.io/kubectl v1.22.6-kubeedge4
-# k8s.io/kubelet => github.com/kubeedge/kubernetes/staging/src/k8s.io/kubelet v1.22.6-kubeedge4
-# k8s.io/kubernetes => github.com/kubeedge/kubernetes v1.22.6-kubeedge4
-# k8s.io/legacy-cloud-providers => github.com/kubeedge/kubernetes/staging/src/k8s.io/legacy-cloud-providers v1.22.6-kubeedge4
-# k8s.io/metrics => github.com/kubeedge/kubernetes/staging/src/k8s.io/metrics v1.22.6-kubeedge4
-# k8s.io/mount-utils => github.com/kubeedge/kubernetes/staging/src/k8s.io/mount-utils v1.22.6-kubeedge4
-# k8s.io/node-api => github.com/kubeedge/kubernetes/staging/src/k8s.io/node-api v1.22.6-kubeedge4
-# k8s.io/pod-security-admission => k8s.io/pod-security-admission v0.22.6
-# k8s.io/repo-infra => github.com/kubeedge/kubernetes/staging/src/k8s.io/repo-infra v1.22.6-kubeedge4
-# k8s.io/sample-apiserver => github.com/kubeedge/kubernetes/staging/src/k8s.io/sample-apiserver v1.22.6-kubeedge4
-# k8s.io/utils v0.0.0 => k8s.io/utils v0.22.6
+# k8s.io/kube-aggregator => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-aggregator v1.22.17-kubeedge1
+# k8s.io/kube-controller-manager => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-controller-manager v1.22.17-kubeedge1
+# k8s.io/kube-proxy => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-proxy v1.22.17-kubeedge1
+# k8s.io/kube-scheduler => github.com/kubeedge/kubernetes/staging/src/k8s.io/kube-scheduler v1.22.17-kubeedge1
+# k8s.io/kubectl => github.com/kubeedge/kubernetes/staging/src/k8s.io/kubectl v1.22.17-kubeedge1
+# k8s.io/kubelet => github.com/kubeedge/kubernetes/staging/src/k8s.io/kubelet v1.22.17-kubeedge1
+# k8s.io/kubernetes => github.com/kubeedge/kubernetes v1.22.17-kubeedge1
+# k8s.io/legacy-cloud-providers => github.com/kubeedge/kubernetes/staging/src/k8s.io/legacy-cloud-providers v1.22.17-kubeedge1
+# k8s.io/metrics => github.com/kubeedge/kubernetes/staging/src/k8s.io/metrics v1.22.17-kubeedge1
+# k8s.io/mount-utils => github.com/kubeedge/kubernetes/staging/src/k8s.io/mount-utils v1.22.17-kubeedge1
+# k8s.io/node-api => github.com/kubeedge/kubernetes/staging/src/k8s.io/node-api v1.22.17-kubeedge1
+# k8s.io/pod-security-admission => k8s.io/pod-security-admission v0.22.17
+# k8s.io/repo-infra => github.com/kubeedge/kubernetes/staging/src/k8s.io/repo-infra v1.22.17-kubeedge1
+# k8s.io/sample-apiserver => github.com/kubeedge/kubernetes/staging/src/k8s.io/sample-apiserver v1.22.17-kubeedge1
+# k8s.io/utils v0.0.0 => k8s.io/utils v0.22.17
# sigs.k8s.io/apiserver-network-proxy/konnectivity-client => sigs.k8s.io/apiserver-network-proxy/konnectivity-client v0.0.27