* Copyright (c) 2026 Huawei Technologies Co., Ltd.
* openFuyao is licensed under Mulan PSL v2.
* You can use this software according to the terms and conditions of the Mulan PSL v2.
* You may obtain a copy of Mulan PSL v2 at:
* http://license.coscl.org.cn/MulanPSL2
* THIS SOFTWARE IS PROVIDED ON AN "AS IS" BASIS, WITHOUT WARRANTIES OF ANY KIND,
* EITHER EXPRESS OR IMPLIED, INCLUDING BUT NOT LIMITED TO NON-INFRINGEMENT,
* MERCHANTABILITY OR FIT FOR A PARTICULAR PURPOSE.
* See the Mulan PSL v2 for more details.
*/
package k8s
import (
"context"
"fmt"
"os"
"path/filepath"
"strings"
"gitcode.com/openFuyao/e2e-auto-test/e2e/framework/executor"
discoveryv1 "k8s.io/api/discovery/v1"
corev1 "k8s.io/api/core/v1"
"k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/client-go/dynamic"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/clientcmd"
)
type K8SClient struct {
RestConfig *rest.Config
Clientset *kubernetes.Clientset
DynamicClient dynamic.Interface
}
func localKubeconfigPath() (string, error) {
if v := strings.TrimSpace(os.Getenv("KUBECONFIG")); v != "" {
return v, nil
}
home, err := os.UserHomeDir()
if err != nil {
return "", fmt.Errorf("get home dir failed: %w", err)
}
return filepath.Join(home, ".kube", "config"), nil
}
func NewK8SClientFromLocalKubeconfig() (*K8SClient, error) {
cfgPath, err := localKubeconfigPath()
if err != nil {
return nil, err
}
restCfg, err := clientcmd.BuildConfigFromFlags("", cfgPath)
if err != nil {
return nil, fmt.Errorf("failed to build kubeconfig (%s): %w", cfgPath, err)
}
applyE2EClientRateLimits(restCfg)
clientset, err := kubernetes.NewForConfig(restCfg)
if err != nil {
return nil, fmt.Errorf("failed to create clientset: %w", err)
}
dynamicClient, err := dynamic.NewForConfig(restCfg)
if err != nil {
return nil, fmt.Errorf("failed to create dynamic client: %w", err)
}
return &K8SClient{RestConfig: restCfg, Clientset: clientset, DynamicClient: dynamicClient}, nil
}
func applyE2EClientRateLimits(cfg *rest.Config) {
if cfg == nil {
return
}
if cfg.QPS <= 0 {
cfg.QPS = 50
}
if cfg.Burst <= 0 {
cfg.Burst = 100
}
}
func NewK8SClientViaSSH(sshClient *executor.SSHExecutor) (*K8SClient, error) {
kubeconfigPath, err := localKubeconfigPath()
if err != nil {
return nil, fmt.Errorf("failed to get kubeconfig path: %w", err)
}
getKubeconfigCmd := fmt.Sprintf("cat %s", kubeconfigPath)
execResult, err := sshClient.Exec(getKubeconfigCmd)
if err != nil {
return nil, fmt.Errorf("failed to get kubeconfig via SSH: %w", err)
}
kubeconfig := execResult.Stdout
tmpFile, err := os.CreateTemp("", "kubeconfig-*")
if err != nil {
return nil, fmt.Errorf("failed to create temp file: %w", err)
}
defer os.Remove(tmpFile.Name())
if _, err := tmpFile.WriteString(kubeconfig); err != nil {
return nil, fmt.Errorf("failed to write kubeconfig: %w", err)
}
tmpFile.Close()
config, err := clientcmd.BuildConfigFromFlags("", tmpFile.Name())
if err != nil {
return nil, fmt.Errorf("failed to build config: %w", err)
}
applyE2EClientRateLimits(config)
clientset, err := kubernetes.NewForConfig(config)
if err != nil {
return nil, fmt.Errorf("failed to create clientset: %w", err)
}
dynamicClient, err := dynamic.NewForConfig(config)
if err != nil {
return nil, fmt.Errorf("failed to create dynamic client: %w", err)
}
return &K8SClient{RestConfig: config, Clientset: clientset, DynamicClient: dynamicClient}, nil
}
func NewK8SDynamicClientViaSSH(sshClient *executor.SSHExecutor) (dynamic.Interface, error) {
kubeconfigPath, err := localKubeconfigPath()
if err != nil {
return nil, fmt.Errorf("failed to get kubeconfig path: %w", err)
}
getKubeconfigCmd := fmt.Sprintf("cat %s", kubeconfigPath)
execResult, err := sshClient.Exec(getKubeconfigCmd)
if err != nil {
return nil, fmt.Errorf("failed to get kubeconfig via SSH: %w", err)
}
kubeconfig := execResult.Stdout
tmpFile, err := os.CreateTemp("", "kubeconfig-*")
if err != nil {
return nil, fmt.Errorf("failed to create temp file: %w", err)
}
defer os.Remove(tmpFile.Name())
if _, err := tmpFile.WriteString(kubeconfig); err != nil {
return nil, fmt.Errorf("failed to write kubeconfig: %w", err)
}
tmpFile.Close()
config, err := clientcmd.BuildConfigFromFlags("", tmpFile.Name())
if err != nil {
return nil, fmt.Errorf("failed to build config: %w", err)
}
dynamicClient, err := dynamic.NewForConfig(config)
if err != nil {
return nil, fmt.Errorf("failed to create dynamic client: %w", err)
}
return dynamicClient, nil
}
func (c *K8SClient) ListPods(ctx context.Context, namespace string, opts metav1.ListOptions) (*corev1.PodList, error) {
return c.Clientset.CoreV1().Pods(namespace).List(ctx, opts)
}
func (c *K8SClient) ListServices(ctx context.Context, namespace string, opts metav1.ListOptions) (*corev1.ServiceList, error) {
return c.Clientset.CoreV1().Services(namespace).List(ctx, opts)
}
func (c *K8SClient) ListEndpointSlices(ctx context.Context, namespace string, opts metav1.ListOptions) (*discoveryv1.EndpointSliceList, error) {
return c.Clientset.DiscoveryV1().EndpointSlices(namespace).List(ctx, opts)
}
func (c *K8SClient) ListNodes(ctx context.Context, opts metav1.ListOptions) (*corev1.NodeList, error) {
return c.Clientset.CoreV1().Nodes().List(ctx, opts)
}
func (c *K8SClient) ListServiceAccounts(ctx context.Context, namespace string, opts metav1.ListOptions) (*corev1.ServiceAccountList, error) {
return c.Clientset.CoreV1().ServiceAccounts(namespace).List(ctx, opts)
}
func (c *K8SClient) DeleteServiceAccount(ctx context.Context, namespace, name string) error {
return c.Clientset.CoreV1().ServiceAccounts(namespace).Delete(ctx, name, metav1.DeleteOptions{})
}
func (c *K8SClient) ListNamespaces(ctx context.Context, opts metav1.ListOptions) (*corev1.NamespaceList, error) {
return c.Clientset.CoreV1().Namespaces().List(ctx, opts)
}
func (c *K8SClient) GetNamespace(ctx context.Context, name string) (*corev1.Namespace, error) {
return c.Clientset.CoreV1().Namespaces().Get(ctx, name, metav1.GetOptions{})
}
func (c *K8SClient) CreateNamespaceIfNotExists(ctx context.Context, namespace string) error {
_, err := c.Clientset.CoreV1().Namespaces().Get(ctx, namespace, metav1.GetOptions{})
if err == nil {
return nil
}
ns := &corev1.Namespace{
ObjectMeta: metav1.ObjectMeta{
Name: namespace,
},
}
_, err = c.Clientset.CoreV1().Namespaces().Create(ctx, ns, metav1.CreateOptions{})
return err
}
func (c *K8SClient) DeleteNamespace(ctx context.Context, namespace string) error {
_, err := c.Clientset.CoreV1().Namespaces().Get(ctx, namespace, metav1.GetOptions{})
if errors.IsNotFound(err) {
return nil
}
if err != nil {
return fmt.Errorf("failed to check namespace: %v", err)
}
return c.Clientset.CoreV1().Namespaces().Delete(ctx, namespace, metav1.DeleteOptions{})
}