diff --git a/pkg/checker/node_checker.go b/pkg/checker/node_checker.go index f0d01ee01..2b00c6032 100644 --- a/pkg/checker/node_checker.go +++ b/pkg/checker/node_checker.go @@ -15,10 +15,13 @@ package checker import ( + "context" "fmt" "os" "text/template" + v1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "github.com/fanux/sealos/pkg/client-go/kubernetes" v2 "github.com/fanux/sealos/pkg/types/v1beta1" "github.com/fanux/sealos/pkg/utils/logger" @@ -32,7 +35,6 @@ const ( ) type NodeChecker struct { - client *kubernetes.Client } type NodeClusterStatus struct { @@ -47,12 +49,11 @@ func (n *NodeChecker) Check(cluster *v2.Cluster, phase string) error { return nil } // checker if all the node is ready - c, err := kubernetes.Newk8sClient() + c, err := kubernetes.NewKubernetesClient("") if err != nil { return err } - n.client = c - nodes, err := n.client.ListNodes() + nodes, err := c.Kubernetes().CoreV1().Nodes().List(context.Background(), v1.ListOptions{}) if err != nil { return err } diff --git a/pkg/checker/pod_checker.go b/pkg/checker/pod_checker.go deleted file mode 100644 index a7dcf4005..000000000 --- a/pkg/checker/pod_checker.go +++ /dev/null @@ -1,127 +0,0 @@ -// Copyright © 2021 sealos. -// -// 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 checker - -import ( - "os" - "text/template" - - "github.com/fanux/sealos/pkg/client-go/kubernetes" - v2 "github.com/fanux/sealos/pkg/types/v1beta1" - "github.com/fanux/sealos/pkg/utils/logger" - - corev1 "k8s.io/api/core/v1" -) - -type PodChecker struct { - client *kubernetes.Client -} - -type PodNamespaceStatus struct { - NamespaceName string - RunningCount uint32 - NotRunningCount uint32 - PodCount uint32 - NotRunningPodList []*corev1.Pod -} - -var PodNamespaceStatusList []PodNamespaceStatus - -func (n *PodChecker) Check(cluster *v2.Cluster, phase string) error { - if phase != PhasePost { - return nil - } - c, err := kubernetes.Newk8sClient() - if err != nil { - return err - } - n.client = c - - namespacePodList, err := n.client.ListAllNamespacesPods() - if err != nil { - return err - } - for _, podNamespace := range namespacePodList { - var runningCount uint32 - var notRunningCount uint32 - var podCount uint32 - var notRunningPodList []*corev1.Pod - for _, pod := range podNamespace.PodList.Items { - if err := getPodReadyStatus(pod); err != nil { - notRunningCount++ - newPod := pod - notRunningPodList = append(notRunningPodList, &newPod) - } else { - runningCount++ - } - } - podCount = runningCount + notRunningCount - podNamespaceStatus := PodNamespaceStatus{ - NamespaceName: podNamespace.Namespace.Name, - RunningCount: runningCount, - NotRunningCount: notRunningCount, - PodCount: podCount, - NotRunningPodList: notRunningPodList, - } - PodNamespaceStatusList = append(PodNamespaceStatusList, podNamespaceStatus) - } - err = n.Output(PodNamespaceStatusList) - if err != nil { - return err - } - return nil -} - -func (n *PodChecker) Output(podNamespaceStatusList []PodNamespaceStatus) error { - t := template.New("pod_checker") - t, err := t.Parse( - `Cluster Pod Status - {{ range . -}} - Namespace: {{ .NamespaceName }} - RunningPod: {{ .RunningCount }}/{{ .PodCount }} - {{ if (gt .NotRunningCount 0) -}} - Not Running Pod List: - {{- range .NotRunningPodList }} - PodName: {{ .Name }} - {{- end }} - {{ end }} - {{- end }} -`) - if err != nil { - panic(err) - } - t = template.Must(t, err) - err = t.Execute(os.Stdout, podNamespaceStatusList) - if err != nil { - logger.Error("pod checkers template can not excute %s", err) - return err - } - return nil -} - -func getPodReadyStatus(pod corev1.Pod) error { - for _, condition := range pod.Status.Conditions { - if condition.Type == "Ready" { - if condition.Status == "True" { - return nil - } - } - } - return &NotFindReadyTypeError{} -} - -func NewPodChecker() Interface { - return &PodChecker{} -} diff --git a/pkg/checker/svc_checker.go b/pkg/checker/svc_checker.go deleted file mode 100644 index 4a9942b80..000000000 --- a/pkg/checker/svc_checker.go +++ /dev/null @@ -1,127 +0,0 @@ -// Copyright © 2021 sealos. -// -// 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 checker - -import ( - "os" - "text/template" - - "github.com/fanux/sealos/pkg/client-go/kubernetes" - v2 "github.com/fanux/sealos/pkg/types/v1beta1" - "github.com/fanux/sealos/pkg/utils/logger" - - corev1 "k8s.io/api/core/v1" -) - -type SvcChecker struct { - client *kubernetes.Client -} - -type SvcNamespaceStatus struct { - NamespaceName string - ServiceCount int - EndpointCount int - UnhealthServiceList []string -} - -type SvcClusterStatus struct { - SvcNamespaceStatusList []*SvcNamespaceStatus -} - -func (n *SvcChecker) Check(cluster *v2.Cluster, phase string) error { - if phase != PhasePost { - return nil - } - // checker if all the node is ready - c, err := kubernetes.Newk8sClient() - if err != nil { - return err - } - n.client = c - - namespaceSvcList, err := n.client.ListAllNamespacesSvcs() - var svcNamespaceStatusList []*SvcNamespaceStatus - if err != nil { - return err - } - for _, svcNamespace := range namespaceSvcList { - serviceCount := len(svcNamespace.ServiceList.Items) - var unhaelthService []string - var endpointCount = 0 - endpointsList, err := n.client.GetEndpointsList(svcNamespace.Namespace.Name) - if err != nil { - break - } - for _, service := range svcNamespace.ServiceList.Items { - if IsExistEndpoint(endpointsList, service.Name) { - endpointCount++ - } else { - unhaelthService = append(unhaelthService, service.Name) - } - } - svcNamespaceStatus := SvcNamespaceStatus{ - NamespaceName: svcNamespace.Namespace.Name, - ServiceCount: serviceCount, - EndpointCount: endpointCount, - UnhealthServiceList: unhaelthService, - } - svcNamespaceStatusList = append(svcNamespaceStatusList, &svcNamespaceStatus) - } - err = n.Output(svcNamespaceStatusList) - if err != nil { - return err - } - return nil -} - -func (n *SvcChecker) Output(svcNamespaceStatusList []*SvcNamespaceStatus) error { - t := template.New("svc_checker") - t, err := t.Parse( - `Cluster Service Status - {{- range . }} - Namespace: {{ .NamespaceName }} - HealthService: {{ .EndpointCount }}/{{ .ServiceCount }} - UnhealthyServiceList: - {{- range .UnhealthyServiceList }} - ServiceName: {{ . }} - {{- end }} - {{- end }} -`) - if err != nil { - panic(err) - } - t = template.Must(t, err) - err = t.Execute(os.Stdout, svcNamespaceStatusList) - if err != nil { - logger.Error("service checkers template can not execute %s", err) - return err - } - return nil -} - -func IsExistEndpoint(endpointList *corev1.EndpointsList, serviceName string) bool { - for _, ep := range endpointList.Items { - if ep.Name == serviceName { - if len(ep.Subsets) > 0 { - return true - } - } - } - return false -} - -func NewSvcChecker() Interface { - return &SvcChecker{} -} diff --git a/pkg/client-go/kubernetes/client.go b/pkg/client-go/kubernetes/client.go index de68fa346..a9374bd30 100644 --- a/pkg/client-go/kubernetes/client.go +++ b/pkg/client-go/kubernetes/client.go @@ -15,187 +15,80 @@ package kubernetes import ( - "context" + "path/filepath" - "github.com/fanux/sealos/pkg/utils/contants" - - "github.com/pkg/errors" - v1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/version" + "k8s.io/client-go/discovery" + "k8s.io/client-go/dynamic" "k8s.io/client-go/kubernetes" + "k8s.io/client-go/rest" "k8s.io/client-go/tools/clientcmd" + "k8s.io/client-go/util/homedir" ) -type Client struct { - client *kubernetes.Clientset +type Client interface { + Kubernetes() kubernetes.Interface + Discovery() discovery.DiscoveryInterface + KubernetesDynamic() dynamic.Interface + Config() *rest.Config } -type NamespacePod struct { - Namespace v1.Namespace - PodList *v1.PodList +type kubernetesClient struct { + // kubernetes client interface + k8s kubernetes.Interface + k8sDynamic dynamic.Interface + // discovery client + discoveryClient *discovery.DiscoveryClient + + config *rest.Config } -type NamespaceSvc struct { - Namespace v1.Namespace - ServiceList *v1.ServiceList -} - -func Newk8sClient() (*Client, error) { - kubeconfig := contants.DefaultKubeConfigFile() - // use the current context in kubeconfig - config, err := clientcmd.BuildConfigFromFlags("", kubeconfig) - if err != nil { - return nil, errors.Wrap(err, "failed to build kube config") - } - - clientSet, err := kubernetes.NewForConfig(config) - if err != nil { - return nil, err - } - - return &Client{ - client: clientSet, - }, nil -} - -func (c *Client) ListNodes() (*v1.NodeList, error) { - nodes, err := c.client.CoreV1().Nodes().List(context.TODO(), metav1.ListOptions{}) - if err != nil { - return nil, errors.Wrapf(err, "failed to get cluster nodes") - } - return nodes, nil -} - -func (c *Client) UpdateNode(node v1.Node) (*v1.Node, error) { - n, err := c.client.CoreV1().Nodes().Update(context.TODO(), &node, metav1.UpdateOptions{}) - if err != nil { - return nil, errors.Wrapf(err, "failed to update cluster node") - } - return n, nil -} - -func (c *Client) DeleteNode(name string) error { - if err := c.client.CoreV1().Nodes().Delete(context.TODO(), name, metav1.DeleteOptions{}); err != nil { - return errors.Wrapf(err, "failed to delete cluster node %s", name) - } - return nil -} - -func (c *Client) listNamespaces() (*v1.NamespaceList, error) { - namespaceList, err := c.client.CoreV1().Namespaces().List(context.TODO(), metav1.ListOptions{}) - if err != nil { - return nil, errors.Wrapf(err, "failed to get namespaces") - } - return namespaceList, nil -} - -func (c *Client) ListNodesByLabel(label string) (*v1.NodeList, error) { - nodes, err := c.client.CoreV1().Nodes().List(context.TODO(), metav1.ListOptions{LabelSelector: label}) - if err != nil { - return nil, errors.Wrapf(err, "failed to get cluster nodes") - } - return nodes, nil -} - -func (c *Client) ListNodeIPByLabel(label string) ([]string, error) { - var ips []string - nodes, err := c.ListNodesByLabel(label) - if err != nil { - return nil, err - } - for _, node := range nodes.Items { - for _, v := range node.Status.Addresses { - if v.Type == v1.NodeInternalIP { - ips = append(ips, v.Address) +// NewKubernetesClient creates a KubernetesClient +func NewKubernetesClient(kubeconfig string) (Client, error) { + config := new(rest.Config) + var err error + if config == nil { + config, err = rest.InClusterConfig() + if err != nil { + if kubeconfig == "" { + kubeconfig = filepath.Join(homedir.HomeDir(), ".kube", "config") + } + config, err = clientcmd.BuildConfigFromFlags("", kubeconfig) + if err != nil { + return nil, err } } } - return ips, nil -} - -func (c *Client) ListAllNamespacesPods() ([]*NamespacePod, error) { - namespaceList, err := c.listNamespaces() + config.QPS = 1e6 + config.Burst = 1e6 + var k kubernetesClient + k.k8s, err = kubernetes.NewForConfig(config) if err != nil { return nil, err } - var namespacePodList []*NamespacePod - for _, ns := range namespaceList.Items { - pods, err := c.client.CoreV1().Pods(ns.Name).List(context.TODO(), metav1.ListOptions{}) - if err != nil { - return nil, errors.Wrapf(err, "failed to get all namespace pods") - } - namespacePod := NamespacePod{ - Namespace: ns, - PodList: pods, - } - namespacePodList = append(namespacePodList, &namespacePod) - } - - return namespacePodList, nil -} - -func (c *Client) ListAllNamespacesSvcs() ([]*NamespaceSvc, error) { - namespaceList, err := c.listNamespaces() + k.discoveryClient, err = discovery.NewDiscoveryClientForConfig(config) if err != nil { return nil, err } - var namespaceSvcList []*NamespaceSvc - for _, ns := range namespaceList.Items { - svcs, err := c.client.CoreV1().Services(ns.Name).List(context.TODO(), metav1.ListOptions{}) - if err != nil { - return nil, errors.Wrapf(err, "failed to get all namespace pods") - } - namespaceSvc := NamespaceSvc{ - Namespace: ns, - ServiceList: svcs, - } - namespaceSvcList = append(namespaceSvcList, &namespaceSvc) - } - return namespaceSvcList, nil -} - -func (c *Client) GetEndpointsList(namespace string) (*v1.EndpointsList, error) { - endpointsList, err := c.client.CoreV1().Endpoints(namespace).List(context.TODO(), metav1.ListOptions{}) - if err != nil { - return nil, errors.Wrapf(err, "failed to get endpoint in namespace %s", namespace) - } - return endpointsList, nil -} - -func (c *Client) ListSvcs(namespace string) (*v1.ServiceList, error) { - svcs, err := c.client.CoreV1().Services(namespace).List(context.TODO(), metav1.ListOptions{}) - if err != nil { - return nil, errors.Wrapf(err, "failed to get all namespace pods") - } - return svcs, nil -} - -func (c *Client) GetClusterVersion() (*version.Info, error) { - info, err := c.client.Discovery().ServerVersion() + k.k8sDynamic, err = dynamic.NewForConfig(config) if err != nil { return nil, err } - return info, nil + k.config = config + return &k, nil } -func (c *Client) ListKubeSystemPodsStatus() (bool, error) { - pods, err := c.client.CoreV1().Pods("kube-system").List(context.TODO(), metav1.ListOptions{}) - if err != nil { - return false, errors.Wrapf(err, "failed to get kube-system namespace pods") - } - // pods.Items maybe nil - if len(pods.Items) == 0 { - return false, nil - } - for _, pod := range pods.Items { - // pod.Status.ContainerStatus == nil because of pod contain initcontainer - if len(pod.Status.ContainerStatuses) == 0 { - continue - } - if !pod.Status.ContainerStatuses[0].Ready { - return false, nil - } - } - return true, nil +func (k *kubernetesClient) Kubernetes() kubernetes.Interface { + return k.k8s +} + +func (k *kubernetesClient) Discovery() discovery.DiscoveryInterface { + return k.discoveryClient +} + +func (k *kubernetesClient) Config() *rest.Config { + return k.config +} + +func (k *kubernetesClient) KubernetesDynamic() dynamic.Interface { + return k.k8sDynamic } diff --git a/pkg/client-go/kubernetes/expansion.go b/pkg/client-go/kubernetes/expansion.go new file mode 100644 index 000000000..762aee5d0 --- /dev/null +++ b/pkg/client-go/kubernetes/expansion.go @@ -0,0 +1,37 @@ +/* +Copyright 2022 cuisongliu@qq.com. + +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 kubernetes + +import ( + "context" + "fmt" + + v1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + clientset "k8s.io/client-go/kubernetes" +) + +var ( + // ControlPlaneComponents defines the control-plane component names + ControlPlaneComponents = []string{KubeAPIServer, KubeControllerManager, KubeScheduler} +) + +// GetStaticPod computes hashes for a single Static Pod resource +func GetStaticPod(client clientset.Interface, nodeName string, component string) (*v1.Pod, error) { + staticPodName := fmt.Sprintf("%s-%s", component, nodeName) + return client.CoreV1().Pods(metav1.NamespaceSystem).Get(context.TODO(), staticPodName, metav1.GetOptions{}) +} diff --git a/pkg/client-go/kubernetes/healthy.go b/pkg/client-go/kubernetes/healthy.go new file mode 100644 index 000000000..7d50b2369 --- /dev/null +++ b/pkg/client-go/kubernetes/healthy.go @@ -0,0 +1,229 @@ +/* +Copyright 2018 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 kubernetes + +import ( + "context" + "fmt" + "net/http" + "path" + "time" + + v1 "k8s.io/api/core/v1" + + "github.com/fanux/sealos/pkg/utils/logger" + + "github.com/pkg/errors" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + netutil "k8s.io/apimachinery/pkg/util/net" + "k8s.io/apimachinery/pkg/util/wait" + clientset "k8s.io/client-go/kubernetes" +) + +// Healthy is an interface for waiting for criteria in Kubernetes to happen +type Healthy interface { + // ForAPI waits for the API Server's /healthz endpoint to become "ok" + ForAPI() error + // ForHealthyKubelet blocks until the kubelet /healthz endpoint returns 'ok' + ForHealthyKubelet(initialTimeout time.Duration, host string) error + // SetTimeout adjusts the timeout to the specified duration + SetTimeout(timeout time.Duration) + // ForHealthyPod pod status + ForHealthyPod(pod *v1.Pod) string +} + +// KubeHealthy is an implementation of Healthy that is backed by a Kubernetes client +type kubeHealthy struct { + client clientset.Interface + timeout time.Duration +} + +// NewKubeHealthy returns a new Healthy object that talks to the given Kubernetes cluster +func NewKubeHealthy(client clientset.Interface, timeout time.Duration) Healthy { + return &kubeHealthy{ + client: client, + timeout: timeout, + } +} + +// ForAPI waits for the API Server's /healthz endpoint to report "ok" +func (w *kubeHealthy) ForAPI() error { + start := time.Now() + return wait.PollImmediate(APICallRetryInterval, w.timeout, func() (bool, error) { + healthStatus := 0 + w.client.Discovery().RESTClient().Get().AbsPath("/healthz").Do(context.TODO()).StatusCode(&healthStatus) + if healthStatus != http.StatusOK { + return false, nil + } + + logger.Debug("[apiclient] All control plane components are healthy after %f seconds\n", time.Since(start).Seconds()) + return true, nil + }) +} + +// ForHealthyKubelet blocks until the kubelet /healthz endpoint returns 'ok',port is 10248 +func (w *kubeHealthy) ForHealthyKubelet(initialTimeout time.Duration, host string) error { + time.Sleep(initialTimeout) + logger.Debug("[kubelet-check] Initial timeout of %v passed.\n", initialTimeout) + return tryRunCommand(func() error { + trans := netutil.SetOldTransportDefaults(&http.Transport{}) + trans.TLSClientConfig.InsecureSkipVerify = true + client := &http.Client{Transport: trans} + + healthzEndpoint := path.Join(fmt.Sprintf("https://%s:%d", host, KubeletHealthzPort), "healthz") + resp, err := client.Get(healthzEndpoint) + if err != nil { + logger.Warn("[kubelet-check] It seems like the kubelet isn't running or healthy.") + logger.Warn("[kubelet-check] The HTTP call equal to 'curl -sSL %s' failed with error: %v.\n", healthzEndpoint, err) + return err + } + defer resp.Body.Close() + if resp.StatusCode != http.StatusOK { + logger.Warn("[kubelet-check] It seems like the kubelet isn't running or healthy.") + logger.Warn("[kubelet-check] The HTTP call equal to 'curl -sSL %s' returned HTTP code %d\n", healthzEndpoint, resp.StatusCode) + return errors.New("the kubelet healthz endpoint is unhealthy") + } + return nil + }, 5) // a failureThreshold of five means waiting for a total of 155 seconds +} + +// SetTimeout adjusts the timeout to the specified duration +func (w *kubeHealthy) SetTimeout(timeout time.Duration) { + w.timeout = timeout +} + +func (w *kubeHealthy) ForHealthyPod(pod *v1.Pod) string { + return printPod(pod) +} + +func printPod(pod *v1.Pod) string { + restarts := 0 + readyContainers := 0 + lastRestartDate := metav1.NewTime(time.Time{}) + + reason := string(pod.Status.Phase) + if pod.Status.Reason != "" { + reason = pod.Status.Reason + } + + initializing := false + for i := range pod.Status.InitContainerStatuses { + container := pod.Status.InitContainerStatuses[i] + restarts += int(container.RestartCount) + if container.LastTerminationState.Terminated != nil { + terminatedDate := container.LastTerminationState.Terminated.FinishedAt + if lastRestartDate.Before(&terminatedDate) { + lastRestartDate = terminatedDate + } + } + switch { + case container.State.Terminated != nil && container.State.Terminated.ExitCode == 0: + continue + case container.State.Terminated != nil: + // initialization is failed + if len(container.State.Terminated.Reason) == 0 { + if container.State.Terminated.Signal != 0 { + reason = fmt.Sprintf("Init:Signal:%d", container.State.Terminated.Signal) + } else { + reason = fmt.Sprintf("Init:ExitCode:%d", container.State.Terminated.ExitCode) + } + } else { + reason = "Init:" + container.State.Terminated.Reason + } + initializing = true + case container.State.Waiting != nil && len(container.State.Waiting.Reason) > 0 && container.State.Waiting.Reason != "PodInitializing": + reason = "Init:" + container.State.Waiting.Reason + initializing = true + default: + reason = fmt.Sprintf("Init:%d/%d", i, len(pod.Spec.InitContainers)) + initializing = true + } + break + } + if !initializing { + restarts = 0 + hasRunning := false + for i := len(pod.Status.ContainerStatuses) - 1; i >= 0; i-- { + container := pod.Status.ContainerStatuses[i] + + restarts += int(container.RestartCount) + if container.LastTerminationState.Terminated != nil { + terminatedDate := container.LastTerminationState.Terminated.FinishedAt + if lastRestartDate.Before(&terminatedDate) { + lastRestartDate = terminatedDate + } + } + if container.State.Waiting != nil && container.State.Waiting.Reason != "" { + reason = container.State.Waiting.Reason + } else if container.State.Terminated != nil && container.State.Terminated.Reason != "" { + reason = container.State.Terminated.Reason + } else if container.State.Terminated != nil && container.State.Terminated.Reason == "" { + if container.State.Terminated.Signal != 0 { + reason = fmt.Sprintf("Signal:%d", container.State.Terminated.Signal) + } else { + reason = fmt.Sprintf("ExitCode:%d", container.State.Terminated.ExitCode) + } + } else if container.Ready && container.State.Running != nil { + hasRunning = true + readyContainers++ + } + } + + // change pod status back to "Running" if there is at least one container still reporting as "Running" status + if reason == "Completed" && hasRunning { + if func() bool { + for _, condition := range pod.Status.Conditions { + if condition.Type == v1.PodReady && condition.Status == v1.ConditionTrue { + return true + } + } + return false + }() { + reason = "Running" + } else { + reason = "NotReady" + } + } + } + + if pod.DeletionTimestamp != nil && pod.Status.Reason == "NodeLost" { + reason = "Unknown" + } else if pod.DeletionTimestamp != nil { + reason = "Terminating" + } + + return reason +} + +// TryRunCommand runs a function a maximum of failureThreshold times, and retries on error. If failureThreshold is hit; the last error is returned +func tryRunCommand(f func() error, failureThreshold int) error { + backoff := wait.Backoff{ + Duration: 5 * time.Second, + Factor: 2, // double the timeout for every failure + Steps: failureThreshold, + } + return wait.ExponentialBackoff(backoff, func() (bool, error) { + err := f() + if err != nil { + // Retry until the timeout + return false, nil //nolint:nilerr + } + // The last f() call was a success, return cleanly + return true, nil + }) +} diff --git a/pkg/client-go/kubernetes/idempotency.go b/pkg/client-go/kubernetes/idempotency.go index 84fe8b2b4..6a437862f 100644 --- a/pkg/client-go/kubernetes/idempotency.go +++ b/pkg/client-go/kubernetes/idempotency.go @@ -34,7 +34,6 @@ import ( "k8s.io/apimachinery/pkg/util/strategicpatch" "k8s.io/apimachinery/pkg/util/wait" clientset "k8s.io/client-go/kubernetes" - clientsetretry "k8s.io/client-go/util/retry" ) const ( @@ -51,104 +50,54 @@ const ( KubeScheduler = "kube-scheduler" ) -var ( - // ControlPlaneComponents defines the control-plane component names - ControlPlaneComponents = []string{KubeAPIServer, KubeControllerManager, KubeScheduler} -) +type Idempotency interface { + CreateOrUpdateConfigMap(cm *v1.ConfigMap) error + CreateOrUpdateSecret(secret *v1.Secret) error + CreateOrUpdateServiceAccount(sa *v1.ServiceAccount) error + CreateOrUpdateDeployment(deploy *apps.Deployment) error + CreateOrUpdateDaemonSet(ds *apps.DaemonSet) error + DeleteDaemonSetForeground(namespace, name string) error + DeleteDeploymentForeground(namespace, name string) error + CreateOrUpdateRole(role *rbac.Role) error + CreateOrUpdateRoleBinding(roleBinding *rbac.RoleBinding) error + CreateOrUpdateClusterRole(clusterRole *rbac.ClusterRole) error + CreateOrUpdateClusterRoleBinding(clusterRoleBinding *rbac.ClusterRoleBinding) error + PatchNode(nodeName string, patchFn func(*v1.Node)) error +} -// ConfigMapMutator is a function that mutates the given ConfigMap and optionally returns an error -type ConfigMapMutator func(*v1.ConfigMap) error +type kubeIdempotency struct { + client clientset.Interface +} -// TODO: We should invent a dynamic mechanism for this using the dynamic client instead of hard-coding these functions per-type +// NewKubeIdempotency returns a new Idempotency object that talks to the given Kubernetes cluster +func NewKubeIdempotency(client clientset.Interface) Idempotency { + return &kubeIdempotency{ + client: client, + } +} // CreateOrUpdateConfigMap creates a ConfigMap if the target resource doesn't exist. If the resource exists already, this function will update the resource instead. -func CreateOrUpdateConfigMap(client clientset.Interface, cm *v1.ConfigMap) error { - if _, err := client.CoreV1().ConfigMaps(cm.ObjectMeta.Namespace).Create(context.TODO(), cm, metav1.CreateOptions{}); err != nil { +func (ki *kubeIdempotency) CreateOrUpdateConfigMap(cm *v1.ConfigMap) error { + if _, err := ki.client.CoreV1().ConfigMaps(cm.ObjectMeta.Namespace).Create(context.TODO(), cm, metav1.CreateOptions{}); err != nil { if !apierrors.IsAlreadyExists(err) { return errors.Wrap(err, "unable to create ConfigMap") } - if _, err := client.CoreV1().ConfigMaps(cm.ObjectMeta.Namespace).Update(context.TODO(), cm, metav1.UpdateOptions{}); err != nil { + if _, err := ki.client.CoreV1().ConfigMaps(cm.ObjectMeta.Namespace).Update(context.TODO(), cm, metav1.UpdateOptions{}); err != nil { return errors.Wrap(err, "unable to update ConfigMap") } } return nil } -// CreateOrMutateConfigMap tries to create the ConfigMap provided as cm. If the resource exists already, the latest version will be fetched from -// the cluster and mutator callback will be called on it, then an Update of the mutated ConfigMap will be performed. This function is resilient -// to conflicts, and a retry will be issued if the ConfigMap was modified on the server between the refresh and the update (while the mutation was -// taking place) -func CreateOrMutateConfigMap(client clientset.Interface, cm *v1.ConfigMap, mutator ConfigMapMutator) error { - var lastError error - err := wait.ExponentialBackoff(wait.Backoff{ - Steps: 20, - Duration: 500 * time.Millisecond, - Factor: 1.0, - Jitter: 0.1, - }, func() (bool, error) { - if _, err := client.CoreV1().ConfigMaps(cm.ObjectMeta.Namespace).Create(context.TODO(), cm, metav1.CreateOptions{}); err != nil { - lastError = err - if apierrors.IsAlreadyExists(err) { - lastError = MutateConfigMap(client, metav1.ObjectMeta{Namespace: cm.ObjectMeta.Namespace, Name: cm.ObjectMeta.Name}, mutator) - return lastError == nil, nil - } - return false, nil - } - return true, nil - }) - if err == nil { - return nil - } - return lastError -} - -// MutateConfigMap takes a ConfigMap Object Meta (namespace and name), retrieves the resource from the server and tries to mutate it -// by calling to the mutator callback, then an Update of the mutated ConfigMap will be performed. This function is resilient -// to conflicts, and a retry will be issued if the ConfigMap was modified on the server between the refresh and the update (while the mutation was -// taking place). -func MutateConfigMap(client clientset.Interface, meta metav1.ObjectMeta, mutator ConfigMapMutator) error { - return clientsetretry.RetryOnConflict(wait.Backoff{ - Steps: 20, - Duration: 500 * time.Millisecond, - Factor: 1.0, - Jitter: 0.1, - }, func() error { - configMap, err := client.CoreV1().ConfigMaps(meta.Namespace).Get(context.TODO(), meta.Name, metav1.GetOptions{}) - if err != nil { - return err - } - if err = mutator(configMap); err != nil { - return errors.Wrap(err, "unable to mutate ConfigMap") - } - _, err = client.CoreV1().ConfigMaps(configMap.ObjectMeta.Namespace).Update(context.TODO(), configMap, metav1.UpdateOptions{}) - return err - }) -} - -// CreateOrRetainConfigMap creates a ConfigMap if the target resource doesn't exist. If the resource exists already, this function will retain the resource instead. -func CreateOrRetainConfigMap(client clientset.Interface, cm *v1.ConfigMap, configMapName string) error { - if _, err := client.CoreV1().ConfigMaps(cm.ObjectMeta.Namespace).Get(context.TODO(), configMapName, metav1.GetOptions{}); err != nil { - if !apierrors.IsNotFound(err) { - return nil - } - if _, err := client.CoreV1().ConfigMaps(cm.ObjectMeta.Namespace).Create(context.TODO(), cm, metav1.CreateOptions{}); err != nil { - if !apierrors.IsAlreadyExists(err) { - return errors.Wrap(err, "unable to create ConfigMap") - } - } - } - return nil -} - // CreateOrUpdateSecret creates a Secret if the target resource doesn't exist. If the resource exists already, this function will update the resource instead. -func CreateOrUpdateSecret(client clientset.Interface, secret *v1.Secret) error { - if _, err := client.CoreV1().Secrets(secret.ObjectMeta.Namespace).Create(context.TODO(), secret, metav1.CreateOptions{}); err != nil { +func (ki *kubeIdempotency) CreateOrUpdateSecret(secret *v1.Secret) error { + if _, err := ki.client.CoreV1().Secrets(secret.ObjectMeta.Namespace).Create(context.TODO(), secret, metav1.CreateOptions{}); err != nil { if !apierrors.IsAlreadyExists(err) { return errors.Wrap(err, "unable to create secret") } - if _, err := client.CoreV1().Secrets(secret.ObjectMeta.Namespace).Update(context.TODO(), secret, metav1.UpdateOptions{}); err != nil { + if _, err := ki.client.CoreV1().Secrets(secret.ObjectMeta.Namespace).Update(context.TODO(), secret, metav1.UpdateOptions{}); err != nil { return errors.Wrap(err, "unable to update secret") } } @@ -156,8 +105,8 @@ func CreateOrUpdateSecret(client clientset.Interface, secret *v1.Secret) error { } // CreateOrUpdateServiceAccount creates a ServiceAccount if the target resource doesn't exist. If the resource exists already, this function will update the resource instead. -func CreateOrUpdateServiceAccount(client clientset.Interface, sa *v1.ServiceAccount) error { - if _, err := client.CoreV1().ServiceAccounts(sa.ObjectMeta.Namespace).Create(context.TODO(), sa, metav1.CreateOptions{}); err != nil { +func (ki *kubeIdempotency) CreateOrUpdateServiceAccount(sa *v1.ServiceAccount) error { + if _, err := ki.client.CoreV1().ServiceAccounts(sa.ObjectMeta.Namespace).Create(context.TODO(), sa, metav1.CreateOptions{}); err != nil { // Note: We don't run .Update here afterwards as that's probably not required // Only thing that could be updated is annotations/labels in .metadata, but we don't use that currently if !apierrors.IsAlreadyExists(err) { @@ -168,42 +117,27 @@ func CreateOrUpdateServiceAccount(client clientset.Interface, sa *v1.ServiceAcco } // CreateOrUpdateDeployment creates a Deployment if the target resource doesn't exist. If the resource exists already, this function will update the resource instead. -func CreateOrUpdateDeployment(client clientset.Interface, deploy *apps.Deployment) error { - if _, err := client.AppsV1().Deployments(deploy.ObjectMeta.Namespace).Create(context.TODO(), deploy, metav1.CreateOptions{}); err != nil { +func (ki *kubeIdempotency) CreateOrUpdateDeployment(deploy *apps.Deployment) error { + if _, err := ki.client.AppsV1().Deployments(deploy.ObjectMeta.Namespace).Create(context.TODO(), deploy, metav1.CreateOptions{}); err != nil { if !apierrors.IsAlreadyExists(err) { return errors.Wrap(err, "unable to create deployment") } - if _, err := client.AppsV1().Deployments(deploy.ObjectMeta.Namespace).Update(context.TODO(), deploy, metav1.UpdateOptions{}); err != nil { + if _, err := ki.client.AppsV1().Deployments(deploy.ObjectMeta.Namespace).Update(context.TODO(), deploy, metav1.UpdateOptions{}); err != nil { return errors.Wrap(err, "unable to update deployment") } } return nil } -// CreateOrRetainDeployment creates a Deployment if the target resource doesn't exist. If the resource exists already, this function will retain the resource instead. -func CreateOrRetainDeployment(client clientset.Interface, deploy *apps.Deployment, deployName string) error { - if _, err := client.AppsV1().Deployments(deploy.ObjectMeta.Namespace).Get(context.TODO(), deployName, metav1.GetOptions{}); err != nil { - if !apierrors.IsNotFound(err) { - return nil - } - if _, err := client.AppsV1().Deployments(deploy.ObjectMeta.Namespace).Create(context.TODO(), deploy, metav1.CreateOptions{}); err != nil { - if !apierrors.IsAlreadyExists(err) { - return errors.Wrap(err, "unable to create deployment") - } - } - } - return nil -} - // CreateOrUpdateDaemonSet creates a DaemonSet if the target resource doesn't exist. If the resource exists already, this function will update the resource instead. -func CreateOrUpdateDaemonSet(client clientset.Interface, ds *apps.DaemonSet) error { - if _, err := client.AppsV1().DaemonSets(ds.ObjectMeta.Namespace).Create(context.TODO(), ds, metav1.CreateOptions{}); err != nil { +func (ki *kubeIdempotency) CreateOrUpdateDaemonSet(ds *apps.DaemonSet) error { + if _, err := ki.client.AppsV1().DaemonSets(ds.ObjectMeta.Namespace).Create(context.TODO(), ds, metav1.CreateOptions{}); err != nil { if !apierrors.IsAlreadyExists(err) { return errors.Wrap(err, "unable to create daemonset") } - if _, err := client.AppsV1().DaemonSets(ds.ObjectMeta.Namespace).Update(context.TODO(), ds, metav1.UpdateOptions{}); err != nil { + if _, err := ki.client.AppsV1().DaemonSets(ds.ObjectMeta.Namespace).Update(context.TODO(), ds, metav1.UpdateOptions{}); err != nil { return errors.Wrap(err, "unable to update daemonset") } } @@ -211,25 +145,25 @@ func CreateOrUpdateDaemonSet(client clientset.Interface, ds *apps.DaemonSet) err } // DeleteDaemonSetForeground deletes the specified DaemonSet in foreground mode; i.e. it blocks until/makes sure all the managed Pods are deleted -func DeleteDaemonSetForeground(client clientset.Interface, namespace, name string) error { +func (ki *kubeIdempotency) DeleteDaemonSetForeground(namespace, name string) error { foregroundDelete := metav1.DeletePropagationForeground - return client.AppsV1().DaemonSets(namespace).Delete(context.TODO(), name, metav1.DeleteOptions{PropagationPolicy: &foregroundDelete}) + return ki.client.AppsV1().DaemonSets(namespace).Delete(context.TODO(), name, metav1.DeleteOptions{PropagationPolicy: &foregroundDelete}) } // DeleteDeploymentForeground deletes the specified Deployment in foreground mode; i.e. it blocks until/makes sure all the managed Pods are deleted -func DeleteDeploymentForeground(client clientset.Interface, namespace, name string) error { +func (ki *kubeIdempotency) DeleteDeploymentForeground(namespace, name string) error { foregroundDelete := metav1.DeletePropagationForeground - return client.AppsV1().Deployments(namespace).Delete(context.TODO(), name, metav1.DeleteOptions{PropagationPolicy: &foregroundDelete}) + return ki.client.AppsV1().Deployments(namespace).Delete(context.TODO(), name, metav1.DeleteOptions{PropagationPolicy: &foregroundDelete}) } // CreateOrUpdateRole creates a Role if the target resource doesn't exist. If the resource exists already, this function will update the resource instead. -func CreateOrUpdateRole(client clientset.Interface, role *rbac.Role) error { - if _, err := client.RbacV1().Roles(role.ObjectMeta.Namespace).Create(context.TODO(), role, metav1.CreateOptions{}); err != nil { +func (ki *kubeIdempotency) CreateOrUpdateRole(role *rbac.Role) error { + if _, err := ki.client.RbacV1().Roles(role.ObjectMeta.Namespace).Create(context.TODO(), role, metav1.CreateOptions{}); err != nil { if !apierrors.IsAlreadyExists(err) { return errors.Wrap(err, "unable to create RBAC role") } - if _, err := client.RbacV1().Roles(role.ObjectMeta.Namespace).Update(context.TODO(), role, metav1.UpdateOptions{}); err != nil { + if _, err := ki.client.RbacV1().Roles(role.ObjectMeta.Namespace).Update(context.TODO(), role, metav1.UpdateOptions{}); err != nil { return errors.Wrap(err, "unable to update RBAC role") } } @@ -237,13 +171,13 @@ func CreateOrUpdateRole(client clientset.Interface, role *rbac.Role) error { } // CreateOrUpdateRoleBinding creates a RoleBinding if the target resource doesn't exist. If the resource exists already, this function will update the resource instead. -func CreateOrUpdateRoleBinding(client clientset.Interface, roleBinding *rbac.RoleBinding) error { - if _, err := client.RbacV1().RoleBindings(roleBinding.ObjectMeta.Namespace).Create(context.TODO(), roleBinding, metav1.CreateOptions{}); err != nil { +func (ki *kubeIdempotency) CreateOrUpdateRoleBinding(roleBinding *rbac.RoleBinding) error { + if _, err := ki.client.RbacV1().RoleBindings(roleBinding.ObjectMeta.Namespace).Create(context.TODO(), roleBinding, metav1.CreateOptions{}); err != nil { if !apierrors.IsAlreadyExists(err) { return errors.Wrap(err, "unable to create RBAC rolebinding") } - if _, err := client.RbacV1().RoleBindings(roleBinding.ObjectMeta.Namespace).Update(context.TODO(), roleBinding, metav1.UpdateOptions{}); err != nil { + if _, err := ki.client.RbacV1().RoleBindings(roleBinding.ObjectMeta.Namespace).Update(context.TODO(), roleBinding, metav1.UpdateOptions{}); err != nil { return errors.Wrap(err, "unable to update RBAC rolebinding") } } @@ -251,13 +185,13 @@ func CreateOrUpdateRoleBinding(client clientset.Interface, roleBinding *rbac.Rol } // CreateOrUpdateClusterRole creates a ClusterRole if the target resource doesn't exist. If the resource exists already, this function will update the resource instead. -func CreateOrUpdateClusterRole(client clientset.Interface, clusterRole *rbac.ClusterRole) error { - if _, err := client.RbacV1().ClusterRoles().Create(context.TODO(), clusterRole, metav1.CreateOptions{}); err != nil { +func (ki *kubeIdempotency) CreateOrUpdateClusterRole(clusterRole *rbac.ClusterRole) error { + if _, err := ki.client.RbacV1().ClusterRoles().Create(context.TODO(), clusterRole, metav1.CreateOptions{}); err != nil { if !apierrors.IsAlreadyExists(err) { return errors.Wrap(err, "unable to create RBAC clusterrole") } - if _, err := client.RbacV1().ClusterRoles().Update(context.TODO(), clusterRole, metav1.UpdateOptions{}); err != nil { + if _, err := ki.client.RbacV1().ClusterRoles().Update(context.TODO(), clusterRole, metav1.UpdateOptions{}); err != nil { return errors.Wrap(err, "unable to update RBAC clusterrole") } } @@ -265,27 +199,27 @@ func CreateOrUpdateClusterRole(client clientset.Interface, clusterRole *rbac.Clu } // CreateOrUpdateClusterRoleBinding creates a ClusterRoleBinding if the target resource doesn't exist. If the resource exists already, this function will update the resource instead. -func CreateOrUpdateClusterRoleBinding(client clientset.Interface, clusterRoleBinding *rbac.ClusterRoleBinding) error { - if _, err := client.RbacV1().ClusterRoleBindings().Create(context.TODO(), clusterRoleBinding, metav1.CreateOptions{}); err != nil { +func (ki *kubeIdempotency) CreateOrUpdateClusterRoleBinding(clusterRoleBinding *rbac.ClusterRoleBinding) error { + if _, err := ki.client.RbacV1().ClusterRoleBindings().Create(context.TODO(), clusterRoleBinding, metav1.CreateOptions{}); err != nil { if !apierrors.IsAlreadyExists(err) { return errors.Wrap(err, "unable to create RBAC clusterrolebinding") } - if _, err := client.RbacV1().ClusterRoleBindings().Update(context.TODO(), clusterRoleBinding, metav1.UpdateOptions{}); err != nil { + if _, err := ki.client.RbacV1().ClusterRoleBindings().Update(context.TODO(), clusterRoleBinding, metav1.UpdateOptions{}); err != nil { return errors.Wrap(err, "unable to update RBAC clusterrolebinding") } } return nil } -// PatchNodeOnce executes patchFn on the node object found by the node name. +// patchNodeOnce executes patchFn on the node object found by the node name. // This is a condition function meant to be used with wait.Poll. false, nil // implies it is safe to try again, an error indicates no more tries should be // made and true indicates success. -func PatchNodeOnce(client clientset.Interface, nodeName string, patchFn func(*v1.Node)) func() (bool, error) { +func (ki *kubeIdempotency) patchNodeOnce(nodeName string, patchFn func(*v1.Node)) func() (bool, error) { return func() (bool, error) { // First get the node object - n, err := client.CoreV1().Nodes().Get(context.TODO(), nodeName, metav1.GetOptions{}) + n, err := ki.client.CoreV1().Nodes().Get(context.TODO(), nodeName, metav1.GetOptions{}) if err != nil { // TODO this should only be for timeouts return false, nil //nolint:nilerr @@ -315,7 +249,7 @@ func PatchNodeOnce(client clientset.Interface, nodeName string, patchFn func(*v1 return false, errors.Wrap(err, "failed to create two way merge patch") } - if _, err := client.CoreV1().Nodes().Patch(context.TODO(), n.Name, types.StrategicMergePatchType, patchBytes, metav1.PatchOptions{}); err != nil { + if _, err := ki.client.CoreV1().Nodes().Patch(context.TODO(), n.Name, types.StrategicMergePatchType, patchBytes, metav1.PatchOptions{}); err != nil { // TODO also check for timeouts if apierrors.IsConflict(err) { logger.Debug("Temporarily unable to update node metadata due to conflict (will retry)") @@ -330,31 +264,9 @@ func PatchNodeOnce(client clientset.Interface, nodeName string, patchFn func(*v1 // PatchNode tries to patch a node using patchFn for the actual mutating logic. // Retries are provided by the wait package. -func PatchNode(client clientset.Interface, nodeName string, patchFn func(*v1.Node)) error { +func (ki *kubeIdempotency) PatchNode(nodeName string, patchFn func(*v1.Node)) error { // wait.Poll will rerun the condition function every interval function if // the function returns false. If the condition function returns an error // then the retries end and the error is returned. - return wait.Poll(APICallRetryInterval, PatchNodeTimeout, PatchNodeOnce(client, nodeName, patchFn)) -} - -// GetConfigMapWithRetry tries to retrieve a ConfigMap using the given client, -// retrying if we get an unexpected error. -// -// TODO: evaluate if this can be done better. Potentially remove the retry if feasible. -func GetConfigMapWithRetry(client clientset.Interface, namespace, name string) (*v1.ConfigMap, error) { - var cm *v1.ConfigMap - var lastError error - err := wait.ExponentialBackoff(clientsetretry.DefaultBackoff, func() (bool, error) { - var err error - cm, err = client.CoreV1().ConfigMaps(namespace).Get(context.TODO(), name, metav1.GetOptions{}) - if err == nil { - return true, nil - } - lastError = err - return false, nil - }) - if err == nil { - return cm, nil - } - return nil, lastError + return wait.Poll(APICallRetryInterval, PatchNodeTimeout, ki.patchNodeOnce(nodeName, patchFn)) } diff --git a/pkg/client-go/kubernetes/wait.go b/pkg/client-go/kubernetes/wait.go deleted file mode 100644 index cd683efda..000000000 --- a/pkg/client-go/kubernetes/wait.go +++ /dev/null @@ -1,243 +0,0 @@ -/* -Copyright 2018 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 kubernetes - -import ( - "context" - "fmt" - "io" - "net/http" - "time" - - "github.com/fanux/sealos/pkg/utils/logger" - - "github.com/pkg/errors" - - v1 "k8s.io/api/core/v1" - apierrors "k8s.io/apimachinery/pkg/api/errors" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - netutil "k8s.io/apimachinery/pkg/util/net" - "k8s.io/apimachinery/pkg/util/wait" - clientset "k8s.io/client-go/kubernetes" -) - -// Waiter is an interface for waiting for criteria in Kubernetes to happen -type Waiter interface { - // WaitForAPI waits for the API Server's /healthz endpoint to become "ok" - WaitForAPI() error - // WaitForPodsWithLabel waits for Pods in the kube-system namespace to become Ready - WaitForPodsWithLabel(kvLabel string) error - // WaitForPodToDisappear waits for the given Pod in the kube-system namespace to be deleted - WaitForPodToDisappear(staticPodName string) error - // WaitForStaticPodSingleHash fetches sha256 hash for the control plane static pod - WaitForStaticPodSingleHash(nodeName string, component string) (string, error) - // WaitForStaticPodHashChange waits for the given static pod component's static pod hash to get updated. - // By doing that we can be sure that the kubelet has restarted the given Static Pod - WaitForStaticPodHashChange(nodeName, component, previousHash string) error - // WaitForStaticPodControlPlaneHashes fetches sha256 hashes for the control plane static pods - WaitForStaticPodControlPlaneHashes(nodeName string) (map[string]string, error) - // WaitForHealthyKubelet blocks until the kubelet /healthz endpoint returns 'ok' - WaitForHealthyKubelet(initialTimeout time.Duration, healthzEndpoint string) error - // SetTimeout adjusts the timeout to the specified duration - SetTimeout(timeout time.Duration) -} - -// KubeWaiter is an implementation of Waiter that is backed by a Kubernetes client -type KubeWaiter struct { - client clientset.Interface - timeout time.Duration - writer io.Writer -} - -// NewKubeWaiter returns a new Waiter object that talks to the given Kubernetes cluster -func NewKubeWaiter(client clientset.Interface, timeout time.Duration, writer io.Writer) Waiter { - return &KubeWaiter{ - client: client, - timeout: timeout, - writer: writer, - } -} - -// WaitForAPI waits for the API Server's /healthz endpoint to report "ok" -func (w *KubeWaiter) WaitForAPI() error { - start := time.Now() - return wait.PollImmediate(APICallRetryInterval, w.timeout, func() (bool, error) { - healthStatus := 0 - w.client.Discovery().RESTClient().Get().AbsPath("/healthz").Do(context.TODO()).StatusCode(&healthStatus) - if healthStatus != http.StatusOK { - return false, nil - } - - logger.Debug("[apiclient] All control plane components are healthy after %f seconds\n", time.Since(start).Seconds()) - return true, nil - }) -} - -// WaitForPodsWithLabel will lookup pods with the given label and wait until they are all -// reporting status as running. -func (w *KubeWaiter) WaitForPodsWithLabel(kvLabel string) error { - lastKnownPodNumber := -1 - return wait.PollImmediate(APICallRetryInterval, w.timeout, func() (bool, error) { - listOpts := metav1.ListOptions{LabelSelector: kvLabel} - pods, err := w.client.CoreV1().Pods(metav1.NamespaceSystem).List(context.TODO(), listOpts) - if err != nil { - fmt.Fprintf(w.writer, "[apiclient] Error getting Pods with label selector %q [%v]\n", kvLabel, err) - return false, nil - } - - if lastKnownPodNumber != len(pods.Items) { - fmt.Fprintf(w.writer, "[apiclient] Found %d Pods for label selector %s\n", len(pods.Items), kvLabel) - lastKnownPodNumber = len(pods.Items) - } - - if len(pods.Items) == 0 { - return false, nil - } - - for _, pod := range pods.Items { - if pod.Status.Phase != v1.PodRunning { - return false, nil - } - } - - return true, nil - }) -} - -// WaitForPodToDisappear blocks until it timeouts or gets a "NotFound" response from the API Server when getting the Static Pod in question -func (w *KubeWaiter) WaitForPodToDisappear(podName string) error { - return wait.PollImmediate(APICallRetryInterval, w.timeout, func() (bool, error) { - _, err := w.client.CoreV1().Pods(metav1.NamespaceSystem).Get(context.TODO(), podName, metav1.GetOptions{}) - if apierrors.IsNotFound(err) { - logger.Warn("[apiclient] The old Pod %q is now removed (which is desired)\n", podName) - return true, nil - } - return false, nil - }) -} - -// WaitForHealthyKubelet blocks until the kubelet /healthz endpoint returns 'ok' -func (w *KubeWaiter) WaitForHealthyKubelet(initialTimeout time.Duration, healthzEndpoint string) error { - time.Sleep(initialTimeout) - logger.Debug("[kubelet-check] Initial timeout of %v passed.\n", initialTimeout) - return TryRunCommand(func() error { - client := &http.Client{Transport: netutil.SetOldTransportDefaults(&http.Transport{})} - resp, err := client.Get(healthzEndpoint) - if err != nil { - logger.Warn("[kubelet-check] It seems like the kubelet isn't running or healthy.") - logger.Warn("[kubelet-check] The HTTP call equal to 'curl -sSL %s' failed with error: %v.\n", healthzEndpoint, err) - return err - } - defer resp.Body.Close() - if resp.StatusCode != http.StatusOK { - logger.Warn("[kubelet-check] It seems like the kubelet isn't running or healthy.") - logger.Warn("[kubelet-check] The HTTP call equal to 'curl -sSL %s' returned HTTP code %d\n", healthzEndpoint, resp.StatusCode) - return errors.New("the kubelet healthz endpoint is unhealthy") - } - return nil - }, 5) // a failureThreshold of five means waiting for a total of 155 seconds -} - -// SetTimeout adjusts the timeout to the specified duration -func (w *KubeWaiter) SetTimeout(timeout time.Duration) { - w.timeout = timeout -} - -// WaitForStaticPodControlPlaneHashes blocks until it timeouts or gets a hash map for all components and their Static Pods -func (w *KubeWaiter) WaitForStaticPodControlPlaneHashes(nodeName string) (map[string]string, error) { - componentHash := "" - var err error - mirrorPodHashes := map[string]string{} - for _, component := range ControlPlaneComponents { - err = wait.PollImmediate(APICallRetryInterval, w.timeout, func() (bool, error) { - componentHash, err = getStaticPodSingleHash(w.client, nodeName, component) - if err != nil { - return false, nil //nolint:nilerr - } - return true, nil - }) - if err != nil { - return nil, err - } - mirrorPodHashes[component] = componentHash - } - - return mirrorPodHashes, nil -} - -// WaitForStaticPodSingleHash blocks until it timeouts or gets a hash for a single component and its Static Pod -func (w *KubeWaiter) WaitForStaticPodSingleHash(nodeName string, component string) (string, error) { - componentPodHash := "" - var err error - err = wait.PollImmediate(APICallRetryInterval, w.timeout, func() (bool, error) { - componentPodHash, err = getStaticPodSingleHash(w.client, nodeName, component) - if err != nil { - return false, nil //nolint:nilerr - } - return true, nil - }) - - return componentPodHash, err -} - -// WaitForStaticPodHashChange blocks until it timeouts or notices that the Mirror Pod (for the Static Pod, respectively) has changed -// This implicitly means this function blocks until the kubelet has restarted the Static Pod in question -func (w *KubeWaiter) WaitForStaticPodHashChange(nodeName, component, previousHash string) error { - return wait.PollImmediate(APICallRetryInterval, w.timeout, func() (bool, error) { - hash, err := getStaticPodSingleHash(w.client, nodeName, component) - if err != nil { - return false, nil //nolint:nilerr - } - // We should continue polling until the UID changes - if hash == previousHash { - return false, nil - } - - return true, nil - }) -} - -// getStaticPodSingleHash computes hashes for a single Static Pod resource -func getStaticPodSingleHash(client clientset.Interface, nodeName string, component string) (string, error) { - staticPodName := fmt.Sprintf("%s-%s", component, nodeName) - staticPod, err := client.CoreV1().Pods(metav1.NamespaceSystem).Get(context.TODO(), staticPodName, metav1.GetOptions{}) - if err != nil { - return "", err - } - - staticPodHash := staticPod.Annotations["kubernetes.io/config.hash"] - logger.Debug("Static pod: %s hash: %s\n", staticPodName, staticPodHash) - return staticPodHash, nil -} - -// TryRunCommand runs a function a maximum of failureThreshold times, and retries on error. If failureThreshold is hit; the last error is returned -func TryRunCommand(f func() error, failureThreshold int) error { - backoff := wait.Backoff{ - Duration: 5 * time.Second, - Factor: 2, // double the timeout for every failure - Steps: failureThreshold, - } - return wait.ExponentialBackoff(backoff, func() (bool, error) { - err := f() - if err != nil { - // Retry until the timeout - return false, nil //nolint:nilerr - } - // The last f() call was a success, return cleanly - return true, nil - }) -} diff --git a/pkg/utils/archive/compress_test.go b/pkg/utils/archive/compress_test.go index 0164d9d51..b210a36c9 100644 --- a/pkg/utils/archive/compress_test.go +++ b/pkg/utils/archive/compress_test.go @@ -19,10 +19,9 @@ import ( "fmt" "io" "os" + "path" "path/filepath" "testing" - - "golang.org/x/sys/unix" ) const basePath = "/tmp" @@ -134,9 +133,11 @@ func TestName(t *testing.T) { //if err != nil { // t.Error(err) //} - err := unix.Setxattr("abc", "trusted.overlay.opaque", []byte{'y'}, 0) - if err != nil { - t.Error(err) - } + //err := unix.Setxattr("abc", "trusted.overlay.opaque", []byte{'y'}, 0) + //if err != nil { + // t.Error(err) + //} //fmt.Println(fm.String()) + data:=path.Join("http://localhost:10250","heathy") + t.Log(data) } diff --git a/vendor/k8s.io/client-go/dynamic/interface.go b/vendor/k8s.io/client-go/dynamic/interface.go new file mode 100644 index 000000000..b08067c34 --- /dev/null +++ b/vendor/k8s.io/client-go/dynamic/interface.go @@ -0,0 +1,61 @@ +/* +Copyright 2016 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 dynamic + +import ( + "context" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/types" + "k8s.io/apimachinery/pkg/watch" +) + +type Interface interface { + Resource(resource schema.GroupVersionResource) NamespaceableResourceInterface +} + +type ResourceInterface interface { + Create(ctx context.Context, obj *unstructured.Unstructured, options metav1.CreateOptions, subresources ...string) (*unstructured.Unstructured, error) + Update(ctx context.Context, obj *unstructured.Unstructured, options metav1.UpdateOptions, subresources ...string) (*unstructured.Unstructured, error) + UpdateStatus(ctx context.Context, obj *unstructured.Unstructured, options metav1.UpdateOptions) (*unstructured.Unstructured, error) + Delete(ctx context.Context, name string, options metav1.DeleteOptions, subresources ...string) error + DeleteCollection(ctx context.Context, options metav1.DeleteOptions, listOptions metav1.ListOptions) error + Get(ctx context.Context, name string, options metav1.GetOptions, subresources ...string) (*unstructured.Unstructured, error) + List(ctx context.Context, opts metav1.ListOptions) (*unstructured.UnstructuredList, error) + Watch(ctx context.Context, opts metav1.ListOptions) (watch.Interface, error) + Patch(ctx context.Context, name string, pt types.PatchType, data []byte, options metav1.PatchOptions, subresources ...string) (*unstructured.Unstructured, error) +} + +type NamespaceableResourceInterface interface { + Namespace(string) ResourceInterface + ResourceInterface +} + +// APIPathResolverFunc knows how to convert a groupVersion to its API path. The Kind field is optional. +// TODO find a better place to move this for existing callers +type APIPathResolverFunc func(kind schema.GroupVersionKind) string + +// LegacyAPIPathResolverFunc can resolve paths properly with the legacy API. +// TODO find a better place to move this for existing callers +func LegacyAPIPathResolverFunc(kind schema.GroupVersionKind) string { + if len(kind.Group) == 0 { + return "/api" + } + return "/apis" +} diff --git a/vendor/k8s.io/client-go/dynamic/scheme.go b/vendor/k8s.io/client-go/dynamic/scheme.go new file mode 100644 index 000000000..3168c872c --- /dev/null +++ b/vendor/k8s.io/client-go/dynamic/scheme.go @@ -0,0 +1,108 @@ +/* +Copyright 2018 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 dynamic + +import ( + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/runtime/serializer" + "k8s.io/apimachinery/pkg/runtime/serializer/json" +) + +var watchScheme = runtime.NewScheme() +var basicScheme = runtime.NewScheme() +var deleteScheme = runtime.NewScheme() +var parameterScheme = runtime.NewScheme() +var deleteOptionsCodec = serializer.NewCodecFactory(deleteScheme) +var dynamicParameterCodec = runtime.NewParameterCodec(parameterScheme) + +var versionV1 = schema.GroupVersion{Version: "v1"} + +func init() { + metav1.AddToGroupVersion(watchScheme, versionV1) + metav1.AddToGroupVersion(basicScheme, versionV1) + metav1.AddToGroupVersion(parameterScheme, versionV1) + metav1.AddToGroupVersion(deleteScheme, versionV1) +} + +// basicNegotiatedSerializer is used to handle discovery and error handling serialization +type basicNegotiatedSerializer struct{} + +func (s basicNegotiatedSerializer) SupportedMediaTypes() []runtime.SerializerInfo { + return []runtime.SerializerInfo{ + { + MediaType: "application/json", + MediaTypeType: "application", + MediaTypeSubType: "json", + EncodesAsText: true, + Serializer: json.NewSerializer(json.DefaultMetaFactory, unstructuredCreater{basicScheme}, unstructuredTyper{basicScheme}, false), + PrettySerializer: json.NewSerializer(json.DefaultMetaFactory, unstructuredCreater{basicScheme}, unstructuredTyper{basicScheme}, true), + StreamSerializer: &runtime.StreamSerializerInfo{ + EncodesAsText: true, + Serializer: json.NewSerializer(json.DefaultMetaFactory, basicScheme, basicScheme, false), + Framer: json.Framer, + }, + }, + } +} + +func (s basicNegotiatedSerializer) EncoderForVersion(encoder runtime.Encoder, gv runtime.GroupVersioner) runtime.Encoder { + return runtime.WithVersionEncoder{ + Version: gv, + Encoder: encoder, + ObjectTyper: unstructuredTyper{basicScheme}, + } +} + +func (s basicNegotiatedSerializer) DecoderToVersion(decoder runtime.Decoder, gv runtime.GroupVersioner) runtime.Decoder { + return decoder +} + +type unstructuredCreater struct { + nested runtime.ObjectCreater +} + +func (c unstructuredCreater) New(kind schema.GroupVersionKind) (runtime.Object, error) { + out, err := c.nested.New(kind) + if err == nil { + return out, nil + } + out = &unstructured.Unstructured{} + out.GetObjectKind().SetGroupVersionKind(kind) + return out, nil +} + +type unstructuredTyper struct { + nested runtime.ObjectTyper +} + +func (t unstructuredTyper) ObjectKinds(obj runtime.Object) ([]schema.GroupVersionKind, bool, error) { + kinds, unversioned, err := t.nested.ObjectKinds(obj) + if err == nil { + return kinds, unversioned, nil + } + if _, ok := obj.(runtime.Unstructured); ok && !obj.GetObjectKind().GroupVersionKind().Empty() { + return []schema.GroupVersionKind{obj.GetObjectKind().GroupVersionKind()}, false, nil + } + return nil, false, err +} + +func (t unstructuredTyper) Recognizes(gvk schema.GroupVersionKind) bool { + return true +} diff --git a/vendor/k8s.io/client-go/dynamic/simple.go b/vendor/k8s.io/client-go/dynamic/simple.go new file mode 100644 index 000000000..9ae320d30 --- /dev/null +++ b/vendor/k8s.io/client-go/dynamic/simple.go @@ -0,0 +1,327 @@ +/* +Copyright 2018 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 dynamic + +import ( + "context" + "fmt" + + "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/apis/meta/v1/unstructured" + "k8s.io/apimachinery/pkg/runtime" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/apimachinery/pkg/types" + "k8s.io/apimachinery/pkg/watch" + "k8s.io/client-go/rest" +) + +type dynamicClient struct { + client *rest.RESTClient +} + +var _ Interface = &dynamicClient{} + +// ConfigFor returns a copy of the provided config with the +// appropriate dynamic client defaults set. +func ConfigFor(inConfig *rest.Config) *rest.Config { + config := rest.CopyConfig(inConfig) + config.AcceptContentTypes = "application/json" + config.ContentType = "application/json" + config.NegotiatedSerializer = basicNegotiatedSerializer{} // this gets used for discovery and error handling types + if config.UserAgent == "" { + config.UserAgent = rest.DefaultKubernetesUserAgent() + } + return config +} + +// NewForConfigOrDie creates a new Interface for the given config and +// panics if there is an error in the config. +func NewForConfigOrDie(c *rest.Config) Interface { + ret, err := NewForConfig(c) + if err != nil { + panic(err) + } + return ret +} + +// NewForConfig creates a new dynamic client or returns an error. +func NewForConfig(inConfig *rest.Config) (Interface, error) { + config := ConfigFor(inConfig) + // for serializing the options + config.GroupVersion = &schema.GroupVersion{} + config.APIPath = "/if-you-see-this-search-for-the-break" + + restClient, err := rest.RESTClientFor(config) + if err != nil { + return nil, err + } + + return &dynamicClient{client: restClient}, nil +} + +type dynamicResourceClient struct { + client *dynamicClient + namespace string + resource schema.GroupVersionResource +} + +func (c *dynamicClient) Resource(resource schema.GroupVersionResource) NamespaceableResourceInterface { + return &dynamicResourceClient{client: c, resource: resource} +} + +func (c *dynamicResourceClient) Namespace(ns string) ResourceInterface { + ret := *c + ret.namespace = ns + return &ret +} + +func (c *dynamicResourceClient) Create(ctx context.Context, obj *unstructured.Unstructured, opts metav1.CreateOptions, subresources ...string) (*unstructured.Unstructured, error) { + outBytes, err := runtime.Encode(unstructured.UnstructuredJSONScheme, obj) + if err != nil { + return nil, err + } + name := "" + if len(subresources) > 0 { + accessor, err := meta.Accessor(obj) + if err != nil { + return nil, err + } + name = accessor.GetName() + if len(name) == 0 { + return nil, fmt.Errorf("name is required") + } + } + + result := c.client.client. + Post(). + AbsPath(append(c.makeURLSegments(name), subresources...)...). + Body(outBytes). + SpecificallyVersionedParams(&opts, dynamicParameterCodec, versionV1). + Do(ctx) + if err := result.Error(); err != nil { + return nil, err + } + + retBytes, err := result.Raw() + if err != nil { + return nil, err + } + uncastObj, err := runtime.Decode(unstructured.UnstructuredJSONScheme, retBytes) + if err != nil { + return nil, err + } + return uncastObj.(*unstructured.Unstructured), nil +} + +func (c *dynamicResourceClient) Update(ctx context.Context, obj *unstructured.Unstructured, opts metav1.UpdateOptions, subresources ...string) (*unstructured.Unstructured, error) { + accessor, err := meta.Accessor(obj) + if err != nil { + return nil, err + } + name := accessor.GetName() + if len(name) == 0 { + return nil, fmt.Errorf("name is required") + } + outBytes, err := runtime.Encode(unstructured.UnstructuredJSONScheme, obj) + if err != nil { + return nil, err + } + + result := c.client.client. + Put(). + AbsPath(append(c.makeURLSegments(name), subresources...)...). + Body(outBytes). + SpecificallyVersionedParams(&opts, dynamicParameterCodec, versionV1). + Do(ctx) + if err := result.Error(); err != nil { + return nil, err + } + + retBytes, err := result.Raw() + if err != nil { + return nil, err + } + uncastObj, err := runtime.Decode(unstructured.UnstructuredJSONScheme, retBytes) + if err != nil { + return nil, err + } + return uncastObj.(*unstructured.Unstructured), nil +} + +func (c *dynamicResourceClient) UpdateStatus(ctx context.Context, obj *unstructured.Unstructured, opts metav1.UpdateOptions) (*unstructured.Unstructured, error) { + accessor, err := meta.Accessor(obj) + if err != nil { + return nil, err + } + name := accessor.GetName() + if len(name) == 0 { + return nil, fmt.Errorf("name is required") + } + + outBytes, err := runtime.Encode(unstructured.UnstructuredJSONScheme, obj) + if err != nil { + return nil, err + } + + result := c.client.client. + Put(). + AbsPath(append(c.makeURLSegments(name), "status")...). + Body(outBytes). + SpecificallyVersionedParams(&opts, dynamicParameterCodec, versionV1). + Do(ctx) + if err := result.Error(); err != nil { + return nil, err + } + + retBytes, err := result.Raw() + if err != nil { + return nil, err + } + uncastObj, err := runtime.Decode(unstructured.UnstructuredJSONScheme, retBytes) + if err != nil { + return nil, err + } + return uncastObj.(*unstructured.Unstructured), nil +} + +func (c *dynamicResourceClient) Delete(ctx context.Context, name string, opts metav1.DeleteOptions, subresources ...string) error { + if len(name) == 0 { + return fmt.Errorf("name is required") + } + deleteOptionsByte, err := runtime.Encode(deleteOptionsCodec.LegacyCodec(schema.GroupVersion{Version: "v1"}), &opts) + if err != nil { + return err + } + + result := c.client.client. + Delete(). + AbsPath(append(c.makeURLSegments(name), subresources...)...). + Body(deleteOptionsByte). + Do(ctx) + return result.Error() +} + +func (c *dynamicResourceClient) DeleteCollection(ctx context.Context, opts metav1.DeleteOptions, listOptions metav1.ListOptions) error { + deleteOptionsByte, err := runtime.Encode(deleteOptionsCodec.LegacyCodec(schema.GroupVersion{Version: "v1"}), &opts) + if err != nil { + return err + } + + result := c.client.client. + Delete(). + AbsPath(c.makeURLSegments("")...). + Body(deleteOptionsByte). + SpecificallyVersionedParams(&listOptions, dynamicParameterCodec, versionV1). + Do(ctx) + return result.Error() +} + +func (c *dynamicResourceClient) Get(ctx context.Context, name string, opts metav1.GetOptions, subresources ...string) (*unstructured.Unstructured, error) { + if len(name) == 0 { + return nil, fmt.Errorf("name is required") + } + result := c.client.client.Get().AbsPath(append(c.makeURLSegments(name), subresources...)...).SpecificallyVersionedParams(&opts, dynamicParameterCodec, versionV1).Do(ctx) + if err := result.Error(); err != nil { + return nil, err + } + retBytes, err := result.Raw() + if err != nil { + return nil, err + } + uncastObj, err := runtime.Decode(unstructured.UnstructuredJSONScheme, retBytes) + if err != nil { + return nil, err + } + return uncastObj.(*unstructured.Unstructured), nil +} + +func (c *dynamicResourceClient) List(ctx context.Context, opts metav1.ListOptions) (*unstructured.UnstructuredList, error) { + result := c.client.client.Get().AbsPath(c.makeURLSegments("")...).SpecificallyVersionedParams(&opts, dynamicParameterCodec, versionV1).Do(ctx) + if err := result.Error(); err != nil { + return nil, err + } + retBytes, err := result.Raw() + if err != nil { + return nil, err + } + uncastObj, err := runtime.Decode(unstructured.UnstructuredJSONScheme, retBytes) + if err != nil { + return nil, err + } + if list, ok := uncastObj.(*unstructured.UnstructuredList); ok { + return list, nil + } + + list, err := uncastObj.(*unstructured.Unstructured).ToList() + if err != nil { + return nil, err + } + return list, nil +} + +func (c *dynamicResourceClient) Watch(ctx context.Context, opts metav1.ListOptions) (watch.Interface, error) { + opts.Watch = true + return c.client.client.Get().AbsPath(c.makeURLSegments("")...). + SpecificallyVersionedParams(&opts, dynamicParameterCodec, versionV1). + Watch(ctx) +} + +func (c *dynamicResourceClient) Patch(ctx context.Context, name string, pt types.PatchType, data []byte, opts metav1.PatchOptions, subresources ...string) (*unstructured.Unstructured, error) { + if len(name) == 0 { + return nil, fmt.Errorf("name is required") + } + result := c.client.client. + Patch(pt). + AbsPath(append(c.makeURLSegments(name), subresources...)...). + Body(data). + SpecificallyVersionedParams(&opts, dynamicParameterCodec, versionV1). + Do(ctx) + if err := result.Error(); err != nil { + return nil, err + } + retBytes, err := result.Raw() + if err != nil { + return nil, err + } + uncastObj, err := runtime.Decode(unstructured.UnstructuredJSONScheme, retBytes) + if err != nil { + return nil, err + } + return uncastObj.(*unstructured.Unstructured), nil +} + +func (c *dynamicResourceClient) makeURLSegments(name string) []string { + url := []string{} + if len(c.resource.Group) == 0 { + url = append(url, "api") + } else { + url = append(url, "apis", c.resource.Group) + } + url = append(url, c.resource.Version) + + if len(c.namespace) > 0 { + url = append(url, "namespaces", c.namespace) + } + url = append(url, c.resource.Resource) + + if len(name) > 0 { + url = append(url, name) + } + + return url +} diff --git a/vendor/k8s.io/client-go/util/retry/OWNERS b/vendor/k8s.io/client-go/util/retry/OWNERS deleted file mode 100644 index dec3e88d6..000000000 --- a/vendor/k8s.io/client-go/util/retry/OWNERS +++ /dev/null @@ -1,4 +0,0 @@ -# See the OWNERS docs at https://go.k8s.io/owners - -reviewers: -- caesarxuchao diff --git a/vendor/k8s.io/client-go/util/retry/util.go b/vendor/k8s.io/client-go/util/retry/util.go deleted file mode 100644 index 15e2722f3..000000000 --- a/vendor/k8s.io/client-go/util/retry/util.go +++ /dev/null @@ -1,105 +0,0 @@ -/* -Copyright 2016 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 retry - -import ( - "time" - - "k8s.io/apimachinery/pkg/api/errors" - "k8s.io/apimachinery/pkg/util/wait" -) - -// DefaultRetry is the recommended retry for a conflict where multiple clients -// are making changes to the same resource. -var DefaultRetry = wait.Backoff{ - Steps: 5, - Duration: 10 * time.Millisecond, - Factor: 1.0, - Jitter: 0.1, -} - -// DefaultBackoff is the recommended backoff for a conflict where a client -// may be attempting to make an unrelated modification to a resource under -// active management by one or more controllers. -var DefaultBackoff = wait.Backoff{ - Steps: 4, - Duration: 10 * time.Millisecond, - Factor: 5.0, - Jitter: 0.1, -} - -// OnError allows the caller to retry fn in case the error returned by fn is retriable -// according to the provided function. backoff defines the maximum retries and the wait -// interval between two retries. -func OnError(backoff wait.Backoff, retriable func(error) bool, fn func() error) error { - var lastErr error - err := wait.ExponentialBackoff(backoff, func() (bool, error) { - err := fn() - switch { - case err == nil: - return true, nil - case retriable(err): - lastErr = err - return false, nil - default: - return false, err - } - }) - if err == wait.ErrWaitTimeout { - err = lastErr - } - return err -} - -// RetryOnConflict is used to make an update to a resource when you have to worry about -// conflicts caused by other code making unrelated updates to the resource at the same -// time. fn should fetch the resource to be modified, make appropriate changes to it, try -// to update it, and return (unmodified) the error from the update function. On a -// successful update, RetryOnConflict will return nil. If the update function returns a -// "Conflict" error, RetryOnConflict will wait some amount of time as described by -// backoff, and then try again. On a non-"Conflict" error, or if it retries too many times -// and gives up, RetryOnConflict will return an error to the caller. -// -// err := retry.RetryOnConflict(retry.DefaultRetry, func() error { -// // Fetch the resource here; you need to refetch it on every try, since -// // if you got a conflict on the last update attempt then you need to get -// // the current version before making your own changes. -// pod, err := c.Pods("mynamespace").Get(name, metav1.GetOptions{}) -// if err ! nil { -// return err -// } -// -// // Make whatever updates to the resource are needed -// pod.Status.Phase = v1.PodFailed -// -// // Try to update -// _, err = c.Pods("mynamespace").UpdateStatus(pod) -// // You have to return err itself here (not wrapped inside another error) -// // so that RetryOnConflict can identify it correctly. -// return err -// }) -// if err != nil { -// // May be conflict if max retries were hit, or may be something unrelated -// // like permissions or a network error -// return err -// } -// ... -// -// TODO: Make Backoff an interface? -func RetryOnConflict(backoff wait.Backoff, fn func() error) error { - return OnError(backoff, errors.IsConflict, fn) -} diff --git a/vendor/modules.txt b/vendor/modules.txt index 8f6717cc6..656fa7562 100644 --- a/vendor/modules.txt +++ b/vendor/modules.txt @@ -749,6 +749,7 @@ k8s.io/client-go/applyconfigurations/storage/v1 k8s.io/client-go/applyconfigurations/storage/v1alpha1 k8s.io/client-go/applyconfigurations/storage/v1beta1 k8s.io/client-go/discovery +k8s.io/client-go/dynamic k8s.io/client-go/kubernetes k8s.io/client-go/kubernetes/scheme k8s.io/client-go/kubernetes/typed/admissionregistration/v1 @@ -814,7 +815,6 @@ k8s.io/client-go/util/connrotation k8s.io/client-go/util/flowcontrol k8s.io/client-go/util/homedir k8s.io/client-go/util/keyutil -k8s.io/client-go/util/retry k8s.io/client-go/util/workqueue # k8s.io/cluster-bootstrap v0.21.0 ## explicit