package podbackend import ( "context" "encoding/json" "fmt" "reflect" corev1 "k8s.io/api/core/v1" apierrors "k8s.io/apimachinery/pkg/api/errors" 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/client-go/dynamic" "k8s.io/client-go/kubernetes" "k8s.io/client-go/rest" ) var identityEntryResource = schema.GroupVersionResource{ Group: "spire.spiffe.io", Version: "v1alpha1", Resource: "clusterstaticentries", } // Client uses client-go's typed client for Pods and its dynamic client for the // SPIRE Operator CRD. type Client struct { Kubernetes kubernetes.Interface Dynamic dynamic.Interface } func NewClient(config *rest.Config) (*Client, error) { kubernetesClient, err := kubernetes.NewForConfig(config) if err != nil { return nil, fmt.Errorf("create Kubernetes client: %w", err) } dynamicClient, err := dynamic.NewForConfig(config) if err != nil { return nil, fmt.Errorf("create Kubernetes dynamic client: %w", err) } return &Client{Kubernetes: kubernetesClient, Dynamic: dynamicClient}, nil } func NewInClusterClient() (*Client, error) { config, err := rest.InClusterConfig() if err != nil { return nil, fmt.Errorf("load in-cluster Kubernetes config: %w", err) } return NewClient(config) } func (c *Client) ListPods(ctx context.Context, namespace, selector string) ([]Pod, error) { list, err := c.Kubernetes.CoreV1().Pods(namespace).List(ctx, metav1.ListOptions{LabelSelector: selector}) if err != nil { return nil, err } pods := make([]Pod, 0, len(list.Items)) for _, item := range list.Items { pods = append(pods, podFromKubernetes(item)) } return pods, nil } func (c *Client) CreatePod(ctx context.Context, manifest PodManifest) (Pod, error) { environment := make([]corev1.EnvVar, 0, len(manifest.Environment)) for name, value := range manifest.Environment { environment = append(environment, corev1.EnvVar{Name: name, Value: value}) } document := &corev1.Pod{ ObjectMeta: metav1.ObjectMeta{ Name: manifest.Name, Namespace: manifest.Namespace, Labels: manifest.Labels, Annotations: manifest.Annotations, }, Spec: corev1.PodSpec{ ServiceAccountName: manifest.ServiceAccount, RestartPolicy: corev1.RestartPolicyNever, Containers: []corev1.Container{{ Name: "executor", Image: manifest.Image, Args: manifest.Args, Env: environment, SecurityContext: &corev1.SecurityContext{Privileged: boolPointer(true)}, VolumeMounts: []corev1.VolumeMount{{ Name: "spire-agent-socket", MountPath: "/run/spire/agent-sockets", ReadOnly: true, }}, }}, Volumes: []corev1.Volume{{ Name: "spire-agent-socket", VolumeSource: corev1.VolumeSource{CSI: &corev1.CSIVolumeSource{ Driver: "csi.spiffe.io", ReadOnly: boolPointer(true), }}, }}, }, } created, err := c.Kubernetes.CoreV1().Pods(manifest.Namespace).Create(ctx, document, metav1.CreateOptions{}) if err != nil { return Pod{}, err } return podFromKubernetes(*created), nil } func (c *Client) DeletePod(ctx context.Context, namespace, name string) error { policy := metav1.DeletePropagationBackground err := c.Kubernetes.CoreV1().Pods(namespace).Delete(ctx, name, metav1.DeleteOptions{PropagationPolicy: &policy}) if apierrors.IsNotFound(err) { return nil } return err } func (c *Client) LabelPod(ctx context.Context, namespace, name string, labels map[string]string) error { patch, err := json.Marshal(map[string]any{"metadata": map[string]any{"labels": labels}}) if err != nil { return err } _, err = c.Kubernetes.CoreV1().Pods(namespace).Patch(ctx, name, types.MergePatchType, patch, metav1.PatchOptions{}) return err } func (c *Client) EnsureIdentityEntry(ctx context.Context, entry IdentityEntry) error { resource := c.Dynamic.Resource(identityEntryResource) existing, err := resource.Get(ctx, entry.Name, metav1.GetOptions{}) if err == nil { existingSpec, _, nestedErr := unstructured.NestedMap(existing.Object, "spec") if nestedErr != nil { return nestedErr } if !reflect.DeepEqual(existingSpec, identityEntryObject(entry).Object["spec"]) { return fmt.Errorf("identity entry %s exists with different selectors or SPIFFE ID", entry.Name) } return nil } if !apierrors.IsNotFound(err) { return err } _, err = resource.Create(ctx, identityEntryObject(entry), metav1.CreateOptions{}) return err } func (c *Client) DeleteIdentityEntry(ctx context.Context, name string) error { policy := metav1.DeletePropagationBackground err := c.Dynamic.Resource(identityEntryResource).Delete(ctx, name, metav1.DeleteOptions{PropagationPolicy: &policy}) if apierrors.IsNotFound(err) { return nil } return err } func identityEntryObject(entry IdentityEntry) *unstructured.Unstructured { selectors := make([]any, len(entry.Selectors)) for index, selector := range entry.Selectors { selectors[index] = selector } return &unstructured.Unstructured{Object: map[string]any{ "apiVersion": "spire.spiffe.io/v1alpha1", "kind": "ClusterStaticEntry", "metadata": map[string]any{ "name": entry.Name, "labels": stringMap(entry.Labels), }, "spec": map[string]any{ "className": entry.ClassName, "parentID": entry.ParentID, "spiffeID": entry.SPIFFEID, "selectors": selectors, }, }} } func podFromKubernetes(pod corev1.Pod) Pod { return Pod{ Name: pod.Name, Namespace: pod.Namespace, UID: string(pod.UID), Phase: string(pod.Status.Phase), Labels: pod.Labels, Annotations: pod.Annotations, } } func boolPointer(value bool) *bool { return &value } func stringMap(values map[string]string) map[string]any { result := make(map[string]any, len(values)) for key, value := range values { result[key] = value } return result }