diff --git a/cmd/cupdate/config.go b/cmd/cupdate/config.go new file mode 100644 index 00000000..09aec9e5 --- /dev/null +++ b/cmd/cupdate/config.go @@ -0,0 +1,176 @@ +package main + +import ( + "encoding/base64" + "encoding/json" + "fmt" + "log/slog" + "os" + "path/filepath" + "strings" + "time" + + "github.com/AlexGustafsson/cupdate/internal/httputil" + "github.com/AlexGustafsson/cupdate/internal/platform/docker" + "github.com/caarlos0/env/v11" + "k8s.io/client-go/rest" +) + +type Config struct { + Log struct { + Level string `env:"LEVEL" envDefault:"info"` + } `envPrefix:"LOG_"` + + API struct { + Address string `env:"ADDRESS" envDefault:"0.0.0.0"` + Port uint16 `env:"PORT" envDefault:"8080"` + } `envPrefix:"API_"` + + Web struct { + Disabled bool `env:"DISABLED"` + Address string `env:"ADDRESS"` + } `envPrefix:"WEB_"` + + HTTP struct { + UserAgent string `env:"USER_AGENT" envDefault:"Cupdate/1.0"` + } `envPrefix:"HTTP_"` + + Cache struct { + Path string `env:"PATH" envDefault:"cachev1.boltdb"` + MaxAge time.Duration `env:"MAX_AGE" envDefault:"24h"` + } `envPrefix:"CACHE_"` + + Database struct { + Path string `env:"PATH" envDefault:"dbv1.sqlite"` + } `envPrefix:"DB_"` + + Processing struct { + Interval time.Duration `env:"INTERVAL" envDefault:"1h"` + Items int `env:"ITEMS" envDefault:"10"` + MinAge time.Duration `env:"MIN_AGE" envDefault:"72h"` + Timeout time.Duration `env:"TIMEOUT" envDefault:"2m"` + QueueSize int `env:"QUEUE_SIZE" envDefault:"50"` + QueueBurst int `env:"QUEUE_BURST" envDefault:"10"` + QueueRate time.Duration `env:"QUEUE_RATE" envDefault:"1m"` + } `envPrefix:"PROCESSING_"` + + Workflow struct { + CleanupMaxAge time.Duration `env:"CLEANUP_MAX_AGE" envDefault:"48h"` + CleanupInterval time.Duration `env:"CLEANUP_INTERVAL" envDefault:"1h"` + } `envPrefix:"WORKFLOW_"` + + Kubernetes struct { + Host string `env:"HOST"` + IncludeOldReplicaSets bool `env:"INCLUDE_OLD_REPLICAS"` + } `envPrefix:"KUBERNETES_"` + + Docker struct { + Hosts []string `env:"HOST"` + IncludeAllContainers bool `env:"INCLUDE_ALL_CONTAINERS"` + TLSPath string `env:"TLS_PATH"` + } `envPrefix:"DOCKER_"` + + OTEL struct { + Target string `env:"TARGET"` + Insecure bool `env:"INSECURE"` + } `envPrefix:"OTEL_"` + + Registry struct { + Secrets string `env:"SECRETS"` + } `envPrefix:"REGISTRY_"` +} + +func (c *Config) LogLevel() (slog.Level, error) { + var logLevel slog.Level + + switch c.Log.Level { + case "debug": + logLevel = slog.LevelDebug + case "info": + logLevel = slog.LevelInfo + case "warn": + logLevel = slog.LevelWarn + case "error": + logLevel = slog.LevelError + default: + return logLevel, fmt.Errorf("invalid log level") + } + + return logLevel, nil +} + +func (c *Config) RegistryAuth() (*httputil.AuthMux, error) { + registryAuth := httputil.NewAuthMux() + + if c.Registry.Secrets != "" { + file, err := os.Open(c.Registry.Secrets) + if err != nil { + return nil, fmt.Errorf("failed to read registry secrets: %w", err) + } + + var dockerConfig *docker.ConfigFile + err = json.NewDecoder(file).Decode(&dockerConfig) + file.Close() + if err != nil { + return nil, fmt.Errorf("failed to parse registry secrets: %w", err) + } + + for k, v := range dockerConfig.HttpHeaders { + registryAuth.SetHeader(k, v) + } + + for pattern, auth := range dockerConfig.Auths { + if auth.Auth == "" { + registryAuth.Handle(pattern, httputil.BasicAuthHandler{ + Username: auth.Username, + Password: auth.Password, + }) + } else { + value, err := base64.StdEncoding.DecodeString(auth.Auth) + if err != nil { + return nil, fmt.Errorf("invalid registry secrets file: %w", err) + } + + username, password, ok := strings.Cut(string(value), ":") + if !ok { + return nil, fmt.Errorf("invalid registry secrets file: invalid auth field") + } + + registryAuth.Handle(pattern, httputil.BasicAuthHandler{ + Username: username, + Password: password, + }) + } + } + } + + return registryAuth, nil +} + +func (c *Config) KubernetesClientConfig() (*rest.Config, error) { + if c.Kubernetes.Host == "" { + return rest.InClusterConfig() + } + + return &rest.Config{ + Host: c.Kubernetes.Host, + }, nil +} + +func (c *Config) DatabaseURI() (string, error) { + absolutePath, err := filepath.Abs(c.Database.Path) + if err != nil { + return "", err + } + + return "file://" + absolutePath, nil +} + +func ParseConfigFromEnv() (*Config, error) { + var config Config + if err := env.ParseWithOptions(&config, env.Options{Prefix: "CUPDATE_"}); err != nil { + return nil, err + } + + return &config, nil +} diff --git a/cmd/cupdate/log.go b/cmd/cupdate/log.go new file mode 100644 index 00000000..e26e8420 --- /dev/null +++ b/cmd/cupdate/log.go @@ -0,0 +1,21 @@ +package main + +import ( + "log/slog" + "os" + + "github.com/AlexGustafsson/cupdate/internal/slogutil" +) + +func InitDefaultLogger() { + options := &slog.HandlerOptions{ + Level: slog.LevelInfo, + } + handler := slogutil.NewHandler(os.Stderr, options) + + logger := slog.New(handler). + With(slog.String("service.version", Version)). + With(slog.String("service.name", "cupdate")) + + slog.SetDefault(logger) +} diff --git a/cmd/cupdate/main.go b/cmd/cupdate/main.go index 42dba9ad..f90d2906 100644 --- a/cmd/cupdate/main.go +++ b/cmd/cupdate/main.go @@ -2,548 +2,207 @@ package main import ( "context" - "encoding/json" + "errors" "fmt" "log/slog" "net/http" - "net/url" "os" "os/signal" - "path/filepath" "strings" + "sync" "time" "github.com/AlexGustafsson/cupdate/internal/api" "github.com/AlexGustafsson/cupdate/internal/cache" - "github.com/AlexGustafsson/cupdate/internal/configutils" "github.com/AlexGustafsson/cupdate/internal/httputil" - "github.com/AlexGustafsson/cupdate/internal/models" "github.com/AlexGustafsson/cupdate/internal/oci" "github.com/AlexGustafsson/cupdate/internal/otelutil" - "github.com/AlexGustafsson/cupdate/internal/platform" - "github.com/AlexGustafsson/cupdate/internal/platform/docker" - "github.com/AlexGustafsson/cupdate/internal/platform/kubernetes" "github.com/AlexGustafsson/cupdate/internal/slogutil" "github.com/AlexGustafsson/cupdate/internal/store" "github.com/AlexGustafsson/cupdate/internal/web" "github.com/AlexGustafsson/cupdate/internal/worker" - "github.com/caarlos0/env/v11" "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus/promhttp" - "golang.org/x/sync/errgroup" - "k8s.io/client-go/rest" ) -// Version is the current Cupdate version. Overwritten at build time. -var Version string = "development build" - -type Config struct { - Log struct { - Level string `env:"LEVEL" envDefault:"info"` - } `envPrefix:"LOG_"` - - API struct { - Address string `env:"ADDRESS" envDefault:"0.0.0.0"` - Port uint16 `env:"PORT" envDefault:"8080"` - } `envPrefix:"API_"` - - Web struct { - Disabled bool `env:"DISABLED"` - Address string `env:"ADDRESS"` - } `envPrefix:"WEB_"` - - HTTP struct { - UserAgent string `env:"USER_AGENT" envDefault:"Cupdate/1.0"` - } `envPrefix:"HTTP_"` - - Cache struct { - Path string `env:"PATH" envDefault:"cachev1.boltdb"` - MaxAge time.Duration `env:"MAX_AGE" envDefault:"24h"` - } `envPrefix:"CACHE_"` - - Database struct { - Path string `env:"PATH" envDefault:"dbv1.sqlite"` - } `envPrefix:"DB_"` - - Processing struct { - Interval time.Duration `env:"INTERVAL" envDefault:"1h"` - Items int `env:"ITEMS" envDefault:"10"` - MinAge time.Duration `env:"MIN_AGE" envDefault:"72h"` - Timeout time.Duration `env:"TIMEOUT" envDefault:"2m"` - QueueSize int `env:"QUEUE_SIZE" envDefault:"50"` - QueueBurst int `env:"QUEUE_BURST" envDefault:"10"` - QueueRate time.Duration `env:"QUEUE_RATE" envDefault:"1m"` - } `envPrefix:"PROCESSING_"` - - Workflow struct { - CleanupMaxAge time.Duration `env:"CLEANUP_MAX_AGE" envDefault:"48h"` - CleanupInterval time.Duration `env:"CLEANUP_INTERVAL" envDefault:"1h"` - } `envPrefix:"WORKFLOW_"` - - Kubernetes struct { - Host string `env:"HOST"` - IncludeOldReplicaSets bool `env:"INCLUDE_OLD_REPLICAS"` - } `envPrefix:"KUBERNETES_"` - - Docker struct { - Hosts []string `env:"HOST"` - IncludeAllContainers bool `env:"INCLUDE_ALL_CONTAINERS"` - TLSPath string `env:"TLS_PATH"` - } `envPrefix:"DOCKER_"` - - OTEL struct { - Target string `env:"TARGET"` - Insecure bool `env:"INSECURE"` - } `envPrefix:"OTEL_"` - - Registry struct { - Secrets string `env:"SECRETS"` - } `envPrefix:"REGISTRY_"` -} +// ErrNonZeroExitCode is a special error which, when caught in main, will exit +// with a non-zero exit code without logging the error, assuming it was already +// handled. +var ErrNonZeroExitCode = errors.New("unexpected error") func main() { - slog.SetDefault(slog.New(slogutil.NewHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelError})).With(slog.String("service.version", Version)).With(slog.String("service.name", "cupdate"))) + InitDefaultLogger() - var config Config - err := env.ParseWithOptions(&config, env.Options{ - Prefix: "CUPDATE_", - }) - if err != nil { - slog.Error("Failed to parse config from environment variables", slog.Any("error", err)) - os.Exit(1) - } + runCtx, cancelRunCtx := context.WithCancel(context.Background()) - var logLevel slog.Level - switch config.Log.Level { - case "debug": - logLevel = slog.LevelDebug - case "info": - logLevel = slog.LevelInfo - case "warn": - logLevel = slog.LevelWarn - case "error": - logLevel = slog.LevelError - default: - slog.Error("Failed to parse config - invalid log level") - os.Exit(1) + var mutex sync.Mutex + shutdownFuncs := []func(context.Context){} + registerShutdownFunc := func(f func(context.Context)) { + mutex.Lock() + defer mutex.Unlock() + + shutdownFuncs = append(shutdownFuncs, f) } - slog.SetDefault(slog.New(slogutil.NewHandler(os.Stderr, &slog.HandlerOptions{Level: logLevel})).With(slog.String("service.version", Version)).With(slog.String("service.name", "cupdate"))) - slog.Debug("Parsed config", slog.Any("config", config)) + signals := make(chan os.Signal, 1) + signal.Notify(signals, os.Interrupt) + caught := 0 + go func() { + for range signals { + caught++ + if caught == 1 { + slog.Info("Caught signal, exiting gracefully") - registryAuth := httputil.NewAuthMux() - if config.Registry.Secrets != "" { - file, err := os.Open(config.Registry.Secrets) - if err != nil { - slog.Error("Failed to read registry secrets", slog.Any("error", err)) - os.Exit(1) - } + shutdownCtx, cancelShutdownCtx := context.WithTimeout(context.Background(), 15*time.Second) + mutex.Lock() + for _, shutdownFunc := range shutdownFuncs { + shutdownFunc(shutdownCtx) + } + cancelShutdownCtx() - var dockerConfig *docker.ConfigFile - err = json.NewDecoder(file).Decode(&dockerConfig) - file.Close() - if err != nil { - slog.Error("Failed to parse registry secrets", slog.Any("error", err)) - os.Exit(1) + // Cancel the run context last as to make sure the cleanup funcs were + // run + cancelRunCtx() + } else { + slog.Warn("Caught signal, exiting now") + os.Exit(1) + } } + }() - for k, v := range dockerConfig.HttpHeaders { - registryAuth.SetHeader(k, v) + if err := run(runCtx, registerShutdownFunc); err != nil { + if err != ErrNonZeroExitCode { + slog.Error("Exiting with non-zero exit code", slog.Any("error", err)) } + os.Exit(1) + } +} - for pattern, auth := range dockerConfig.Auths { - if auth.Auth == "" { - registryAuth.Handle(pattern, httputil.BasicAuthHandler{ - Username: auth.Username, - Password: auth.Password, - }) - } else { - registryAuth.Handle(pattern, httputil.BearerToken(auth.Auth)) - } - } +func run(ctx context.Context, registerShutdownFunc func(func(context.Context))) error { + slog.SetDefault(slog.New(slogutil.NewHandler(os.Stderr, &slog.HandlerOptions{Level: slog.LevelError})).With(slog.String("service.version", Version)).With(slog.String("service.name", "cupdate"))) + + config, err := ParseConfigFromEnv() + if err != nil { + slog.Error("Failed to parse config from environment variables", slog.Any("error", err)) + return ErrNonZeroExitCode + } + + logLevel, err := config.LogLevel() + if err != nil { + slog.Error("Failed to parse config - invalid log level") + return ErrNonZeroExitCode } + slog.SetLogLoggerLevel(logLevel) - ctx, cancel := context.WithCancel(context.Background()) + slog.Debug("Parsed config", slog.Any("config", config)) + + registryAuth, err := config.RegistryAuth() + if err != nil { + slog.Error("Failed to parse registry auth", slog.Any("error", err)) + return ErrNonZeroExitCode + } if config.OTEL.Target != "" { shutdown, err := otelutil.Init(ctx, config.OTEL.Target, config.OTEL.Insecure) if err != nil { slog.ErrorContext(ctx, "Failed to initialize otel", slog.Any("error", err)) - os.Exit(1) + return ErrNonZeroExitCode } - // TODO: This won't be invoked on exit, make it part of the shutdown - // procedure defer shutdown(ctx) } - // Set up the configured platform (Docker if specified, auto discovery of - // Kubernetes otherwise) - var targetPlatform platform.Grapher - if len(config.Docker.Hosts) == 0 { - var kubernetesConfig *rest.Config - if config.Kubernetes.Host == "" { - var err error - kubernetesConfig, err = rest.InClusterConfig() - if err != nil { - slog.ErrorContext(ctx, "Failed to configure Kubernetes client", slog.Any("error", err)) - os.Exit(1) - } - } else { - kubernetesConfig = &rest.Config{ - Host: config.Kubernetes.Host, - } - } - - targetPlatform, err = kubernetes.NewPlatform(kubernetesConfig, &kubernetes.Options{IncludeOldReplicaSets: config.Kubernetes.IncludeOldReplicaSets}) - if err != nil { - slog.ErrorContext(ctx, "Failed to create kubernetes source", slog.Any("error", err)) - os.Exit(1) - } - } else { - graphers := make([]platform.Grapher, 0) - for _, host := range config.Docker.Hosts { - options := &docker.Options{ - IncludeAllContainers: config.Docker.IncludeAllContainers, - } - - if config.Docker.TLSPath != "" { - uri, err := url.Parse(host) - if err != nil { - slog.ErrorContext(ctx, "Failed to parse docker URI", slog.Any("error", err), slog.String("host", host)) - os.Exit(1) - } - - tlsConfig, err := configutils.LoadTLSConfig( - filepath.Join(config.Docker.TLSPath, uri.Hostname()), - config.Docker.TLSPath, - ) - if err != nil { - slog.ErrorContext(ctx, "Failed to read docker TLS files", slog.Any("error", err), slog.String("host", host)) - } - - options.TLSClientConfig = tlsConfig - } - - platform, err := docker.NewPlatform(ctx, host, options) - if err != nil { - slog.ErrorContext(ctx, "Failed to create docker source", slog.Any("error", err), slog.String("host", host)) - os.Exit(1) - } - - graphers = append(graphers, platform) - } - targetPlatform = &platform.CompoundGrapher{ - Graphers: graphers, - } + targetPlatform, err := ConfigurePlatformGrapher(ctx, config) + if err != nil { + slog.Error("Failed to configure platform", slog.Any("error", err)) + return ErrNonZeroExitCode } cache, err := cache.NewDiskCache(config.Cache.Path) if err != nil { - slog.ErrorContext(ctx, "Failed to create disk cache", slog.Any("error", err)) - os.Exit(1) + slog.Error("Failed to create disk cache", slog.Any("error", err)) + return ErrNonZeroExitCode + } + defer func() { + if err := cache.Close(); err != nil { + slog.Error("Failed to close cache", slog.Any("error", err)) + } + }() + if err := prometheus.DefaultRegisterer.Register(cache); err != nil { + slog.Error("Failed to register prometheus metrics for cache", slog.Any("error", err)) + return ErrNonZeroExitCode } - prometheus.DefaultRegisterer.MustRegister(cache) - absoluteDatabasePath, err := filepath.Abs(config.Database.Path) + databaseURI, err := config.DatabaseURI() if err != nil { - slog.ErrorContext(ctx, "Failed to resolve database path", slog.Any("error", err)) - os.Exit(1) + slog.Error("Failed to resolve database path", slog.Any("error", err)) + return ErrNonZeroExitCode } - if err := store.Initialize(ctx, "file://"+absoluteDatabasePath); err != nil { - slog.ErrorContext(ctx, "Failed to initialize database", slog.Any("error", err)) - os.Exit(1) + if err := store.Initialize(ctx, databaseURI); err != nil { + slog.Error("Failed to initialize database", slog.Any("error", err)) + return ErrNonZeroExitCode } - readStore, err := store.New("file://"+absoluteDatabasePath, true) + readStore, err := store.New(databaseURI, true) if err != nil { - slog.ErrorContext(ctx, "Failed to load database", slog.Any("error", err)) - os.Exit(1) + slog.Error("Failed to open database for reading", slog.Any("error", err)) + return ErrNonZeroExitCode } - writeStore, err := store.New("file://"+absoluteDatabasePath, false) + defer func() { + if err := readStore.Close(); err != nil { + slog.Error("Failed to close read-only database", slog.Any("error", err)) + } + }() + + writeStore, err := store.New(databaseURI, false) if err != nil { - slog.ErrorContext(ctx, "Failed to load database", slog.Any("error", err)) - os.Exit(1) + slog.Error("Failed to open database for writing", slog.Any("error", err)) + return ErrNonZeroExitCode } + defer func() { + if err := writeStore.Close(); err != nil { + slog.Error("Failed to close writable database", slog.Any("error", err)) + } + }() - var wg errgroup.Group - - processQueue := worker.NewQueue[oci.Reference](config.Processing.QueueBurst, config.Processing.QueueRate) - prometheus.DefaultRegisterer.MustRegister(processQueue) - - wg.Go(func() error { - ticker := time.NewTicker(config.Processing.Interval) - defer ticker.Stop() - defer processQueue.Close() - - for { - select { - case <-ctx.Done(): - return ctx.Err() - case <-ticker.C: - ctx, cancel := context.WithTimeout(ctx, 30*time.Second) - - slog.DebugContext(ctx, "Identifying old references to process") - images, err := readStore.ListRawImages(ctx, &store.ListRawImagesOptions{ - NotUpdatedSince: time.Now().Add(-config.Processing.MinAge), - Limit: config.Processing.Items, - }) - if err != nil { - slog.ErrorContext(ctx, "Failed to process old references", slog.Any("error", err)) - cancel() - continue - } - - for _, image := range images { - reference, err := oci.ParseReference(image.Reference) - if err != nil { - slog.ErrorContext(ctx, "Unexpectedly failed to parse reference from store", slog.Any("error", err), slog.String("reference", image.Reference)) - cancel() - return err - } - - processQueue.Push(reference) - } + workerQueue := worker.NewQueue[oci.Reference](config.Processing.QueueBurst, config.Processing.QueueRate) + if err := prometheus.DefaultRegisterer.Register(workerQueue); err != nil { + slog.Error("Failed to register prometheus metrics for worker queue", slog.Any("error", err)) + return ErrNonZeroExitCode + } - cancel() - } - } - }) + imageWorkerScheduler := worker.NewImageWorkerScheduler(readStore) httpClient := httputil.NewClient(cache, config.Cache.MaxAge) httpClient.UserAgent = config.HTTP.UserAgent - prometheus.DefaultRegisterer.MustRegister(httpClient) - - worker := worker.New(httpClient, writeStore, registryAuth) - prometheus.DefaultRegisterer.MustRegister(worker) - - wg.Go(func() error { - for reference := range processQueue.Pull() { - ctx, cancel := context.WithTimeout(ctx, config.Processing.Timeout) - err := worker.ProcessRawImage(ctx, reference) - cancel() - if err != nil { - slog.ErrorContext(ctx, "Failed to process queued raw image", slog.Any("error", err), slog.String("reference", reference.String())) - } - } - - return nil - }) - - wg.Go(func() error { - slog.InfoContext(ctx, "Starting platform grapher") - - grapher, ok := targetPlatform.(platform.ContinuousGrapher) - if !ok { - slog.DebugContext(ctx, "Platform lacks native continuous graphing support. Falling back to polling", slog.Duration("interval", config.Processing.Interval)) - grapher = &platform.PollGrapher{ - Grapher: targetPlatform, - Interval: config.Processing.Interval, - } - } - - graphs, err := grapher.GraphContinuously(ctx) - if err != nil { - slog.ErrorContext(ctx, "Failed to start graphing platform", slog.Any("error", err)) - return err - } - - for graph := range graphs { - slog.DebugContext(ctx, "Got updated platform graph") - - // Delete ignored images / trees - graph.DeleteFunc(func(n platform.Node) bool { - return n.Labels().Ignore() - }) - - roots := graph.Roots() - - for _, root := range roots { - imageNode := root.(platform.ImageNode) - - subgraph := graph.Subgraph(root.ID()) - - edges := subgraph.Edges() - nodes := subgraph.Nodes() - - var namespaceNode *platform.Node - - mappedNodes := make(map[string]models.GraphNode) - for _, node := range nodes { - switch n := node.(type) { - case kubernetes.Resource: - mappedNodes[node.ID()] = models.GraphNode{ - Domain: "kubernetes", - Type: string(n.Kind()), - Name: n.Name(), - Labels: n.Labels().RemoveUnsupported(), - InternalLabels: n.InternalLabels(), - } - if node.Type() == "kubernetes/"+kubernetes.ResourceKindCoreV1Namespace { - namespaceNode = &node - } - case docker.Resource: - mappedNodes[node.ID()] = models.GraphNode{ - Domain: "docker", - Type: string(n.Kind()), - Name: n.Name(), - Labels: n.Labels().RemoveUnsupported(), - InternalLabels: n.InternalLabels(), - } - if node.Type() == "docker/"+docker.ResourceKindSwarmNamespace || node.Type() == "docker/"+docker.ResourceKindComposeProject { - namespaceNode = &node - } - case platform.ImageNode: - // This node is added later on - default: - panic(fmt.Sprintf("mapping unimplemented node type: %s", node.Type())) - } - } - - // Resolve labels for the image node. The nearest label takes precedence - resolvedLabels := make(map[string]string) - queue := []string{root.ID()} - for len(queue) > 0 { - id := queue[0] - queue = queue[1:] - - for k, v := range mappedNodes[id].Labels { - _, ok := resolvedLabels[k] - if !ok { - resolvedLabels[k] = v - } - } - - for adjacent, isChild := range edges[id] { - if isChild { - queue = append(queue, adjacent) - } - } - } - mappedNodes[root.ID()] = models.GraphNode{ - Domain: "oci", - Type: "image", - Name: imageNode.Reference.String(), - Labels: resolvedLabels, - InternalLabels: nil, - } - - tags := []string{} - - // Set tags for resources - if namespaceNode != nil { - children := edges[(*namespaceNode).ID()] - for childID, isParent := range children { - if isParent { - continue - } - - var childNode *platform.Node - for _, node := range nodes { - var n = node - if node.ID() == childID { - childNode = &n - break - } - } - - if childNode != nil { - switch resource := (*childNode).(type) { - case kubernetes.Resource: - kind := resource.Kind() - if kind.IsSupported() { - tags = append(tags, kubernetes.TagName(resource.Kind())) - } - case docker.Resource: - tags = append(tags, docker.TagName(resource.Kind())) - } - } - } - } - - mappedGraph := models.Graph{ - Edges: edges, - Nodes: mappedNodes, - } - - rawImage := &models.RawImage{ - Reference: imageNode.Reference.String(), - Tags: tags, - Graph: mappedGraph, - } - - // TODO: Do this inside of the worker as well? - slog.DebugContext(ctx, "Inserting raw image", slog.String("reference", rawImage.Reference)) - inserted, err := writeStore.InsertRawImage(context.TODO(), rawImage) - if err != nil { - slog.ErrorContext(ctx, "Failed to insert raw image", slog.Any("error", err)) - return err - } - - // Try to schedule the image for processing - if inserted { - slog.DebugContext(ctx, "Raw image inserted for first time - scheduling for processing") - processQueue.Push(imageNode.Reference) - } - } - - allReferences := make([]string, 0) - for _, root := range roots { - imageNode := root.(platform.ImageNode) - allReferences = append(allReferences, imageNode.Reference.String()) - } - - slog.DebugContext(ctx, "Cleaning up removed images") - removed, err := writeStore.DeleteNonPresent(context.TODO(), allReferences) - if err == nil { - slog.DebugContext(ctx, "Cleaned up removed images successfully", slog.Int64("removed", removed)) - } else { - slog.ErrorContext(ctx, "Failed to clean up removed images", slog.Any("error", err)) - } - } - - return nil - }) - - wg.Go(func() error { - ticker := time.NewTicker(config.Workflow.CleanupInterval) - defer ticker.Stop() + if err := prometheus.DefaultRegisterer.Register(httpClient); err != nil { + slog.Error("Failed to register prometheus metrics for HTTP client", slog.Any("error", err)) + return ErrNonZeroExitCode + } - for { - select { - case <-ctx.Done(): - return ctx.Err() - case <-ticker.C: - ctx, cancel := context.WithTimeout(ctx, 30*time.Second) - slog.DebugContext(ctx, "Cleaning up old workflow runs") - removed, err := writeStore.DeleteWorkflowRuns(context.TODO(), time.Now().Add(-config.Workflow.CleanupMaxAge)) - cancel() - if err == nil { - slog.DebugContext(ctx, "Cleaned up old workflow runs successfully", slog.Int64("removed", removed)) - } else { - slog.ErrorContext(ctx, "Failed to clean up old workflow runs", slog.Any("error", err)) - } - } - } - }) + imageWorker := worker.NewImageWorker(httpClient, writeStore, registryAuth) + if err := prometheus.DefaultRegisterer.Register(imageWorker); err != nil { + slog.Error("Failed to register prometheus metrics for worker", slog.Any("error", err)) + return ErrNonZeroExitCode + } - mux := http.NewServeMux() + serveMux := http.NewServeMux() - apiServer := api.NewServer(readStore, worker.Hub, processQueue) + apiServer := api.NewServer(readStore, imageWorker.Hub, workerQueue) apiServer.WebAddress = config.Web.Address - mux.Handle("/api/v1/", apiServer) - mux.Handle("/metrics", promhttp.Handler()) - mux.Handle("/livez", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + serveMux.Handle("/api/v1/", apiServer) + serveMux.Handle("/metrics", promhttp.Handler()) + serveMux.Handle("/livez", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { w.WriteHeader(http.StatusOK) })) - mux.Handle("/readyz", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + serveMux.Handle("/readyz", http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { // TODO: Figure out what checks to have w.WriteHeader(http.StatusOK) })) if !config.Web.Disabled { - mux.Handle("/", web.MustNewEmbeddedServer()) + serveMux.Handle("/", web.MustNewEmbeddedServer()) } httpServer := &http.Server{ @@ -558,56 +217,79 @@ func main() { writer = gzip } - mux.ServeHTTP(writer, r) + serveMux.ServeHTTP(writer, r) }), } - - wg.Go(func() error { - slog.InfoContext(ctx, "Starting HTTP server") - err := httpServer.ListenAndServe() - if err != nil && err != http.ErrServerClosed { - slog.ErrorContext(ctx, "Failed to serve", slog.Any("error", err)) - return err + registerShutdownFunc(func(ctx context.Context) { + if err := httpServer.Shutdown(ctx); err != nil { + slog.Error("Failed to shutdown HTTP server", slog.Any("error", err)) + // Fallthrough } - return nil }) - signals := make(chan os.Signal, 1) - signal.Notify(signals, os.Interrupt) - caught := 0 + var wg sync.WaitGroup + + // Run the worker scheduler to push jobs to the worker queue + wg.Add(1) go func() { - for range signals { - caught++ - if caught == 1 { - slog.InfoContext(ctx, "Caught signal, exiting gracefully") - if err := httpServer.Close(); err != nil { - slog.ErrorContext(ctx, "Failed to close server", slog.Any("error", err)) - // Fallthrough - } - if err := cache.Close(); err != nil { - slog.ErrorContext(ctx, "Failed to close cache", slog.Any("error", err)) - // Fallthrough - } - if err := readStore.Close(); err != nil { - slog.ErrorContext(ctx, "Failed to close read store", slog.Any("error", err)) - // Fallthrough - } - if err := writeStore.Close(); err != nil { - slog.ErrorContext(ctx, "Failed to close write store", slog.Any("error", err)) - // Fallthrough - } - // Cancel goroutines started in main last as to block on all of the - // above calls + defer wg.Done() + defer workerQueue.Close() + + imageWorkerScheduler.PushTo(ctx, workerQueue, config.Processing.Interval, config.Processing.MinAge, config.Processing.Items) + }() + + // Run the worker to handle pull jobs from the worker queue + wg.Add(1) + go func() { + defer wg.Done() + + imageWorker.PullFrom(ctx, workerQueue, config.Processing.Timeout) + }() + + // Keep available images up-to-date by reacting on changes made to the + // platform + slog.DebugContext(ctx, "Starting platform grapher") + wg.Add(1) + go func() { + defer wg.Done() + + }() + + // Start cleaning up on an interval + wg.Add(1) + go func() { + ticker := time.NewTicker(config.Workflow.CleanupInterval) + defer ticker.Stop() + + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + ctx, cancel := context.WithTimeout(ctx, 30*time.Second) + slog.Debug("Cleaning up old workflow runs") + removed, err := writeStore.DeleteWorkflowRuns(ctx, time.Now().Add(-config.Workflow.CleanupMaxAge)) cancel() - } else { - slog.InfoContext(ctx, "Caught signal, exiting now") - os.Exit(1) + if err == nil { + slog.Debug("Cleaned up old workflow runs successfully", slog.Int64("removed", removed)) + } else { + slog.Error("Failed to clean up old workflow runs", slog.Any("error", err)) + } } } }() - if err := wg.Wait(); err != nil && err != ctx.Err() { - slog.ErrorContext(ctx, "Failed to run", slog.Any("error", err)) - os.Exit(1) - } + // Start HTTP server + wg.Add(1) + go func() { + defer wg.Done() + + if err := httpServer.ListenAndServe(); err != nil && err != http.ErrServerClosed { + slog.Error("Failed to serve", slog.Any("error", err)) + // Fallthrough + } + }() + + wg.Wait() + return ctx.Err() } diff --git a/cmd/cupdate/platform.go b/cmd/cupdate/platform.go new file mode 100644 index 00000000..1e1ce62f --- /dev/null +++ b/cmd/cupdate/platform.go @@ -0,0 +1,71 @@ +package main + +import ( + "context" + "fmt" + "log/slog" + "net/url" + "path/filepath" + + "github.com/AlexGustafsson/cupdate/internal/configutils" + "github.com/AlexGustafsson/cupdate/internal/platform" + "github.com/AlexGustafsson/cupdate/internal/platform/docker" + "github.com/AlexGustafsson/cupdate/internal/platform/kubernetes" +) + +func ConfigurePlatformGrapher(ctx context.Context, config *Config) (platform.ContinuousGrapher, error) { + // Set up the configured platform (Docker if specified, auto discovery of + // Kubernetes otherwise) + if len(config.Docker.Hosts) == 0 { + kubernetesConfig, err := config.KubernetesClientConfig() + if err != nil { + return nil, err + } + + platform, err := kubernetes.NewPlatform(kubernetesConfig, &kubernetes.Options{IncludeOldReplicaSets: config.Kubernetes.IncludeOldReplicaSets}) + if err != nil { + return nil, err + } + + return platform, nil + } + + graphers := make([]platform.Grapher, 0) + for _, host := range config.Docker.Hosts { + options := &docker.Options{ + IncludeAllContainers: config.Docker.IncludeAllContainers, + } + + if config.Docker.TLSPath != "" { + uri, err := url.Parse(host) + if err != nil { + return nil, fmt.Errorf("failed to parse docker host '%s': %w", host, err) + } + + tlsConfig, err := configutils.LoadTLSConfig( + filepath.Join(config.Docker.TLSPath, uri.Hostname()), + config.Docker.TLSPath, + ) + if err != nil { + return nil, fmt.Errorf("failed to read docker host '%s' TLS files: %w", host, err) + } + + options.TLSClientConfig = tlsConfig + } + + platform, err := docker.NewPlatform(ctx, host, options) + if err != nil { + return nil, err + } + + graphers = append(graphers, platform) + } + + slog.Debug("Platform lacks native continuous graphing support. Falling back to polling", slog.Duration("interval", config.Processing.Interval)) + return &platform.PollGrapher{ + Grapher: &platform.CompoundGrapher{ + Graphers: graphers, + }, + Interval: config.Processing.Interval, + }, nil +} diff --git a/cmd/cupdate/version.go b/cmd/cupdate/version.go new file mode 100644 index 00000000..9d63c923 --- /dev/null +++ b/cmd/cupdate/version.go @@ -0,0 +1,4 @@ +package main + +// Version is the current Cupdate version. Overwritten at build time. +var Version string = "development build" diff --git a/internal/api/server.go b/internal/api/server.go index 9cb816a3..1ac5ffc4 100644 --- a/internal/api/server.go +++ b/internal/api/server.go @@ -27,13 +27,13 @@ var ( type Server struct { api *store.Store - hub *events.Hub[worker.Event] + hub *events.Hub[worker.ImageEvent] mux *http.ServeMux WebAddress string } -func NewServer(api *store.Store, hub *events.Hub[worker.Event], processQueue *worker.Queue[oci.Reference]) *Server { +func NewServer(api *store.Store, hub *events.Hub[worker.ImageEvent], processQueue *worker.Queue[oci.Reference]) *Server { s := &Server{ api: api, hub: hub, @@ -338,11 +338,11 @@ func NewServer(api *store.Store, hub *events.Hub[worker.Event], processQueue *wo for event := range s.hub.Subscribe(ctx) { var eventType models.EventType switch event.Type { - case worker.EventTypeUpdated: + case worker.ImageEventTypeUpdated: eventType = models.EventTypeImageUpdated - case worker.EventTypeProcessed: + case worker.ImageEventTypeProcessed: eventType = models.EventTypeImageProcessed - case worker.EventTypeNewVersionAvailable: + case worker.ImageEventTypeNewVersionAvailable: eventType = models.EventTypeImageNewVersionAvailable } diff --git a/internal/httputil/auth.go b/internal/httputil/auth.go index 972bb17a..0a53bbb1 100644 --- a/internal/httputil/auth.go +++ b/internal/httputil/auth.go @@ -36,15 +36,6 @@ func (h BasicAuthHandler) HandleAuth(r *http.Request) error { return nil } -// BearerToken auths requests using the Bearer authorization scheme. -type BearerToken string - -func (t BearerToken) HandleAuth(r *http.Request) error { - r.Header.Set("Authorization", "Bearer "+string(t)) - - return nil -} - // AuthMux is an HTTP auth multiplexer. It matches URLs of auth requests against // a list of registered patterns and calls the handler for the pattern that most // closely matches the request. diff --git a/internal/worker/graphworker.go b/internal/worker/graphworker.go new file mode 100644 index 00000000..aaf2b2fb --- /dev/null +++ b/internal/worker/graphworker.go @@ -0,0 +1,192 @@ +package worker + +import ( + "context" + "fmt" + "log/slog" + + "github.com/AlexGustafsson/cupdate/internal/models" + "github.com/AlexGustafsson/cupdate/internal/oci" + "github.com/AlexGustafsson/cupdate/internal/platform" + "github.com/AlexGustafsson/cupdate/internal/platform/docker" + "github.com/AlexGustafsson/cupdate/internal/platform/kubernetes" + "github.com/AlexGustafsson/cupdate/internal/store" +) + +type GraphWorker struct { + grapher platform.ContinuousGrapher + store *store.Store + queue *Queue[oci.Reference] +} + +func NewGraphWorker(grapher platform.ContinuousGrapher, store *store.Store, queue *Queue[oci.Reference]) *GraphWorker { + return &GraphWorker{ + grapher: grapher, + store: store, + queue: queue, + } +} + +func (w *GraphWorker) GraphContinuously(ctx context.Context) error { + graphs, err := w.grapher.GraphContinuously(ctx) + if err != nil { + slog.Error("Failed to start graphing platform", slog.Any("error", err)) + return err + } + + for graph := range graphs { + slog.Debug("Got updated platform graph") + + // Delete ignored images / trees + graph.DeleteFunc(func(n platform.Node) bool { + return n.Labels().Ignore() + }) + + roots := graph.Roots() + + for _, root := range roots { + imageNode := root.(platform.ImageNode) + + subgraph := graph.Subgraph(root.ID()) + + edges := subgraph.Edges() + nodes := subgraph.Nodes() + + var namespaceNode *platform.Node + + mappedNodes := make(map[string]models.GraphNode) + for _, node := range nodes { + switch n := node.(type) { + case kubernetes.Resource: + mappedNodes[node.ID()] = models.GraphNode{ + Domain: "kubernetes", + Type: string(n.Kind()), + Name: n.Name(), + Labels: n.Labels().RemoveUnsupported(), + InternalLabels: n.InternalLabels(), + } + if node.Type() == "kubernetes/"+kubernetes.ResourceKindCoreV1Namespace { + namespaceNode = &node + } + case docker.Resource: + mappedNodes[node.ID()] = models.GraphNode{ + Domain: "docker", + Type: string(n.Kind()), + Name: n.Name(), + Labels: n.Labels().RemoveUnsupported(), + InternalLabels: n.InternalLabels(), + } + if node.Type() == "docker/"+docker.ResourceKindSwarmNamespace || node.Type() == "docker/"+docker.ResourceKindComposeProject { + namespaceNode = &node + } + case platform.ImageNode: + // This node is added later on + default: + panic(fmt.Sprintf("mapping unimplemented node type: %s", node.Type())) + } + } + + // Resolve labels for the image node. The nearest label takes precedence + resolvedLabels := make(map[string]string) + queue := []string{root.ID()} + for len(queue) > 0 { + id := queue[0] + queue = queue[1:] + + for k, v := range mappedNodes[id].Labels { + _, ok := resolvedLabels[k] + if !ok { + resolvedLabels[k] = v + } + } + + for adjacent, isChild := range edges[id] { + if isChild { + queue = append(queue, adjacent) + } + } + } + mappedNodes[root.ID()] = models.GraphNode{ + Domain: "oci", + Type: "image", + Name: imageNode.Reference.String(), + Labels: resolvedLabels, + InternalLabels: nil, + } + + tags := []string{} + + // Set tags for resources + if namespaceNode != nil { + children := edges[(*namespaceNode).ID()] + for childID, isParent := range children { + if isParent { + continue + } + + var childNode *platform.Node + for _, node := range nodes { + var n = node + if node.ID() == childID { + childNode = &n + break + } + } + + if childNode != nil { + switch resource := (*childNode).(type) { + case kubernetes.Resource: + kind := resource.Kind() + if kind.IsSupported() { + tags = append(tags, kubernetes.TagName(resource.Kind())) + } + case docker.Resource: + tags = append(tags, docker.TagName(resource.Kind())) + } + } + } + } + + mappedGraph := models.Graph{ + Edges: edges, + Nodes: mappedNodes, + } + + rawImage := &models.RawImage{ + Reference: imageNode.Reference.String(), + Tags: tags, + Graph: mappedGraph, + } + + // TODO: Do this inside of the worker as well? + slog.DebugContext(ctx, "Inserting raw image", slog.String("reference", rawImage.Reference)) + inserted, err := w.store.InsertRawImage(context.TODO(), rawImage) + if err != nil { + slog.ErrorContext(ctx, "Failed to insert raw image", slog.Any("error", err)) + continue + } + + // Try to schedule the image for processing + if inserted { + slog.DebugContext(ctx, "Raw image inserted for first time - scheduling for processing") + w.queue.Push(imageNode.Reference) + } + } + + allReferences := make([]string, 0) + for _, root := range roots { + imageNode := root.(platform.ImageNode) + allReferences = append(allReferences, imageNode.Reference.String()) + } + + slog.DebugContext(ctx, "Cleaning up removed images") + removed, err := w.store.DeleteNonPresent(context.TODO(), allReferences) + if err == nil { + slog.DebugContext(ctx, "Cleaned up removed images successfully", slog.Int64("removed", removed)) + } else { + slog.ErrorContext(ctx, "Failed to clean up removed images", slog.Any("error", err)) + } + } + + return nil +} diff --git a/internal/worker/worker.go b/internal/worker/imageworker.go similarity index 84% rename from internal/worker/worker.go rename to internal/worker/imageworker.go index 1b3bca5e..c89a5e82 100644 --- a/internal/worker/worker.go +++ b/internal/worker/imageworker.go @@ -17,32 +17,32 @@ import ( "github.com/prometheus/client_golang/prometheus" ) -var _ prometheus.Collector = (*Worker)(nil) +var _ prometheus.Collector = (*ImageWorker)(nil) -// EventType is the type of an event. -type EventType string +// ImageEventType is the type of an event. +type ImageEventType string const ( - // EventTypeUpdated is emitted whenever data of an image is updated. - EventTypeUpdated EventType = "updated" - // EventTypeProcessed is emitted whenever an image was processed. - EventTypeProcessed EventType = "processed" - // EventTypeNewVersionAvailable is emitted whenever the latest available + // ImageEventTypeUpdated is emitted whenever data of an image is updated. + ImageEventTypeUpdated ImageEventType = "updated" + // ImageEventTypeProcessed is emitted whenever an image was processed. + ImageEventTypeProcessed ImageEventType = "processed" + // ImageEventTypeNewVersionAvailable is emitted whenever the latest available // version of an image changes. - EventTypeNewVersionAvailable EventType = "newVersionAvailable" + ImageEventTypeNewVersionAvailable ImageEventType = "newVersionAvailable" ) -// Event describes a Worker event. -type Event struct { +// ImageEvent describes an [ImageWorker] event. +type ImageEvent struct { Reference string - Type EventType + Type ImageEventType } -// Worker processes raw container image entries, running the image workflow and -// storing the result to the state store. -// The worker produces events of the type [Event]. -type Worker struct { - *events.Hub[Event] +// ImageWorker processes raw container image entries, running the image workflow +// and storing the result to the state store. +// The worker produces events of the type [Image]. +type ImageWorker struct { + *events.Hub[ImageEvent] httpClient httputil.Requester store *store.Store @@ -53,9 +53,9 @@ type Worker struct { processingGauge prometheus.Gauge } -func New(httpClient httputil.Requester, store *store.Store, registryAuth *httputil.AuthMux) *Worker { - return &Worker{ - Hub: events.NewHub[Event](), +func NewImageWorker(httpClient httputil.Requester, store *store.Store, registryAuth *httputil.AuthMux) *ImageWorker { + return &ImageWorker{ + Hub: events.NewHub[ImageEvent](), httpClient: httpClient, store: store, @@ -79,8 +79,21 @@ func New(httpClient httputil.Requester, store *store.Store, registryAuth *httput } } +// PullFrom continuously processes raw images by pulling them from a [Queue]. +func (w *ImageWorker) PullFrom(ctx context.Context, queue *Queue[oci.Reference], timeout time.Duration) { + for reference := range queue.Pull() { + ctx, cancel := context.WithTimeout(ctx, timeout) + err := w.ProcessRawImage(ctx, reference) + cancel() + if err != nil { + slog.ErrorContext(ctx, "Failed to process queued raw image", slog.Any("error", err), slog.String("reference", reference.String())) + // Fallthrough + } + } +} + // ProcessRawImage processes a raw image by the specified reference. -func (w *Worker) ProcessRawImage(ctx context.Context, reference oci.Reference) error { +func (w *ImageWorker) ProcessRawImage(ctx context.Context, reference oci.Reference) error { start := time.Now() w.processingGauge.Inc() defer w.processingGauge.Dec() @@ -309,9 +322,9 @@ func (w *Worker) ProcessRawImage(ctx context.Context, reference oci.Reference) e log.DebugContext(ctx, "Updated image date", slog.Int("changes", len(changes))) // TODO: Group changes, create an event specifying the time. That way the // browser can ignore the event if it already updated after the time? - w.Broadcast(ctx, Event{ + w.Broadcast(ctx, ImageEvent{ Reference: reference.String(), - Type: EventTypeUpdated, + Type: ImageEventTypeUpdated, }) // TODO: Have another readonly job for going over the changes made to @@ -322,28 +335,28 @@ func (w *Worker) ProcessRawImage(ctx context.Context, reference oci.Reference) e } if result.LatestReference != "" && result.LatestReference != result.Reference { - w.Broadcast(ctx, Event{ + w.Broadcast(ctx, ImageEvent{ Reference: reference.String(), - Type: EventTypeNewVersionAvailable, + Type: ImageEventTypeNewVersionAvailable, }) } - w.Broadcast(ctx, Event{ + w.Broadcast(ctx, ImageEvent{ Reference: reference.String(), - Type: EventTypeProcessed, + Type: ImageEventTypeProcessed, }) return nil } // Collect implements prometheus.Collector. -func (w *Worker) Collect(ch chan<- prometheus.Metric) { +func (w *ImageWorker) Collect(ch chan<- prometheus.Metric) { w.processedCounter.Collect(ch) w.processingDuration.Collect(ch) w.processingGauge.Collect(ch) } // Describe implements prometheus.Collector. -func (w *Worker) Describe(descs chan<- *prometheus.Desc) { +func (w *ImageWorker) Describe(descs chan<- *prometheus.Desc) { prometheus.DescribeByCollect(w, descs) } diff --git a/internal/worker/imageworkerscheduler.go b/internal/worker/imageworkerscheduler.go new file mode 100644 index 00000000..a98e8a97 --- /dev/null +++ b/internal/worker/imageworkerscheduler.go @@ -0,0 +1,59 @@ +package worker + +import ( + "context" + "log/slog" + "time" + + "github.com/AlexGustafsson/cupdate/internal/oci" + "github.com/AlexGustafsson/cupdate/internal/store" +) + +type ImageWorkerScheduler struct { + store *store.Store +} + +func NewImageWorkerScheduler(store *store.Store) *ImageWorkerScheduler { + return &ImageWorkerScheduler{ + store: store, + } +} + +// PushTo continuously pushes old references to the queue at a set interval. +func (s *ImageWorkerScheduler) PushTo(ctx context.Context, queue *Queue[oci.Reference], interval time.Duration, minAge time.Duration, batchSize int) { + ticker := time.NewTicker(interval) + defer ticker.Stop() + + for { + select { + case <-ctx.Done(): + return + case <-ticker.C: + ctx, cancel := context.WithTimeout(context.Background(), 30*time.Second) + s.queueBatch(ctx, queue, time.Now().Add(-minAge), batchSize) + cancel() + } + } +} + +func (s *ImageWorkerScheduler) queueBatch(ctx context.Context, queue *Queue[oci.Reference], notUpdatedSince time.Time, batchSize int) { + slog.DebugContext(ctx, "Identifying old references to process") + images, err := s.store.ListRawImages(ctx, &store.ListRawImagesOptions{ + NotUpdatedSince: notUpdatedSince, + Limit: batchSize, + }) + if err != nil { + slog.ErrorContext(ctx, "Failed to process old references", slog.Any("error", err)) + return + } + + for _, image := range images { + reference, err := oci.ParseReference(image.Reference) + if err != nil { + slog.ErrorContext(ctx, "Unexpectedly failed to parse reference from store", slog.Any("error", err), slog.String("reference", image.Reference)) + return + } + + queue.Push(reference) + } +}