fix pre-commit and cleanup
This commit is contained in:
@@ -1,223 +1,223 @@
|
||||
package kubectlcmds
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"regexp"
|
||||
"bytes"
|
||||
"context"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"os"
|
||||
"path/filepath"
|
||||
"regexp"
|
||||
|
||||
"github.com/go-logr/logr"
|
||||
"github.com/go-logr/zapr"
|
||||
"github.com/truecharts/public/clustertool/pkg/helper"
|
||||
"go.uber.org/zap"
|
||||
"go.uber.org/zap/zapcore"
|
||||
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
|
||||
"k8s.io/client-go/tools/clientcmd"
|
||||
"k8s.io/client-go/util/homedir"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
"sigs.k8s.io/kustomize/api/krusty"
|
||||
"sigs.k8s.io/kustomize/kyaml/filesys"
|
||||
"sigs.k8s.io/kustomize/kyaml/kio"
|
||||
"sigs.k8s.io/yaml"
|
||||
"github.com/go-logr/logr"
|
||||
"github.com/go-logr/zapr"
|
||||
"github.com/truecharts/public/clustertool/pkg/helper"
|
||||
"go.uber.org/zap"
|
||||
"go.uber.org/zap/zapcore"
|
||||
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
|
||||
"k8s.io/client-go/tools/clientcmd"
|
||||
"k8s.io/client-go/util/homedir"
|
||||
"sigs.k8s.io/controller-runtime/pkg/client"
|
||||
"sigs.k8s.io/kustomize/api/krusty"
|
||||
"sigs.k8s.io/kustomize/kyaml/filesys"
|
||||
"sigs.k8s.io/kustomize/kyaml/kio"
|
||||
"sigs.k8s.io/yaml"
|
||||
)
|
||||
|
||||
// getKubeClient initializes and returns a controller-runtime client.Client
|
||||
func getKubeClient() (client.Client, error) {
|
||||
// Load kubeconfig from the default location
|
||||
kubeconfig := filepath.Join(homedir.HomeDir(), ".kube", "config")
|
||||
config, err := clientcmd.BuildConfigFromFlags("", kubeconfig)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to load kubeconfig: %v", err)
|
||||
}
|
||||
// Load kubeconfig from the default location
|
||||
kubeconfig := filepath.Join(homedir.HomeDir(), ".kube", "config")
|
||||
config, err := clientcmd.BuildConfigFromFlags("", kubeconfig)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to load kubeconfig: %v", err)
|
||||
}
|
||||
|
||||
// Create a controller-runtime client
|
||||
c, err := client.New(config, client.Options{})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to create Kubernetes client: %v", err)
|
||||
}
|
||||
// Create a controller-runtime client
|
||||
c, err := client.New(config, client.Options{})
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to create Kubernetes client: %v", err)
|
||||
}
|
||||
|
||||
return c, nil
|
||||
return c, nil
|
||||
}
|
||||
|
||||
// setupLogger initializes a logger that writes to a buffer and returns both
|
||||
func setupLogger() (logr.Logger, *bytes.Buffer, error) {
|
||||
// Create a buffer to capture logs
|
||||
var buf bytes.Buffer
|
||||
// Create a buffer to capture logs
|
||||
var buf bytes.Buffer
|
||||
|
||||
// Create a WriteSyncer to write to the buffer
|
||||
writeSyncer := zapcore.AddSync(&buf)
|
||||
// Create a WriteSyncer to write to the buffer
|
||||
writeSyncer := zapcore.AddSync(&buf)
|
||||
|
||||
// Configure zap to use console encoder for readability
|
||||
encoderCfg := zap.NewProductionEncoderConfig()
|
||||
encoderCfg.EncodeTime = zapcore.ISO8601TimeEncoder
|
||||
encoder := zapcore.NewConsoleEncoder(encoderCfg)
|
||||
// Configure zap to use console encoder for readability
|
||||
encoderCfg := zap.NewProductionEncoderConfig()
|
||||
encoderCfg.EncodeTime = zapcore.ISO8601TimeEncoder
|
||||
encoder := zapcore.NewConsoleEncoder(encoderCfg)
|
||||
|
||||
// Set log level to Info
|
||||
level := zapcore.InfoLevel
|
||||
// Set log level to Info
|
||||
level := zapcore.InfoLevel
|
||||
|
||||
// Create zap core
|
||||
core := zapcore.NewCore(encoder, writeSyncer, level)
|
||||
// Create zap core
|
||||
core := zapcore.NewCore(encoder, writeSyncer, level)
|
||||
|
||||
// Create zap logger
|
||||
zapLogger := zap.New(core)
|
||||
// Create zap logger
|
||||
zapLogger := zap.New(core)
|
||||
|
||||
// Wrap zap logger with zapr to get a logr.Logger interface
|
||||
log := zapr.NewLogger(zapLogger)
|
||||
// Wrap zap logger with zapr to get a logr.Logger interface
|
||||
log := zapr.NewLogger(zapLogger)
|
||||
|
||||
return log, &buf, nil
|
||||
return log, &buf, nil
|
||||
}
|
||||
|
||||
// applyYAML applies the given YAML data to the Kubernetes cluster using the provided client and logger
|
||||
func applyYAML(k8sClient client.Client, yamlData []byte, log logr.Logger) error {
|
||||
// Parse the YAML into KIO nodes
|
||||
reader := kio.ByteReader{
|
||||
Reader: bytes.NewReader(yamlData),
|
||||
}
|
||||
nodes, err := reader.Read()
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to parse YAML: %v", err)
|
||||
}
|
||||
// Parse the YAML into KIO nodes
|
||||
reader := kio.ByteReader{
|
||||
Reader: bytes.NewReader(yamlData),
|
||||
}
|
||||
nodes, err := reader.Read()
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to parse YAML: %v", err)
|
||||
}
|
||||
|
||||
// Apply each node to the cluster
|
||||
for _, node := range nodes {
|
||||
obj := &unstructured.Unstructured{}
|
||||
if err := yaml.Unmarshal([]byte(node.MustString()), obj); err != nil {
|
||||
return fmt.Errorf("failed to unmarshal node: %v", err)
|
||||
}
|
||||
if err := k8sClient.Patch(context.TODO(), obj, client.Apply, client.FieldOwner("kustomize-controller")); err != nil {
|
||||
return fmt.Errorf("failed to apply object: %v", err)
|
||||
}
|
||||
log.Info("Successfully applied object", "object", obj.GetName(), "kind", obj.GetKind(), "namespace", obj.GetNamespace())
|
||||
}
|
||||
// Apply each node to the cluster
|
||||
for _, node := range nodes {
|
||||
obj := &unstructured.Unstructured{}
|
||||
if err := yaml.Unmarshal([]byte(node.MustString()), obj); err != nil {
|
||||
return fmt.Errorf("failed to unmarshal node: %v", err)
|
||||
}
|
||||
if err := k8sClient.Patch(context.TODO(), obj, client.Apply, client.FieldOwner("kustomize-controller")); err != nil {
|
||||
return fmt.Errorf("failed to apply object: %v", err)
|
||||
}
|
||||
log.Info("Successfully applied object", "object", obj.GetName(), "kind", obj.GetKind(), "namespace", obj.GetNamespace())
|
||||
}
|
||||
|
||||
return nil
|
||||
return nil
|
||||
}
|
||||
|
||||
// filterLogOutput filters the log data by removing strings that match any of the provided regex patterns
|
||||
func filterLogOutput(logData string) (string, error) {
|
||||
filteredLog := logData
|
||||
for _, pattern := range helper.KubeFilterStr {
|
||||
re, err := regexp.Compile(pattern)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("invalid regex pattern '%s': %v", pattern, err)
|
||||
}
|
||||
filteredLog = re.ReplaceAllString(filteredLog, "")
|
||||
}
|
||||
return filteredLog, nil
|
||||
filteredLog := logData
|
||||
for _, pattern := range helper.KubeFilterStr {
|
||||
re, err := regexp.Compile(pattern)
|
||||
if err != nil {
|
||||
return "", fmt.Errorf("invalid regex pattern '%s': %v", pattern, err)
|
||||
}
|
||||
filteredLog = re.ReplaceAllString(filteredLog, "")
|
||||
}
|
||||
return filteredLog, nil
|
||||
}
|
||||
|
||||
// KubectlApply applies a YAML file to the Kubernetes cluster and filters the logs
|
||||
func KubectlApply(ctx context.Context, filePath string) error {
|
||||
// Check if the file exists
|
||||
if _, err := os.Stat(filePath); os.IsNotExist(err) {
|
||||
return fmt.Errorf("file does not exist: %s", filePath)
|
||||
}
|
||||
// Check if the file exists
|
||||
if _, err := os.Stat(filePath); os.IsNotExist(err) {
|
||||
return fmt.Errorf("file does not exist: %s", filePath)
|
||||
}
|
||||
|
||||
// Read the YAML file
|
||||
yamlData, err := ioutil.ReadFile(filePath)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to read YAML file: %v", err)
|
||||
}
|
||||
// Read the YAML file
|
||||
yamlData, err := ioutil.ReadFile(filePath)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to read YAML file: %v", err)
|
||||
}
|
||||
|
||||
// Initialize logger and buffer
|
||||
log, buf, err := setupLogger()
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to set up logger: %v", err)
|
||||
}
|
||||
// Initialize logger and buffer
|
||||
log, buf, err := setupLogger()
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to set up logger: %v", err)
|
||||
}
|
||||
|
||||
// Initialize Kubernetes client
|
||||
k8sClient, err := getKubeClient()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// Initialize Kubernetes client
|
||||
k8sClient, err := getKubeClient()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Apply the YAML to the cluster
|
||||
if err := applyYAML(k8sClient, yamlData, log); err != nil {
|
||||
return fmt.Errorf("failed to apply YAML: %v", err)
|
||||
}
|
||||
// Apply the YAML to the cluster
|
||||
if err := applyYAML(k8sClient, yamlData, log); err != nil {
|
||||
return fmt.Errorf("failed to apply YAML: %v", err)
|
||||
}
|
||||
|
||||
// Get log output from buffer
|
||||
logOutput := buf.String()
|
||||
// Get log output from buffer
|
||||
logOutput := buf.String()
|
||||
|
||||
// Filter the logs
|
||||
filteredLog, err := filterLogOutput(logOutput)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to filter logs: %v", err)
|
||||
}
|
||||
// Filter the logs
|
||||
filteredLog, err := filterLogOutput(logOutput)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to filter logs: %v", err)
|
||||
}
|
||||
|
||||
// Output filtered logs
|
||||
fmt.Println(filteredLog)
|
||||
// Output filtered logs
|
||||
fmt.Println(filteredLog)
|
||||
|
||||
return nil
|
||||
return nil
|
||||
}
|
||||
|
||||
// KubectlApplyKustomize applies a kustomize directory or file to the Kubernetes cluster and filters the logs
|
||||
func KubectlApplyKustomize(ctx context.Context, filePath string) error {
|
||||
// Check if the path exists
|
||||
if _, err := os.Stat(filePath); os.IsNotExist(err) {
|
||||
return fmt.Errorf("path does not exist: %s", filePath)
|
||||
}
|
||||
// Check if the path exists
|
||||
if _, err := os.Stat(filePath); os.IsNotExist(err) {
|
||||
return fmt.Errorf("path does not exist: %s", filePath)
|
||||
}
|
||||
|
||||
// Determine if the path is a directory or a file
|
||||
fileInfo, err := os.Stat(filePath)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to stat path: %v", err)
|
||||
}
|
||||
// Determine if the path is a directory or a file
|
||||
fileInfo, err := os.Stat(filePath)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to stat path: %v", err)
|
||||
}
|
||||
|
||||
var kustomizePath string
|
||||
if fileInfo.IsDir() {
|
||||
// If it's a directory, use it as the kustomize path
|
||||
kustomizePath = filePath
|
||||
} else {
|
||||
// If it's a file, use its directory as the kustomize path
|
||||
kustomizePath = filepath.Dir(filePath)
|
||||
}
|
||||
var kustomizePath string
|
||||
if fileInfo.IsDir() {
|
||||
// If it's a directory, use it as the kustomize path
|
||||
kustomizePath = filePath
|
||||
} else {
|
||||
// If it's a file, use its directory as the kustomize path
|
||||
kustomizePath = filepath.Dir(filePath)
|
||||
}
|
||||
|
||||
// Process kustomize to get the YAML output
|
||||
fSys := filesys.MakeFsOnDisk()
|
||||
k := krusty.MakeKustomizer(krusty.MakeDefaultOptions())
|
||||
resMap, err := k.Run(fSys, kustomizePath)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to run kustomize: %v", err)
|
||||
}
|
||||
// Process kustomize to get the YAML output
|
||||
fSys := filesys.MakeFsOnDisk()
|
||||
k := krusty.MakeKustomizer(krusty.MakeDefaultOptions())
|
||||
resMap, err := k.Run(fSys, kustomizePath)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to run kustomize: %v", err)
|
||||
}
|
||||
|
||||
// Convert ResMap to YAML
|
||||
output, err := resMap.AsYaml()
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to convert ResMap to YAML: %v", err)
|
||||
}
|
||||
// Convert ResMap to YAML
|
||||
output, err := resMap.AsYaml()
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to convert ResMap to YAML: %v", err)
|
||||
}
|
||||
|
||||
// Initialize logger and buffer
|
||||
log, buf, err := setupLogger()
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to set up logger: %v", err)
|
||||
}
|
||||
// Initialize logger and buffer
|
||||
log, buf, err := setupLogger()
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to set up logger: %v", err)
|
||||
}
|
||||
|
||||
// Initialize Kubernetes client
|
||||
k8sClient, err := getKubeClient()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
// Initialize Kubernetes client
|
||||
k8sClient, err := getKubeClient()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Apply the YAML to the cluster
|
||||
if err := applyYAML(k8sClient, output, log); err != nil {
|
||||
return fmt.Errorf("failed to apply YAML: %v", err)
|
||||
}
|
||||
// Apply the YAML to the cluster
|
||||
if err := applyYAML(k8sClient, output, log); err != nil {
|
||||
return fmt.Errorf("failed to apply YAML: %v", err)
|
||||
}
|
||||
|
||||
// Get log output from buffer
|
||||
logOutput := buf.String()
|
||||
// Get log output from buffer
|
||||
logOutput := buf.String()
|
||||
|
||||
// Filter the logs
|
||||
filteredLog, err := filterLogOutput(logOutput)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to filter logs: %v", err)
|
||||
}
|
||||
// Filter the logs
|
||||
filteredLog, err := filterLogOutput(logOutput)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to filter logs: %v", err)
|
||||
}
|
||||
|
||||
// Output filtered logs
|
||||
fmt.Println(filteredLog)
|
||||
// Output filtered logs
|
||||
fmt.Println(filteredLog)
|
||||
|
||||
return nil
|
||||
return nil
|
||||
}
|
||||
|
||||
@@ -1,85 +1,85 @@
|
||||
package kubectlcmds
|
||||
|
||||
import (
|
||||
"context"
|
||||
"time"
|
||||
"context"
|
||||
"time"
|
||||
|
||||
"github.com/rs/zerolog/log"
|
||||
certificatesv1 "k8s.io/api/certificates/v1"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
"k8s.io/client-go/rest"
|
||||
"k8s.io/client-go/tools/clientcmd"
|
||||
"github.com/rs/zerolog/log"
|
||||
certificatesv1 "k8s.io/api/certificates/v1"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
"k8s.io/client-go/rest"
|
||||
"k8s.io/client-go/tools/clientcmd"
|
||||
)
|
||||
|
||||
// getClientset creates a Kubernetes clientset from the in-cluster config or kubeconfig file
|
||||
func GetClientset() (*kubernetes.Clientset, error) {
|
||||
// use the current context in kubeconfig
|
||||
config, err := clientcmd.BuildConfigFromFlags("", clientcmd.RecommendedHomeFile)
|
||||
if err != nil {
|
||||
config, err = rest.InClusterConfig()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
// use the current context in kubeconfig
|
||||
config, err := clientcmd.BuildConfigFromFlags("", clientcmd.RecommendedHomeFile)
|
||||
if err != nil {
|
||||
config, err = rest.InClusterConfig()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
}
|
||||
|
||||
// create the clientset
|
||||
clientset, err := kubernetes.NewForConfig(config)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
// create the clientset
|
||||
clientset, err := kubernetes.NewForConfig(config)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return clientset, nil
|
||||
return clientset, nil
|
||||
}
|
||||
|
||||
// Example function to approve pending CSRs
|
||||
func ApprovePendingCertificates(clientset *kubernetes.Clientset, stopCh <-chan struct{}) {
|
||||
log.Info().Msg("Waiting to approve certificates...")
|
||||
log.Info().Msg("Waiting to approve certificates...")
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-stopCh:
|
||||
log.Info().Msg("Stopping certificate approval...")
|
||||
return
|
||||
default:
|
||||
// Get the list of pending CSRs
|
||||
csrList, err := clientset.CertificatesV1().CertificateSigningRequests().List(context.TODO(), metav1.ListOptions{})
|
||||
if err != nil {
|
||||
log.Info().Msgf("Error getting CSRs: %v", err)
|
||||
time.Sleep(5 * time.Second)
|
||||
continue
|
||||
}
|
||||
for {
|
||||
select {
|
||||
case <-stopCh:
|
||||
log.Info().Msg("Stopping certificate approval...")
|
||||
return
|
||||
default:
|
||||
// Get the list of pending CSRs
|
||||
csrList, err := clientset.CertificatesV1().CertificateSigningRequests().List(context.TODO(), metav1.ListOptions{})
|
||||
if err != nil {
|
||||
log.Info().Msgf("Error getting CSRs: %v", err)
|
||||
time.Sleep(5 * time.Second)
|
||||
continue
|
||||
}
|
||||
|
||||
// Approve pending CSRs
|
||||
for _, csr := range csrList.Items {
|
||||
if csr.Status.Conditions == nil || len(csr.Status.Conditions) == 0 {
|
||||
// Create a copy of the CSR object
|
||||
csrCopy := csr.DeepCopy()
|
||||
// Approve pending CSRs
|
||||
for _, csr := range csrList.Items {
|
||||
if csr.Status.Conditions == nil || len(csr.Status.Conditions) == 0 {
|
||||
// Create a copy of the CSR object
|
||||
csrCopy := csr.DeepCopy()
|
||||
|
||||
// Prepare approval conditions
|
||||
conditions := []certificatesv1.CertificateSigningRequestCondition{
|
||||
{
|
||||
Type: certificatesv1.CertificateApproved,
|
||||
Reason: "AutoApproved",
|
||||
Message: "This CSR was approved automatically by controller.",
|
||||
LastUpdateTime: metav1.Now(),
|
||||
Status: "True",
|
||||
},
|
||||
}
|
||||
csrCopy.Status.Conditions = conditions
|
||||
// Prepare approval conditions
|
||||
conditions := []certificatesv1.CertificateSigningRequestCondition{
|
||||
{
|
||||
Type: certificatesv1.CertificateApproved,
|
||||
Reason: "AutoApproved",
|
||||
Message: "This CSR was approved automatically by controller.",
|
||||
LastUpdateTime: metav1.Now(),
|
||||
Status: "True",
|
||||
},
|
||||
}
|
||||
csrCopy.Status.Conditions = conditions
|
||||
|
||||
// Update approval for the CSR using the copied object
|
||||
_, err := clientset.CertificatesV1().CertificateSigningRequests().UpdateApproval(context.TODO(), csr.Name, csrCopy, metav1.UpdateOptions{})
|
||||
if err != nil {
|
||||
log.Info().Msgf("Error approving CSR %s: %v\n", csr.Name, err)
|
||||
} else {
|
||||
log.Info().Msgf("Approved CSR", csr.Name)
|
||||
}
|
||||
}
|
||||
}
|
||||
// Update approval for the CSR using the copied object
|
||||
_, err := clientset.CertificatesV1().CertificateSigningRequests().UpdateApproval(context.TODO(), csr.Name, csrCopy, metav1.UpdateOptions{})
|
||||
if err != nil {
|
||||
log.Info().Msgf("Error approving CSR %s: %v\n", csr.Name, err)
|
||||
} else {
|
||||
log.Info().Msgf("Approved CSR", csr.Name)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Sleep for 5 seconds before checking again
|
||||
time.Sleep(5 * time.Second)
|
||||
}
|
||||
}
|
||||
// Sleep for 5 seconds before checking again
|
||||
time.Sleep(5 * time.Second)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,77 +1,77 @@
|
||||
package kubectlcmds
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
"context"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/rs/zerolog/log"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
"k8s.io/client-go/tools/clientcmd"
|
||||
"github.com/rs/zerolog/log"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
"k8s.io/client-go/tools/clientcmd"
|
||||
)
|
||||
|
||||
func CheckStatus(requiredPods []string, excludePod []string, timeout time.Duration) error {
|
||||
// Load kubeconfig from the default location
|
||||
kubeconfig := clientcmd.NewDefaultClientConfigLoadingRules().GetDefaultFilename()
|
||||
config, err := clientcmd.BuildConfigFromFlags("", kubeconfig)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error loading kubeconfig: %w", err)
|
||||
}
|
||||
// Load kubeconfig from the default location
|
||||
kubeconfig := clientcmd.NewDefaultClientConfigLoadingRules().GetDefaultFilename()
|
||||
config, err := clientcmd.BuildConfigFromFlags("", kubeconfig)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error loading kubeconfig: %w", err)
|
||||
}
|
||||
|
||||
// Create clientset
|
||||
clientset, err := kubernetes.NewForConfig(config)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error creating clientset: %w", err)
|
||||
}
|
||||
// Create clientset
|
||||
clientset, err := kubernetes.NewForConfig(config)
|
||||
if err != nil {
|
||||
return fmt.Errorf("error creating clientset: %w", err)
|
||||
}
|
||||
|
||||
// Maximum duration to wait (15 minutes)
|
||||
maxDuration := timeout * time.Minute
|
||||
endTime := time.Now().Add(maxDuration)
|
||||
// Maximum duration to wait (15 minutes)
|
||||
maxDuration := timeout * time.Minute
|
||||
endTime := time.Now().Add(maxDuration)
|
||||
|
||||
for time.Now().Before(endTime) {
|
||||
// Get pods in all namespaces
|
||||
pods, err := clientset.CoreV1().Pods("").List(context.TODO(), metav1.ListOptions{})
|
||||
if err != nil {
|
||||
// return fmt.Errorf("error listing pods: %w", err)
|
||||
}
|
||||
for time.Now().Before(endTime) {
|
||||
// Get pods in all namespaces
|
||||
pods, err := clientset.CoreV1().Pods("").List(context.TODO(), metav1.ListOptions{})
|
||||
if err != nil {
|
||||
// return fmt.Errorf("error listing pods: %w", err)
|
||||
}
|
||||
|
||||
// Check if the required pods are both present and running
|
||||
requiredPodsMap := make(map[string]bool)
|
||||
for _, pod := range requiredPods {
|
||||
requiredPodsMap[pod] = false
|
||||
}
|
||||
// Check if the required pods are both present and running
|
||||
requiredPodsMap := make(map[string]bool)
|
||||
for _, pod := range requiredPods {
|
||||
requiredPodsMap[pod] = false
|
||||
}
|
||||
|
||||
for _, pod := range pods.Items {
|
||||
for _, requiredPod := range requiredPods {
|
||||
for _, excludePod := range excludePod {
|
||||
if strings.Contains(pod.Name, excludePod) {
|
||||
requiredPodsMap[requiredPod] = true
|
||||
}
|
||||
}
|
||||
if strings.Contains(pod.Name, requiredPod) && pod.Status.Phase == "Running" {
|
||||
requiredPodsMap[requiredPod] = true
|
||||
}
|
||||
}
|
||||
}
|
||||
for _, pod := range pods.Items {
|
||||
for _, requiredPod := range requiredPods {
|
||||
for _, excludePod := range excludePod {
|
||||
if strings.Contains(pod.Name, excludePod) {
|
||||
requiredPodsMap[requiredPod] = true
|
||||
}
|
||||
}
|
||||
if strings.Contains(pod.Name, requiredPod) && pod.Status.Phase == "Running" {
|
||||
requiredPodsMap[requiredPod] = true
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
allRunning := true
|
||||
for _, isRunning := range requiredPodsMap {
|
||||
if !isRunning {
|
||||
allRunning = false
|
||||
break
|
||||
}
|
||||
}
|
||||
allRunning := true
|
||||
for _, isRunning := range requiredPodsMap {
|
||||
if !isRunning {
|
||||
allRunning = false
|
||||
break
|
||||
}
|
||||
}
|
||||
|
||||
if allRunning {
|
||||
log.Info().Msg("All required pods are running")
|
||||
return nil
|
||||
}
|
||||
if allRunning {
|
||||
log.Info().Msg("All required pods are running")
|
||||
return nil
|
||||
}
|
||||
|
||||
// Wait for 5 seconds before checking again
|
||||
time.Sleep(5 * time.Second)
|
||||
}
|
||||
// Wait for 5 seconds before checking again
|
||||
time.Sleep(5 * time.Second)
|
||||
}
|
||||
|
||||
return fmt.Errorf("timeout: not all required pods are running after 15 minutes")
|
||||
return fmt.Errorf("timeout: not all required pods are running after 15 minutes")
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user