package k8s import ( "bytes" "context" "fmt" "io" batchv1 "k8s.io/api/batch/v1" corev1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/client-go/kubernetes" "k8s.io/client-go/tools/clientcmd" ) type k8sClient struct { clientset kubernetes.Interface } func NewClient(kubeconfig string) (K8sClient, error) { config, err := clientcmd.BuildConfigFromFlags("", kubeconfig) if err != nil { return nil, fmt.Errorf("building config: %w", err) } clientset, err := kubernetes.NewForConfig(config) if err != nil { return nil, fmt.Errorf("creating clientset: %w", err) } return &k8sClient{clientset: clientset}, nil } func NewClientFromInterface(clientset kubernetes.Interface) K8sClient { return &k8sClient{clientset: clientset} } func (c *k8sClient) CreateJob(spec *JobSpec) (*Job, error) { ctx := context.Background() envVars := make([]corev1.EnvVar, 0, len(spec.EnvVars)) for k, v := range spec.EnvVars { envVars = append(envVars, corev1.EnvVar{Name: k, Value: v}) } backoffLimit := int32(0) job := &batchv1.Job{ ObjectMeta: metav1.ObjectMeta{ Name: spec.Name, Namespace: spec.Namespace, }, Spec: batchv1.JobSpec{ BackoffLimit: &backoffLimit, Template: corev1.PodTemplateSpec{ Spec: corev1.PodSpec{ RestartPolicy: corev1.RestartPolicyNever, Containers: []corev1.Container{ { Name: "toolchain", Image: spec.Image, Env: envVars, }, }, }, }, }, } if spec.TimeoutSecond > 0 { timeout := int64(spec.TimeoutSecond) job.Spec.ActiveDeadlineSeconds = &timeout } result, err := c.clientset.BatchV1().Jobs(spec.Namespace).Create(ctx, job, metav1.CreateOptions{}) if err != nil { return nil, fmt.Errorf("creating job: %w", err) } return &Job{ Name: result.Name, Namespace: result.Namespace, Status: JobPending, }, nil } func (c *k8sClient) GetJob(name, namespace string) (*Job, error) { ctx := context.Background() result, err := c.clientset.BatchV1().Jobs(namespace).Get(ctx, name, metav1.GetOptions{}) if err != nil { return nil, fmt.Errorf("getting job: %w", err) } status := JobPending if result.Status.Succeeded > 0 { status = JobSucceeded } else if result.Status.Failed > 0 { status = JobFailed } else if result.Status.Active > 0 { status = JobRunning } job := &Job{ Name: result.Name, Namespace: result.Namespace, Status: status, } if status == JobFailed { logs, _ := c.getPodLogs(name, namespace) job.Logs = logs } return job, nil } func (c *k8sClient) DeleteJob(name, namespace string) error { ctx := context.Background() propagation := metav1.DeletePropagationForeground return c.clientset.BatchV1().Jobs(namespace).Delete(ctx, name, metav1.DeleteOptions{ PropagationPolicy: &propagation, }) } func (c *k8sClient) getPodLogs(jobName, namespace string) (string, error) { ctx := context.Background() pods, err := c.clientset.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{ LabelSelector: fmt.Sprintf("job-name=%s", jobName), }) if err != nil || len(pods.Items) == 0 { return "", err } pod := pods.Items[0] req := c.clientset.CoreV1().Pods(namespace).GetLogs(pod.Name, &corev1.PodLogOptions{}) stream, err := req.Stream(ctx) if err != nil { return "", err } defer stream.Close() var buf bytes.Buffer if _, err := io.Copy(&buf, stream); err != nil { return "", err } return buf.String(), nil }