194 lines
6.1 KiB
Go
194 lines
6.1 KiB
Go
package podbackend
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"reflect"
|
|
|
|
corev1 "k8s.io/api/core/v1"
|
|
apierrors "k8s.io/apimachinery/pkg/api/errors"
|
|
"k8s.io/apimachinery/pkg/api/resource"
|
|
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},
|
|
{Name: "docker-data", MountPath: "/var/lib/docker"},
|
|
},
|
|
}},
|
|
Volumes: []corev1.Volume{
|
|
{
|
|
Name: "spire-agent-socket",
|
|
VolumeSource: corev1.VolumeSource{CSI: &corev1.CSIVolumeSource{
|
|
Driver: "csi.spiffe.io", ReadOnly: boolPointer(true),
|
|
}},
|
|
},
|
|
{
|
|
Name: "docker-data",
|
|
VolumeSource: corev1.VolumeSource{EmptyDir: &corev1.EmptyDirVolumeSource{
|
|
SizeLimit: resourceQuantity("20Gi"),
|
|
}},
|
|
},
|
|
},
|
|
},
|
|
}
|
|
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 resourceQuantity(value string) *resource.Quantity {
|
|
quantity := resource.MustParse(value)
|
|
return &quantity
|
|
}
|
|
|
|
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
|
|
}
|