diff options
| author | KubeEdge Bot <48982446+kubeedge-bot@users.noreply.github.com> | 2023-11-28 17:38:41 +0800 |
|---|---|---|
| committer | GitHub <noreply@github.com> | 2023-11-28 17:38:41 +0800 |
| commit | 38e537a9d8b8617598ec5562b89f5f571b1c285d (patch) | |
| tree | 68e3282f28630876f17361d7d03a7cfa6fc17f38 | |
| parent | Merge pull request #5210 from Onion-of-dreamed/automated-cherry-pick-of-#5107... (diff) | |
| parent | update version for k8s v1.22.17 (diff) | |
| download | kubeedge-1.12.6.tar.gz | |
Merge pull request #5214 from Shelley-BaoYue/bump-k8s-1.22.17v1.12.6
[release-1.12] Bump k8s to 1.22.17
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 } @@ -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 ) @@ -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 |
