|
| 1 | +// Copyright Amazon.com, Inc. or its affiliates. All Rights Reserved. |
| 2 | +// SPDX-License-Identifier: MIT |
| 3 | + |
| 4 | +package k8smetadata |
| 5 | + |
| 6 | +import ( |
| 7 | + "context" |
| 8 | + "math/rand" |
| 9 | + "time" |
| 10 | + |
| 11 | + "go.opentelemetry.io/collector/component" |
| 12 | + "go.opentelemetry.io/collector/extension" |
| 13 | + "go.uber.org/atomic" |
| 14 | + "go.uber.org/zap" |
| 15 | + "k8s.io/client-go/informers" |
| 16 | + "k8s.io/client-go/kubernetes" |
| 17 | + "k8s.io/client-go/tools/clientcmd" |
| 18 | + |
| 19 | + "github.com/aws/amazon-cloudwatch-agent/internal/k8sCommon/k8sclient" |
| 20 | +) |
| 21 | + |
| 22 | +const ( |
| 23 | + deletionDelay = 2 * time.Minute |
| 24 | + jitterKubernetesAPISeconds = 10 |
| 25 | +) |
| 26 | + |
| 27 | +type KubernetesMetadata struct { |
| 28 | + logger *zap.Logger |
| 29 | + config *Config |
| 30 | + ready atomic.Bool |
| 31 | + safeStopCh *k8sclient.SafeChannel |
| 32 | + endpointSliceWatcher *k8sclient.EndpointSliceWatcher |
| 33 | +} |
| 34 | + |
| 35 | +var _ extension.Extension = (*KubernetesMetadata)(nil) |
| 36 | + |
| 37 | +func jitterSleep(seconds int) { |
| 38 | + jitter := time.Duration(rand.Intn(seconds)) * time.Second // nolint:gosec |
| 39 | + time.Sleep(jitter) |
| 40 | +} |
| 41 | + |
| 42 | +func (e *KubernetesMetadata) Start(_ context.Context, _ component.Host) error { |
| 43 | + e.logger.Debug("Starting k8smetadata extension...") |
| 44 | + |
| 45 | + config, err := clientcmd.BuildConfigFromFlags("", "") |
| 46 | + if err != nil { |
| 47 | + e.logger.Error("Failed to create config", zap.Error(err)) |
| 48 | + } |
| 49 | + |
| 50 | + clientset, err := kubernetes.NewForConfig(config) |
| 51 | + if err != nil { |
| 52 | + e.logger.Error("Failed to create kubernetes client", zap.Error(err)) |
| 53 | + } |
| 54 | + |
| 55 | + // jitter calls to the kubernetes api (a precaution to prevent overloading api server) |
| 56 | + jitterSleep(jitterKubernetesAPISeconds) |
| 57 | + |
| 58 | + timedDeleter := &k8sclient.TimedDeleter{Delay: deletionDelay} |
| 59 | + sharedInformerFactory := informers.NewSharedInformerFactory(clientset, 0) |
| 60 | + |
| 61 | + e.endpointSliceWatcher = k8sclient.NewEndpointSliceWatcher(e.logger, sharedInformerFactory, timedDeleter) |
| 62 | + e.safeStopCh = &k8sclient.SafeChannel{Ch: make(chan struct{}), Closed: false} |
| 63 | + |
| 64 | + e.endpointSliceWatcher.Run(e.safeStopCh.Ch) |
| 65 | + |
| 66 | + e.endpointSliceWatcher.WaitForCacheSync(e.safeStopCh.Ch) |
| 67 | + |
| 68 | + e.logger.Debug("EndpointSlice cache synced, extension fully started") |
| 69 | + e.ready.Store(true) |
| 70 | + |
| 71 | + return nil |
| 72 | +} |
| 73 | + |
| 74 | +func (e *KubernetesMetadata) Shutdown(_ context.Context) error { |
| 75 | + if e.safeStopCh != nil { |
| 76 | + e.safeStopCh.Close() |
| 77 | + } |
| 78 | + return nil |
| 79 | +} |
| 80 | + |
| 81 | +func (e *KubernetesMetadata) GetPodMetadata(ip string) k8sclient.PodMetadata { |
| 82 | + if ip == "" { |
| 83 | + e.logger.Debug("GetPodMetadata: no IP provided") |
| 84 | + return k8sclient.PodMetadata{} |
| 85 | + } |
| 86 | + pm, ok := e.endpointSliceWatcher.IPToPodMetadata.Load(ip) |
| 87 | + if !ok { |
| 88 | + e.logger.Debug("GetPodMetadata: no mapping found for IP", zap.String("ip", ip)) |
| 89 | + return k8sclient.PodMetadata{} |
| 90 | + } |
| 91 | + metadata := pm.(k8sclient.PodMetadata) |
| 92 | + e.logger.Debug("GetPodMetadata: found metadata", |
| 93 | + zap.String("ip", ip), |
| 94 | + zap.String("workload", metadata.Workload), |
| 95 | + zap.String("namespace", metadata.Namespace), |
| 96 | + zap.String("node", metadata.Node), |
| 97 | + ) |
| 98 | + return metadata |
| 99 | +} |
0 commit comments