basic pulling, pinning, and unpinning is working.
This commit is contained in:
@@ -1,13 +1,9 @@
|
||||
package main
|
||||
|
||||
type Image struct {
|
||||
Name string
|
||||
Pinned bool
|
||||
}
|
||||
|
||||
// ContainerRuntime defines the interface for managing containerd images
|
||||
type ContainerRuntime interface {
|
||||
List() ([]Image, error)
|
||||
// List images: returns map of container image to whether or not it is pinned
|
||||
List() (map[string]bool, error)
|
||||
Pin(imageRef string) error
|
||||
Unpin(imageRef string) error
|
||||
Pull(imageRef string) error
|
||||
|
||||
+38
-34
@@ -2,12 +2,9 @@ package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"crypto/rand"
|
||||
"encoding/hex"
|
||||
"crypto/md5"
|
||||
"encoding/base64"
|
||||
"fmt"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/containerd/containerd"
|
||||
"github.com/containerd/containerd/errdefs"
|
||||
"github.com/containerd/containerd/images"
|
||||
@@ -22,8 +19,8 @@ type Containerd struct {
|
||||
namespace string
|
||||
}
|
||||
|
||||
// NewImageManager creates a new Containerd
|
||||
func NewImageManager(socketPath, namespace string) (*Containerd, error) {
|
||||
// NewContainerd creates a new Containerd
|
||||
func NewContainerd(socketPath, namespace string) (*Containerd, error) {
|
||||
client, err := containerd.New(socketPath)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("failed to connect to containerd at %s: %w", socketPath, err)
|
||||
@@ -40,7 +37,7 @@ func NewImageManager(socketPath, namespace string) (*Containerd, error) {
|
||||
}
|
||||
|
||||
// List returns all images with their pinned status
|
||||
func (m *Containerd) List() ([]Image, error) {
|
||||
func (m *Containerd) List() (map[string]bool, error) {
|
||||
// Get all images
|
||||
images, err := m.client.ImageService().List(m.ctx)
|
||||
if err != nil {
|
||||
@@ -63,12 +60,9 @@ func (m *Containerd) List() ([]Image, error) {
|
||||
}
|
||||
|
||||
// Create the result list
|
||||
var result []Image
|
||||
var result = make(map[string]bool)
|
||||
for _, img := range images {
|
||||
result = append(result, Image{
|
||||
Name: img.Name,
|
||||
Pinned: pinnedImages[img.Name],
|
||||
})
|
||||
result[img.Name] = pinnedImages[img.Name]
|
||||
}
|
||||
|
||||
return result, nil
|
||||
@@ -77,7 +71,7 @@ func (m *Containerd) List() ([]Image, error) {
|
||||
// Pin creates a lease for an image to prevent garbage collection
|
||||
func (m *Containerd) Pin(imageRef string) error {
|
||||
// Create a unique lease ID based on image reference
|
||||
leaseID := fmt.Sprintf("pin-%s-%d", generateID(), time.Now().UnixNano())
|
||||
leaseID := fmt.Sprintf("pin-%s", generateID(imageRef))
|
||||
|
||||
// Get the image to validate it exists
|
||||
_, err := m.client.ImageService().Get(m.ctx, imageRef)
|
||||
@@ -85,6 +79,14 @@ func (m *Containerd) Pin(imageRef string) error {
|
||||
return fmt.Errorf("failed to get image %s: %w", imageRef, err)
|
||||
}
|
||||
|
||||
leaseList, err := m.findLeases(imageRef)
|
||||
if err != nil {
|
||||
return fmt.Errorf("Failed to get leases for image %s: %v", imageRef, err)
|
||||
}
|
||||
if len(leaseList) > 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
// Create a new lease
|
||||
opts := []leases.Opt{
|
||||
leases.WithID(leaseID),
|
||||
@@ -103,30 +105,34 @@ func (m *Containerd) Pin(imageRef string) error {
|
||||
|
||||
// Unpin removes a lease for an image allowing garbage collection
|
||||
func (m *Containerd) Unpin(imageRef string) error {
|
||||
leases, err := m.findLeases(imageRef)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
for _, lease := range leases {
|
||||
if err := m.client.LeasesService().Delete(m.ctx, lease); err != nil {
|
||||
return fmt.Errorf("failed to delete lease %s: %w", lease.ID, err)
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *Containerd) findLeases(imageRef string) ([]leases.Lease, error) {
|
||||
// List all leases
|
||||
leaseList, err := m.client.LeasesService().List(m.ctx)
|
||||
if err != nil {
|
||||
return fmt.Errorf("failed to list leases: %w", err)
|
||||
return nil, fmt.Errorf("failed to list leases: %w", err)
|
||||
}
|
||||
|
||||
// Find leases that reference our image
|
||||
var foundLease bool
|
||||
var leases = make([]leases.Lease, 0)
|
||||
for _, lease := range leaseList {
|
||||
// Check if this lease has a label referencing our image
|
||||
if label, ok := lease.Labels["containerd.io/gc.ref.content.image"]; ok && label == imageRef {
|
||||
// Delete this lease
|
||||
if err := m.client.LeasesService().Delete(m.ctx, lease); err != nil {
|
||||
return fmt.Errorf("failed to delete lease %s: %w", lease.ID, err)
|
||||
}
|
||||
foundLease = true
|
||||
leases = append(leases, lease)
|
||||
}
|
||||
}
|
||||
|
||||
if !foundLease {
|
||||
return fmt.Errorf("no lease found for image %s", imageRef)
|
||||
}
|
||||
|
||||
return nil
|
||||
return leases, nil
|
||||
}
|
||||
|
||||
// Pull pulls an image from a registry
|
||||
@@ -173,11 +179,9 @@ func (m *Containerd) Close() error {
|
||||
}
|
||||
|
||||
// generateID creates a random unique ID
|
||||
func generateID() string {
|
||||
b := make([]byte, 16)
|
||||
if _, err := rand.Read(b); err != nil {
|
||||
// Fall back to a timestamp-based ID in the unlikely event of rand failure
|
||||
return strings.Replace(fmt.Sprintf("%d", time.Now().UnixNano()), "-", "", -1)
|
||||
}
|
||||
return hex.EncodeToString(b)
|
||||
func generateID(image string) string {
|
||||
|
||||
md5Hash := md5.Sum([]byte(image))
|
||||
md5Base64 := base64.StdEncoding.EncodeToString(md5Hash[:])
|
||||
return md5Base64
|
||||
}
|
||||
|
||||
@@ -0,0 +1,29 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"k8s.io/client-go/kubernetes"
|
||||
"k8s.io/client-go/tools/clientcmd"
|
||||
"log"
|
||||
)
|
||||
|
||||
func GetKubernetesConnection() (*kubernetes.Clientset, string) {
|
||||
loadingRules := clientcmd.NewDefaultClientConfigLoadingRules()
|
||||
configOverrides := &clientcmd.ConfigOverrides{}
|
||||
kubeConfig := clientcmd.NewNonInteractiveDeferredLoadingClientConfig(loadingRules, configOverrides)
|
||||
|
||||
config, err := kubeConfig.ClientConfig()
|
||||
if err != nil {
|
||||
log.Panicln(err.Error())
|
||||
}
|
||||
|
||||
clientset, err := kubernetes.NewForConfig(config)
|
||||
if err != nil {
|
||||
log.Panicln(err.Error())
|
||||
}
|
||||
|
||||
namespace, _, err := kubeConfig.Namespace()
|
||||
if err != nil {
|
||||
log.Panicf("Could not get namespace")
|
||||
}
|
||||
return clientset, namespace
|
||||
}
|
||||
+31
-39
@@ -1,53 +1,45 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"log"
|
||||
goflags "flag"
|
||||
"github.com/spf13/cobra"
|
||||
"k8s.io/klog/v2"
|
||||
"os"
|
||||
)
|
||||
|
||||
func main() {
|
||||
// Default containerd socket path
|
||||
socketPath := "/run/containerd/containerd.sock"
|
||||
// Default namespace for Docker with containerd
|
||||
namespace := "moby"
|
||||
klogFlags := goflags.NewFlagSet("", goflags.PanicOnError)
|
||||
klog.InitFlags(klogFlags)
|
||||
|
||||
// Create the image manager
|
||||
imgManager, err := NewImageManager(socketPath, namespace)
|
||||
if err != nil {
|
||||
log.Fatalf("Failed to create image manager: %v", err)
|
||||
var kubernetesNamespace string
|
||||
var socketPath string
|
||||
var containerdNamespace string
|
||||
|
||||
cmd := &cobra.Command{
|
||||
Use: "kube-fetcher",
|
||||
Short: "Fetch images on a kubernetes node",
|
||||
Long: `
|
||||
Queries k8s for all running pods and makes sure that all
|
||||
images referenced in pods are made available on the local k8s node and pinned
|
||||
so they don't get garbage collected'`,
|
||||
RunE: func(cmd *cobra.Command, args []string) error {
|
||||
err := pullAndPin(kubernetesNamespace, socketPath, containerdNamespace)
|
||||
return err
|
||||
},
|
||||
}
|
||||
|
||||
img := "docker.io/library/alpine:3.20.3"
|
||||
err = imgManager.Pull(img)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
cmd.PersistentFlags().StringVar(&kubernetesNamespace, "kubernetes-namespace",
|
||||
"", "Kubernetes containerdNamespace to inspect (default is all namespaces)")
|
||||
cmd.PersistentFlags().StringVar(&socketPath, "socket",
|
||||
"/run/containerd/containerd.sock", "Containerd socket")
|
||||
cmd.PersistentFlags().StringVar(&containerdNamespace, "containerd-namespace",
|
||||
"k8s.io", "Containerd namespace to use")
|
||||
cmd.Flags().AddGoFlagSet(klogFlags)
|
||||
|
||||
imgs, err := imgManager.List()
|
||||
err := cmd.Execute()
|
||||
if err != nil {
|
||||
panic(err)
|
||||
klog.Errorf("Error: %v", err)
|
||||
os.Exit(1)
|
||||
}
|
||||
fmt.Printf("Images: %v\n", imgs)
|
||||
|
||||
err = imgManager.Pin(img)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
imgs, err = imgManager.List()
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
fmt.Printf("Images: %v\n", imgs)
|
||||
|
||||
err = imgManager.Unpin(img)
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
||||
imgs, err = imgManager.List()
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
fmt.Printf("Images: %v\n", imgs)
|
||||
|
||||
}
|
||||
|
||||
@@ -0,0 +1,140 @@
|
||||
package main
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/client-go/kubernetes"
|
||||
"k8s.io/klog/v2"
|
||||
"strings"
|
||||
)
|
||||
|
||||
func canonicalizeImageName(image string) string {
|
||||
// Check if image has a tag
|
||||
var name, tag string
|
||||
if parts := strings.Split(image, ":"); len(parts) > 1 {
|
||||
name = parts[0]
|
||||
tag = parts[1]
|
||||
} else {
|
||||
name = image
|
||||
tag = "latest" // Default tag
|
||||
}
|
||||
|
||||
// Check if image has a host
|
||||
var host, remainder string
|
||||
if strings.Contains(name, "/") {
|
||||
parts := strings.SplitN(name, "/", 2)
|
||||
|
||||
// Determine if the first part is a host
|
||||
// A host must contain a dot or colon (for registries with ports)
|
||||
if strings.Contains(parts[0], ".") || strings.Contains(parts[0], ":") {
|
||||
host = parts[0]
|
||||
remainder = parts[1]
|
||||
} else {
|
||||
// No host specified, use default docker.io
|
||||
host = "docker.io"
|
||||
remainder = name
|
||||
}
|
||||
} else {
|
||||
// No host and no /, use default docker.io and library
|
||||
host = "docker.io"
|
||||
remainder = "library/" + name
|
||||
}
|
||||
|
||||
// Handle the case when remainder doesn't specify library but it's not a docker.io official image
|
||||
if host == "docker.io" && !strings.Contains(remainder, "/") {
|
||||
remainder = "library/" + remainder
|
||||
} else if host == "docker.io" && !strings.HasPrefix(remainder, "library/") {
|
||||
// Check if it's not already in the library
|
||||
if parts := strings.SplitN(remainder, "/", 2); len(parts) == 2 && parts[0] != "library" {
|
||||
// Not a library image and not already properly formatted
|
||||
// This is a docker.io user image (e.g., user/image)
|
||||
// We keep it as is
|
||||
}
|
||||
}
|
||||
|
||||
return fmt.Sprintf("%s/%s:%s", host, remainder, tag)
|
||||
}
|
||||
|
||||
func getContainers(clientset *kubernetes.Clientset, kubernetesNamespace string) map[string]bool {
|
||||
pods, err := clientset.CoreV1().Pods(kubernetesNamespace).List(context.Background(),
|
||||
metav1.ListOptions{})
|
||||
if err != nil {
|
||||
panic(err)
|
||||
}
|
||||
|
||||
containers := make(map[string]bool)
|
||||
|
||||
for _, pod := range pods.Items {
|
||||
klog.V(3).Infof("%s/%s\n", pod.Namespace, pod.Name)
|
||||
for _, container := range pod.Spec.InitContainers {
|
||||
klog.V(3).Infof(" %s\n", container.Image)
|
||||
containers[canonicalizeImageName(container.Image)] = true
|
||||
}
|
||||
for _, container := range pod.Spec.Containers {
|
||||
klog.V(3).Infof(" %s\n", container.Image)
|
||||
containers[canonicalizeImageName(container.Image)] = true
|
||||
}
|
||||
}
|
||||
return containers
|
||||
}
|
||||
|
||||
func pullAndPin(kubernetesNamespace, socketPath, containerdNamespace string) error {
|
||||
|
||||
// Create the image manager
|
||||
containerd, err := NewContainerd(socketPath, containerdNamespace)
|
||||
if err != nil {
|
||||
klog.Fatalf("Failed to create image manager: %v", err)
|
||||
}
|
||||
|
||||
imgs, err := containerd.List()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
clientset, _ := GetKubernetesConnection()
|
||||
|
||||
containers := getContainers(clientset, kubernetesNamespace)
|
||||
for container, _ := range containers {
|
||||
klog.V(3).Infof("Found container %s\n", container)
|
||||
}
|
||||
|
||||
// unpin images that are not used
|
||||
for container, pinned := range imgs {
|
||||
if !containers[container] && pinned {
|
||||
klog.Infof("Unpinning %s\n", container)
|
||||
err := containerd.Unpin(container)
|
||||
if err != nil {
|
||||
klog.Warningf(" error: %v", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Pull images that are used
|
||||
for container := range containers {
|
||||
if _, found := imgs[container]; !found {
|
||||
klog.Infof("Pulling %s\n", container)
|
||||
err := containerd.Pull(container)
|
||||
if err != nil {
|
||||
klog.Warningf("error: %v", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
imgs, err = containerd.List()
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
// Pin images that are used and present
|
||||
for container := range containers {
|
||||
if pinned, found := imgs[container]; found && !pinned {
|
||||
klog.Infof("Pinning %s\n", container)
|
||||
err := containerd.Pin(container)
|
||||
if err != nil {
|
||||
klog.Warningf(" error: %v", err)
|
||||
}
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
Reference in New Issue
Block a user