diff --git a/pkg/operator/apis/monitoring/v1alpha1/types.go b/pkg/operator/apis/monitoring/v1alpha1/types.go index fc4122807f3f..a9f0b6396694 100644 --- a/pkg/operator/apis/monitoring/v1alpha1/types.go +++ b/pkg/operator/apis/monitoring/v1alpha1/types.go @@ -4,6 +4,7 @@ import ( prom_v1 "github.com/prometheus-operator/prometheus-operator/pkg/apis/monitoring/v1" v1 "k8s.io/api/core/v1" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "sigs.k8s.io/controller-runtime/pkg/client" ) // +kubebuilder:object:root=true @@ -24,6 +25,7 @@ type GrafanaAgent struct { // MetricsInstanceSelector returns a selector to find MetricsInstances. func (a *GrafanaAgent) MetricsInstanceSelector() ObjectSelector { return ObjectSelector{ + ObjectType: &MetricsInstance{}, ParentNamespace: a.Namespace, NamespaceSelector: a.Spec.Metrics.InstanceNamespaceSelector, Labels: a.Spec.Metrics.InstanceSelector, @@ -33,6 +35,7 @@ func (a *GrafanaAgent) MetricsInstanceSelector() ObjectSelector { // LogsInstanceSelector returns a selector to find LogsInstances. func (a *GrafanaAgent) LogsInstanceSelector() ObjectSelector { return ObjectSelector{ + ObjectType: &LogsInstance{}, ParentNamespace: a.Namespace, NamespaceSelector: a.Spec.Logs.InstanceNamespaceSelector, Labels: a.Spec.Logs.InstanceSelector, @@ -153,6 +156,7 @@ type GrafanaAgentSpec struct { // resource hierarchy. When NamespaceSelector is nil, objects should be // searched directly in the ParentNamespace. type ObjectSelector struct { + ObjectType client.Object ParentNamespace string NamespaceSelector *metav1.LabelSelector Labels *metav1.LabelSelector diff --git a/pkg/operator/apis/monitoring/v1alpha1/types_logs.go b/pkg/operator/apis/monitoring/v1alpha1/types_logs.go index 4b5f46b4290d..d32b76c48b20 100644 --- a/pkg/operator/apis/monitoring/v1alpha1/types_logs.go +++ b/pkg/operator/apis/monitoring/v1alpha1/types_logs.go @@ -100,9 +100,10 @@ type LogsInstance struct { Spec LogsInstanceSpec `json:"spec,omitempty"` } -// PodLogsInstanceSelector returns the selector to discover PodLogs. -func (i *LogsInstance) PodLogsInstanceSelector() ObjectSelector { +// PodLogsSelector returns the selector to discover PodLogs. +func (i *LogsInstance) PodLogsSelector() ObjectSelector { return ObjectSelector{ + ObjectType: &PodLogs{}, ParentNamespace: i.Namespace, NamespaceSelector: i.Spec.PodLogsNamespaceSelector, Labels: i.Spec.PodLogsSelector, diff --git a/pkg/operator/apis/monitoring/v1alpha1/types_metrics.go b/pkg/operator/apis/monitoring/v1alpha1/types_metrics.go index 9b62046f5c1b..5a20cc22d3cf 100644 --- a/pkg/operator/apis/monitoring/v1alpha1/types_metrics.go +++ b/pkg/operator/apis/monitoring/v1alpha1/types_metrics.go @@ -182,6 +182,7 @@ type MetricsInstance struct { // ServiceMonitorSelector returns a selector to find ServiceMonitors. func (p *MetricsInstance) ServiceMonitorSelector() ObjectSelector { return ObjectSelector{ + ObjectType: &prom_v1.ServiceMonitor{}, ParentNamespace: p.Namespace, NamespaceSelector: p.Spec.ServiceMonitorNamespaceSelector, Labels: p.Spec.ServiceMonitorSelector, @@ -191,6 +192,7 @@ func (p *MetricsInstance) ServiceMonitorSelector() ObjectSelector { // PodMonitorSelector returns a selector to find PodMonitors. func (p *MetricsInstance) PodMonitorSelector() ObjectSelector { return ObjectSelector{ + ObjectType: &prom_v1.PodMonitor{}, ParentNamespace: p.Namespace, NamespaceSelector: p.Spec.PodMonitorNamespaceSelector, Labels: p.Spec.PodMonitorSelector, @@ -200,6 +202,7 @@ func (p *MetricsInstance) PodMonitorSelector() ObjectSelector { // ProbeSelector returns a selector to find Probes. func (p *MetricsInstance) ProbeSelector() ObjectSelector { return ObjectSelector{ + ObjectType: &prom_v1.Probe{}, ParentNamespace: p.Namespace, NamespaceSelector: p.Spec.ProbeNamespaceSelector, Labels: p.Spec.ProbeSelector, diff --git a/pkg/operator/build_hierarchy.go b/pkg/operator/build_hierarchy.go new file mode 100644 index 000000000000..cab0b212644b --- /dev/null +++ b/pkg/operator/build_hierarchy.go @@ -0,0 +1,282 @@ +package operator + +import ( + "context" + "fmt" + + "github.com/go-kit/log" + "github.com/go-kit/log/level" + grafana "github.com/grafana/agent/pkg/operator/apis/monitoring/v1alpha1" + "github.com/grafana/agent/pkg/operator/assets" + "github.com/grafana/agent/pkg/operator/config" + "github.com/grafana/agent/pkg/operator/hierarchy" + prom "github.com/prometheus-operator/prometheus-operator/pkg/apis/monitoring/v1" + corev1 "k8s.io/api/core/v1" + v1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/api/meta" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/runtime" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/apiutil" +) + +// buildHierarchy constructs a resource hierarchy starting from root. +func buildHierarchy(ctx context.Context, l log.Logger, cli client.Client, root *grafana.GrafanaAgent) (deployment config.Deployment, watchers []hierarchy.Watcher, err error) { + deployment.Agent = root + + // search is used throughout BuildHierarchy, where it will perform a list for + // a set of objects in the hierarchy and populate the watchers return + // variable. + search := func(resources []hierarchyResource) error { + for _, res := range resources { + sel, err := res.Find(ctx, cli) + if err != nil { + gvk, _ := apiutil.GVKForObject(res.List, cli.Scheme()) + return fmt.Errorf("failed to find %q resource: %w", gvk.String(), err) + } + + watchers = append(watchers, hierarchy.Watcher{ + Object: res.Selector.ObjectType, + Owner: client.ObjectKeyFromObject(root), + Selector: sel, + }) + } + return nil + } + + // Root resources + var ( + metricInstances grafana.MetricsInstanceList + logsInstances grafana.LogsInstanceList + ) + var roots = []hierarchyResource{ + {List: &metricInstances, Selector: root.MetricsInstanceSelector()}, + {List: &logsInstances, Selector: root.LogsInstanceSelector()}, + } + if err := search(roots); err != nil { + return deployment, nil, err + } + + // Metrics resources + for _, metricsInst := range metricInstances.Items { + var ( + serviceMonitors prom.ServiceMonitorList + podMonitors prom.PodMonitorList + probes prom.ProbeList + ) + var children = []hierarchyResource{ + {List: &serviceMonitors, Selector: metricsInst.ServiceMonitorSelector()}, + {List: &podMonitors, Selector: metricsInst.PodMonitorSelector()}, + {List: &probes, Selector: metricsInst.ProbeSelector()}, + } + if err := search(children); err != nil { + return deployment, nil, err + } + + deployment.Metrics = append(deployment.Metrics, config.MetricsInstance{ + Instance: metricsInst, + ServiceMonitors: filterServiceMonitors(l, root, &serviceMonitors).Items, + PodMonitors: podMonitors.Items, + Probes: probes.Items, + }) + } + + // Logs resources + for _, logsInst := range logsInstances.Items { + var ( + podLogs grafana.PodLogsList + ) + var children = []hierarchyResource{ + {List: &podLogs, Selector: logsInst.PodLogsSelector()}, + } + if err := search(children); err != nil { + return deployment, nil, err + } + + deployment.Logs = append(deployment.Logs, config.LogInstance{ + Instance: logsInst, + PodLogs: podLogs.Items, + }) + } + + // Finally, find all referenced secrets + secrets, secretWatchers, err := buildSecrets(ctx, cli, deployment) + if err != nil { + return deployment, nil, fmt.Errorf("failed to discover secrets: %w", err) + } + deployment.Secrets = secrets + watchers = append(watchers, secretWatchers...) + + return deployment, watchers, nil +} + +type hierarchyResource struct { + List client.ObjectList // List to populate + Selector grafana.ObjectSelector // Raw selector to use for list +} + +func (hr *hierarchyResource) Find(ctx context.Context, cli client.Client) (hierarchy.Selector, error) { + sel, err := toSelector(hr.Selector) + if err != nil { + return nil, fmt.Errorf("failed to build selector: %w", err) + } + err = hierarchy.List(ctx, cli, hr.List, sel) + if err != nil { + return nil, fmt.Errorf("failed to list resources: %w", err) + } + return sel, nil +} + +func toSelector(os grafana.ObjectSelector) (hierarchy.Selector, error) { + var res hierarchy.LabelsSelector + res.NamespaceName = os.ParentNamespace + + if os.NamespaceSelector != nil { + sel, err := metav1.LabelSelectorAsSelector(os.NamespaceSelector) + if err != nil { + return nil, fmt.Errorf("invalid namespace selector: %w", err) + } + res.NamespaceLabels = sel + } + + sel, err := metav1.LabelSelectorAsSelector(os.Labels) + if err != nil { + return nil, fmt.Errorf("invalid label selector: %w", err) + } + res.Labels = sel + return &res, nil +} + +func filterServiceMonitors(l log.Logger, root *grafana.GrafanaAgent, list *prom.ServiceMonitorList) *prom.ServiceMonitorList { + items := make([]*prom.ServiceMonitor, 0, len(list.Items)) + +Item: + for _, item := range list.Items { + if root.Spec.Metrics.ArbitraryFSAccessThroughSMs.Deny { + for _, ep := range item.Spec.Endpoints { + err := testForArbitraryFSAccess(ep) + if err == nil { + continue + } + + level.Warn(l).Log( + "msg", "skipping service monitor", + "agent", client.ObjectKeyFromObject(root), + "servicemonitor", client.ObjectKeyFromObject(item), + "err", err, + ) + continue Item + } + } + items = append(items, item) + } + + return &prom.ServiceMonitorList{ + TypeMeta: list.TypeMeta, + ListMeta: *list.ListMeta.DeepCopy(), + Items: items, + } +} + +func testForArbitraryFSAccess(e prom.Endpoint) error { + if e.BearerTokenFile != "" { + return fmt.Errorf("it accesses file system via bearer token file which is disallowed via GrafanaAgent specification") + } + + if e.TLSConfig == nil { + return nil + } + + if e.TLSConfig.CAFile != "" || e.TLSConfig.CertFile != "" || e.TLSConfig.KeyFile != "" { + return fmt.Errorf("it accesses file system via TLS config which is disallowed via GrafanaAgent specification") + } + + return nil +} + +func buildSecrets(ctx context.Context, cli client.Client, deploy config.Deployment) (secrets assets.SecretStore, watchers []hierarchy.Watcher, err error) { + secrets = make(assets.SecretStore) + + // KeySelector caches to make sure we don't create duplicate watchers. + var ( + usedSecretSelectors = map[hierarchy.KeySelector]struct{}{} + usedConfigMapSelectors = map[hierarchy.KeySelector]struct{}{} + ) + + for _, ref := range deploy.AssetReferences() { + var ( + objectList client.ObjectList + sel hierarchy.KeySelector + ) + + switch { + case ref.Reference.Secret != nil: + objectList = &corev1.SecretList{} + sel = hierarchy.KeySelector{ + Namespace: ref.Namespace, + Name: ref.Reference.Secret.Name, + } + case ref.Reference.ConfigMap != nil: + objectList = &corev1.ConfigMapList{} + sel = hierarchy.KeySelector{ + Namespace: ref.Namespace, + Name: ref.Reference.ConfigMap.Name, + } + } + + gvk, _ := apiutil.GVKForObject(objectList, cli.Scheme()) + if err := hierarchy.List(ctx, cli, objectList, &sel); err != nil { + return nil, nil, fmt.Errorf("failed to find %q resource: %w", gvk.String(), err) + } + + err := meta.EachListItem(objectList, func(o runtime.Object) error { + var value string + + switch o := o.(type) { + case *corev1.Secret: + rawValue, ok := o.Data[ref.Reference.Secret.Key] + if !ok { + return fmt.Errorf("no key %s in Secret %s", ref.Reference.ConfigMap.Key, o.Name) + } + value = string(rawValue) + case *corev1.ConfigMap: + rawValue, ok := o.BinaryData[ref.Reference.ConfigMap.Key] + if !ok { + return fmt.Errorf("no key %s in ConfigMap %s", ref.Reference.ConfigMap.Key, o.Name) + } + value = string(rawValue) + } + + secrets[assets.KeyForSelector(ref.Namespace, &ref.Reference)] = value + return nil + }) + if err != nil { + return nil, nil, fmt.Errorf("failed to iterate over %q list: %w", gvk.String(), err) + } + + switch { + case ref.Reference.Secret != nil: + if _, used := usedSecretSelectors[sel]; used { + continue + } + watchers = append(watchers, hierarchy.Watcher{ + Object: &v1.Secret{}, + Owner: client.ObjectKeyFromObject(deploy.Agent), + Selector: &sel, + }) + usedSecretSelectors[sel] = struct{}{} + case ref.Reference.ConfigMap != nil: + if _, used := usedConfigMapSelectors[sel]; used { + continue + } + watchers = append(watchers, hierarchy.Watcher{ + Object: &v1.ConfigMap{}, + Owner: client.ObjectKeyFromObject(deploy.Agent), + Selector: &sel, + }) + usedConfigMapSelectors[sel] = struct{}{} + } + } + + return secrets, watchers, nil +} diff --git a/pkg/operator/build_hierarchy_test.go b/pkg/operator/build_hierarchy_test.go new file mode 100644 index 000000000000..b77adb761c95 --- /dev/null +++ b/pkg/operator/build_hierarchy_test.go @@ -0,0 +1,166 @@ +//go:build !nonetwork && !nodocker && !race +// +build !nonetwork,!nodocker,!race + +package operator + +import ( + "context" + "fmt" + "sync" + "testing" + "time" + + grafana "github.com/grafana/agent/pkg/operator/apis/monitoring/v1alpha1" + "github.com/grafana/agent/pkg/operator/hierarchy" + "github.com/grafana/agent/pkg/util" + "github.com/grafana/agent/pkg/util/k8s" + "github.com/grafana/agent/pkg/util/structwalk" + prom "github.com/prometheus-operator/prometheus-operator/pkg/apis/monitoring/v1" + "github.com/stretchr/testify/require" + v1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/labels" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/apiutil" +) + +// Test_buildHierarchy checks that an entire resource hierarchy can be +// discovered. +func Test_buildHierarchy(t *testing.T) { + var wg sync.WaitGroup + defer wg.Wait() + + ctx, cancel := context.WithTimeout(context.Background(), 3*time.Minute) + defer cancel() + + l := util.TestLogger(t) + cluster := NewTestCluster(ctx, t, l) + cli := cluster.Client() + + resources := k8s.NewResourceSet(l, cluster) + defer resources.Stop() + require.NoError(t, resources.AddFile(ctx, "./testdata/test-resource-hierarchy.yaml")) + + // Get root resource + var root grafana.GrafanaAgent + err := cli.Get(ctx, client.ObjectKey{Namespace: "default", Name: "grafana-agent-example"}, &root) + require.NoError(t, err) + + deployment, watchers, err := buildHierarchy(ctx, l, cli, &root) + require.NoError(t, err) + + // Check resources in hierarchy + { + expectedResources := []string{ + "GrafanaAgent/grafana-agent-example", + "MetricsInstance/primary", + "LogsInstance/primary", + "PodMonitor/grafana-agents", + "PodLogs/grafana-agents", + } + var gotResources []string + structwalk.Walk(&resourceWalker{ + onResource: func(c client.Object) { + gvk, _ := apiutil.GVKForObject(c, cli.Scheme()) + + key := fmt.Sprintf("%s/%s", gvk.Kind, c.GetName()) + gotResources = append(gotResources, key) + }, + }, deployment) + + require.ElementsMatch(t, expectedResources, gotResources) + } + + // Check secrets + { + expectedSecrets := []string{ + "/secrets/default/prometheus-fake-credentials/fakeUsername", + "/secrets/default/prometheus-fake-credentials/fakePassword", + } + var actualSecrets []string + for key := range deployment.Secrets { + actualSecrets = append(actualSecrets, string(key)) + } + + require.ElementsMatch(t, expectedSecrets, actualSecrets) + } + + // Check configured watchers + { + expectedWatchers := []hierarchy.Watcher{ + { + Object: &grafana.MetricsInstance{}, + Owner: client.ObjectKey{Namespace: "default", Name: "grafana-agent-example"}, + Selector: &hierarchy.LabelsSelector{ + NamespaceName: "default", + Labels: labels.SelectorFromSet(labels.Set{"agent": "grafana-agent-example"}), + }, + }, + { + Object: &grafana.LogsInstance{}, + Owner: client.ObjectKey{Namespace: "default", Name: "grafana-agent-example"}, + Selector: &hierarchy.LabelsSelector{ + NamespaceName: "default", + Labels: labels.SelectorFromSet(labels.Set{"agent": "grafana-agent-example"}), + }, + }, + { + Object: &prom.ServiceMonitor{}, + Owner: client.ObjectKey{Namespace: "default", Name: "grafana-agent-example"}, + Selector: &hierarchy.LabelsSelector{ + NamespaceName: "default", + NamespaceLabels: labels.Everything(), + Labels: labels.SelectorFromSet(labels.Set{"instance": "primary"}), + }, + }, + { + Object: &prom.PodMonitor{}, + Owner: client.ObjectKey{Namespace: "default", Name: "grafana-agent-example"}, + Selector: &hierarchy.LabelsSelector{ + NamespaceName: "default", + NamespaceLabels: labels.Everything(), + Labels: labels.SelectorFromSet(labels.Set{"instance": "primary"}), + }, + }, + { + Object: &prom.Probe{}, + Owner: client.ObjectKey{Namespace: "default", Name: "grafana-agent-example"}, + Selector: &hierarchy.LabelsSelector{ + NamespaceName: "default", + Labels: labels.Nothing(), + }, + }, + { + Object: &grafana.PodLogs{}, + Owner: client.ObjectKey{Namespace: "default", Name: "grafana-agent-example"}, + Selector: &hierarchy.LabelsSelector{ + NamespaceName: "default", + NamespaceLabels: labels.Everything(), + Labels: labels.SelectorFromSet(labels.Set{"instance": "primary"}), + }, + }, + { + Object: &v1.Secret{}, + Owner: client.ObjectKey{Namespace: "default", Name: "grafana-agent-example"}, + Selector: &hierarchy.KeySelector{ + Namespace: "default", + Name: "prometheus-fake-credentials", + }, + }, + } + require.ElementsMatch(t, expectedWatchers, watchers) + } +} + +type resourceWalker struct { + onResource func(c client.Object) +} + +func (w *resourceWalker) Visit(v interface{}) (next structwalk.Visitor) { + if v == nil { + return nil + } + if obj, ok := v.(client.Object); ok { + w.onResource(obj) + } + return w +} diff --git a/pkg/operator/config/config.go b/pkg/operator/config/config.go index 011ff65b7a79..50678200ab62 100644 --- a/pkg/operator/config/config.go +++ b/pkg/operator/config/config.go @@ -53,6 +53,8 @@ type Deployment struct { // Logs is the set of logging instances discovered from the root Agent // resource. Logs []LogInstance + // Secrets that can be referenced in the deployment. + Secrets assets.SecretStore } // DeepCopy creates a deep copy of d. diff --git a/pkg/operator/config/config_references.go b/pkg/operator/config/config_references.go index 390064505306..f1ba0956957b 100644 --- a/pkg/operator/config/config_references.go +++ b/pkg/operator/config/config_references.go @@ -17,13 +17,18 @@ type AssetReference struct { // the deployment. Every used secret and configmap should then be loaded into // an assets.SecretStore. func (d *Deployment) AssetReferences() []AssetReference { + return AssetReferences(d) +} + +// AssetReferences returns all secret or configmap selectors used throughout v. +func AssetReferences(v interface{}) []AssetReference { var refs []AssetReference w := assetReferencesWalker{ addReference: func(ar AssetReference) { refs = append(refs, ar) }, } - structwalk.Walk(&w, d) + structwalk.Walk(&w, v) return refs } diff --git a/pkg/operator/deployment_builder.go b/pkg/operator/deployment_builder.go deleted file mode 100644 index 0045381bdc6f..000000000000 --- a/pkg/operator/deployment_builder.go +++ /dev/null @@ -1,290 +0,0 @@ -package operator - -import ( - "context" - "fmt" - - "github.com/go-kit/log" - "github.com/go-kit/log/level" - grafana "github.com/grafana/agent/pkg/operator/apis/monitoring/v1alpha1" - "github.com/grafana/agent/pkg/operator/assets" - "github.com/grafana/agent/pkg/operator/config" - prom "github.com/prometheus-operator/prometheus-operator/pkg/apis/monitoring/v1" - "k8s.io/apimachinery/pkg/api/meta" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/runtime" - "sigs.k8s.io/controller-runtime/pkg/client" -) - -type deploymentBuilder struct { - client.Client - - Logger log.Logger - Config *Config - Agent *grafana.GrafanaAgent - Secrets assets.SecretStore - - // ResourceSelectors is filled as objects are found and can be used to - // trigger future reconciles. - ResourceSelectors map[secondaryResource][]resourceSelector -} - -func (b *deploymentBuilder) Build(ctx context.Context) (config.Deployment, error) { - metricsInstanceSel, err := b.buildResourceSelector(b.Agent.MetricsInstanceSelector()) - if err != nil { - return config.Deployment{}, fmt.Errorf("failed to build MetricsInstance selector: %w", err) - } - b.addSelector(resourcePromInstance, metricsInstanceSel) - - rootMetricInstances, err := b.getMetricsInstances(ctx, metricsInstanceSel) - if err != nil { - return config.Deployment{}, err - } - metricInstances := make([]config.MetricsInstance, 0, len(rootMetricInstances)) - - for _, inst := range rootMetricInstances { - // Get resource selectors for ServiceMonitors, PodMonitors, Probes - var ( - sMonSel, pMonSel, probeSel resourceSelector - ) - getters := []struct { - name string - rsel *resourceSelector - osel grafana.ObjectSelector - }{ - {"ServiceMonitor", &sMonSel, inst.ServiceMonitorSelector()}, - {"PodMonitor", &pMonSel, inst.PodMonitorSelector()}, - {"Probe", &probeSel, inst.ProbeSelector()}, - } - for _, g := range getters { - var err error - *g.rsel, err = b.buildResourceSelector(g.osel) - if err != nil { - return config.Deployment{}, fmt.Errorf("failed to build %s selector: %w", g.name, err) - } - } - - // Use resource selectors to look up objects - sMons, err := b.getServiceMonitors(ctx, sMonSel) - if err != nil { - return config.Deployment{}, fmt.Errorf("unable to fetch ServiceMonitors: %w", err) - } - pMons, err := b.getPodMonitors(ctx, pMonSel) - if err != nil { - return config.Deployment{}, fmt.Errorf("unable to fetch PodMonitors: %w", err) - } - probes, err := b.getProbes(ctx, probeSel) - if err != nil { - return config.Deployment{}, fmt.Errorf("unable to fetch Probes: %w", err) - } - - metricInstances = append(metricInstances, config.MetricsInstance{ - Instance: inst, - ServiceMonitors: sMons, - PodMonitors: pMons, - Probes: probes, - }) - - b.addSelector(resourceServiceMonitor, sMonSel) - b.addSelector(resourcePodMonitor, pMonSel) - b.addSelector(resourceProbe, probeSel) - } - - logsInstanceSel, err := b.buildResourceSelector(b.Agent.LogsInstanceSelector()) - if err != nil { - return config.Deployment{}, fmt.Errorf("failed to build LogsInstance selector: %w", err) - } - b.addSelector(resourceLogsInstance, logsInstanceSel) - - rootLogsInstances, err := b.getLogsInstances(ctx, logsInstanceSel) - if err != nil { - return config.Deployment{}, err - } - logsInstances := make([]config.LogInstance, 0, len(rootLogsInstances)) - - for _, inst := range rootLogsInstances { - podLogsSel, err := b.buildResourceSelector(inst.PodLogsInstanceSelector()) - if err != nil { - return config.Deployment{}, fmt.Errorf("failed to build PodLogs selector: %w", err) - } - podLogs, err := b.getPodLogs(ctx, podLogsSel) - if err != nil { - return config.Deployment{}, fmt.Errorf("unable to fetch PodLogs: %w", err) - } - - logsInstances = append(logsInstances, config.LogInstance{ - Instance: inst, - PodLogs: podLogs, - }) - - b.addSelector(resourcePodLogs, podLogsSel) - } - - return config.Deployment{ - Agent: b.Agent, - Metrics: metricInstances, - Logs: logsInstances, - }, nil -} - -// buildResourceSelector builds a selector for discovering objects with a -// namespace selector and an object selector. If the namespace -// selector, it will default to finding everything in the parent -// namespace. -func (b *deploymentBuilder) buildResourceSelector(sel grafana.ObjectSelector) (resourceSelector, error) { - var namespaceFilter string - if sel.NamespaceSelector == nil { - // When there's no namespaceSelector defined, default to looking in the - // current namespace and matching everything within that namespace. - namespaceFilter = sel.ParentNamespace - sel.NamespaceSelector = &metav1.LabelSelector{} - } - - nsLabels, err := metav1.LabelSelectorAsSelector(sel.NamespaceSelector) - if err != nil { - return nil, fmt.Errorf("failed to convert object namespace selector into label selector: %w", err) - } - objectLabels, err := metav1.LabelSelectorAsSelector(sel.Labels) - if err != nil { - return nil, fmt.Errorf("failed to convert object selector into label selector: %w", err) - } - - var ss []resourceSelector - - if namespaceFilter != "" { - ss = append(ss, &namespaceSelector{Namespace: namespaceFilter}) - } - ss = append(ss, &namespaceLabelSelector{Selector: nsLabels}) - ss = append(ss, &labelSelector{Selector: objectLabels}) - return &multiSelector{Selectors: ss}, nil -} - -// addSelector registers the given selector for the resource to use for update -// tracking. -func (b *deploymentBuilder) addSelector(res secondaryResource, sel resourceSelector) { - b.ResourceSelectors[res] = append(b.ResourceSelectors[res], sel) -} - -func (b *deploymentBuilder) getMetricsInstances(ctx context.Context, sel resourceSelector) ([]*grafana.MetricsInstance, error) { - var list grafana.MetricsInstanceList - if err := b.list(ctx, &list, sel); err != nil { - return nil, fmt.Errorf("unable to discover MetricsInstance: %w", err) - } - return list.Items, nil -} - -// list finds all objects for sel and fills them into list. -func (b *deploymentBuilder) list( - ctx context.Context, - list client.ObjectList, - sel resourceSelector, -) error { - - var lo client.ListOptions - sel.SetListOptions(&lo) - if err := b.List(ctx, list, &lo); err != nil { - return fmt.Errorf("failed to list objects: %w", err) - } - - elements, err := meta.ExtractList(list) - if err != nil { - return fmt.Errorf("failed to get list: %w", err) - } - - filteredElements := make([]runtime.Object, 0, len(elements)) - for _, e := range elements { - o, ok := e.(client.Object) - if !ok { - return fmt.Errorf("unexpected object returned") - } - - if sel.Matches(b.Logger, b.Client, o) { - filteredElements = append(filteredElements, e) - } - } - - if err := meta.SetList(list, filteredElements); err != nil { - return fmt.Errorf("failed to populate list of objects: %w", err) - } - return nil -} - -func (b *deploymentBuilder) getServiceMonitors(ctx context.Context, sel resourceSelector) ([]*prom.ServiceMonitor, error) { - var list prom.ServiceMonitorList - if err := b.list(ctx, &list, sel); err != nil { - return nil, fmt.Errorf("unable to discover ServiceMonitors: %w", err) - } - - items := make([]*prom.ServiceMonitor, 0, len(list.Items)) -Item: - for _, item := range list.Items { - if b.Agent.Spec.Metrics.ArbitraryFSAccessThroughSMs.Deny { - for _, ep := range item.Spec.Endpoints { - err := testForArbitraryFSAccess(ep) - if err == nil { - continue - } - - level.Warn(b.Logger).Log( - "msg", "skipping service monitor", - "agent", client.ObjectKeyFromObject(b.Agent), - "servicemonitor", client.ObjectKeyFromObject(item), - "err", err, - ) - - continue Item - } - } - items = append(items, item) - } - - return items, nil -} - -func testForArbitraryFSAccess(e prom.Endpoint) error { - if e.BearerTokenFile != "" { - return fmt.Errorf("it accesses file system via bearer token file which is disallowed via GrafanaAgent specification") - } - - if e.TLSConfig == nil { - return nil - } - - if e.TLSConfig.CAFile != "" || e.TLSConfig.CertFile != "" || e.TLSConfig.KeyFile != "" { - return fmt.Errorf("it accesses file system via TLS config which is disallowed via GrafanaAgent specification") - } - - return nil -} - -func (b *deploymentBuilder) getPodMonitors(ctx context.Context, sel resourceSelector) ([]*prom.PodMonitor, error) { - var list prom.PodMonitorList - if err := b.list(ctx, &list, sel); err != nil { - return nil, fmt.Errorf("unable to discover PodMonitors: %w", err) - } - return list.Items, nil -} - -func (b *deploymentBuilder) getProbes(ctx context.Context, sel resourceSelector) ([]*prom.Probe, error) { - var list prom.ProbeList - if err := b.list(ctx, &list, sel); err != nil { - return nil, fmt.Errorf("unable to discover Probes: %w", err) - } - return list.Items, nil -} - -func (b *deploymentBuilder) getLogsInstances(ctx context.Context, sel resourceSelector) ([]*grafana.LogsInstance, error) { - var list grafana.LogsInstanceList - if err := b.list(ctx, &list, sel); err != nil { - return nil, fmt.Errorf("unable to discover LogsInstances: %w", err) - } - return list.Items, nil -} - -func (b *deploymentBuilder) getPodLogs(ctx context.Context, sel resourceSelector) ([]*grafana.PodLogs, error) { - var list grafana.PodLogsList - if err := b.list(ctx, &list, sel); err != nil { - return nil, fmt.Errorf("unable to discover PodLogs: %w", err) - } - return list.Items, nil -} diff --git a/pkg/operator/hierarchy/hierarchy.go b/pkg/operator/hierarchy/hierarchy.go new file mode 100644 index 000000000000..05c5658bb3d7 --- /dev/null +++ b/pkg/operator/hierarchy/hierarchy.go @@ -0,0 +1,149 @@ +// Package hierarchy provides tools to discover a resource hierarchy. A +// resource hierarchy is made when a resource has a set of rules to discover +// other resources. +package hierarchy + +import ( + "context" + "fmt" + "sync" + "time" + + "github.com/go-kit/log" + "github.com/go-kit/log/level" + "k8s.io/apimachinery/pkg/runtime/schema" + "k8s.io/client-go/util/workqueue" + "sigs.k8s.io/controller-runtime/pkg/client" + "sigs.k8s.io/controller-runtime/pkg/client/apiutil" + "sigs.k8s.io/controller-runtime/pkg/event" + "sigs.k8s.io/controller-runtime/pkg/handler" + "sigs.k8s.io/controller-runtime/pkg/reconcile" +) + +// Notifier can be attached to a controller and generate reconciles when +// objects inside of a resource hierarchy change. +type Notifier struct { + log log.Logger + client client.Client + + watchersMut sync.RWMutex + watchers map[schema.GroupVersionKind][]Watcher +} + +// Watcher is something watching for changes to a resource. +type Watcher struct { + Object client.Object // Object to watch for events against. + Owner client.ObjectKey // Owner to receive a reconcile for. + Selector Selector // Selector to use to match changed objects. +} + +// NewNotifier creates a new Notifier which uses the provided client for +// performing hierarchy lookups. +func NewNotifier(l log.Logger, cli client.Client) *Notifier { + return &Notifier{ + log: l, + client: cli, + watchers: make(map[schema.GroupVersionKind][]Watcher), + } +} + +// EventHandler returns an event handler that can be given to +// controller.Watches. +// +// controller.Watches should be called once per type in the resource hierarchy. +// Each call to controller.Watches should use the same Notifier. +func (n *Notifier) EventHandler() handler.EventHandler { + // TODO(rfratto): It's possible to create a custom implementation of + // source.Source so we wouldn't have to call controller.Watches a bunch of + // times. I played around a little with an implementation but it was going to + // be a lot of work to dynamically spin up/down informers, so I put it aside + // for now. Maybe it's an improvement for the future. + return ¬ifierEventHandler{Notifier: n} +} + +// Notify configures reconciles to be generated for a set of watchers when +// watched resources change. +// +// Notify appends to the list of watchers. To remove out notifications for a +// specific owner, call StopNotify. +func (n *Notifier) Notify(watchers ...Watcher) error { + n.watchersMut.Lock() + defer n.watchersMut.Unlock() + + for _, w := range watchers { + gvk, err := apiutil.GVKForObject(w.Object, n.client.Scheme()) + if err != nil { + return fmt.Errorf("could not get GVK: %w", err) + } + + n.watchers[gvk] = append(n.watchers[gvk], w) + } + + return nil +} + +// StopNotify removes all watches for a specific owner. +func (n *Notifier) StopNotify(owner client.ObjectKey) { + n.watchersMut.Lock() + defer n.watchersMut.Unlock() + + for key, watchers := range n.watchers { + rem := make([]Watcher, 0, len(watchers)) + for _, w := range watchers { + if w.Owner != owner { + rem = append(rem, w) + } + } + n.watchers[key] = rem + } +} + +type notifierEventHandler struct { + *Notifier +} + +var _ handler.EventHandler = (*notifierEventHandler)(nil) + +func (h *notifierEventHandler) Create(ev event.CreateEvent, q workqueue.RateLimitingInterface) { + h.handleEvent(ev.Object, q) +} + +func (h *notifierEventHandler) Update(ev event.UpdateEvent, q workqueue.RateLimitingInterface) { + h.handleEvent(ev.ObjectOld, q) + h.handleEvent(ev.ObjectNew, q) +} + +func (h *notifierEventHandler) Delete(ev event.DeleteEvent, q workqueue.RateLimitingInterface) { + h.handleEvent(ev.Object, q) +} + +func (h *notifierEventHandler) Generic(ev event.GenericEvent, q workqueue.RateLimitingInterface) { + h.handleEvent(ev.Object, q) +} + +func (h *notifierEventHandler) handleEvent(obj client.Object, q workqueue.RateLimitingInterface) { + h.watchersMut.RLock() + defer h.watchersMut.RUnlock() + + gvk, err := apiutil.GVKForObject(obj, h.client.Scheme()) + if err != nil { + level.Error(h.log).Log("msg", "failed to get gvk for object", "err", err) + return + } + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) + defer cancel() + + // Iterate through all of the watchers for the gvk and check to see if we + // should trigger a reconcile. + for _, watcher := range h.watchers[gvk] { + matches, err := watcher.Selector.Matches(ctx, h.client, obj) + if err != nil { + level.Error(h.log).Log("msg", "failed to handle notifier event", "err", err) + return + } + if matches { + q.Add(reconcile.Request{NamespacedName: watcher.Owner}) + } + } +} diff --git a/pkg/operator/hierarchy/hierarchy_test.go b/pkg/operator/hierarchy/hierarchy_test.go new file mode 100644 index 000000000000..a844b647f961 --- /dev/null +++ b/pkg/operator/hierarchy/hierarchy_test.go @@ -0,0 +1,144 @@ +//go:build !nonetwork && !nodocker && !race +// +build !nonetwork,!nodocker,!race + +package hierarchy + +import ( + "context" + "testing" + "time" + + "github.com/go-kit/log" + "github.com/grafana/agent/pkg/util/k8s" + "github.com/stretchr/testify/require" + v1 "k8s.io/api/core/v1" + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/labels" + "k8s.io/apimachinery/pkg/types" + "k8s.io/client-go/util/workqueue" + "sigs.k8s.io/controller-runtime/pkg/event" +) + +// TestNotifier tests that notifier properly handles events for changed +// objects. +func TestNotifier(t *testing.T) { + l := log.NewNopLogger() + + ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute) + defer cancel() + + cluster, err := k8s.NewCluster(ctx, k8s.Options{}) + require.NoError(t, err) + defer cluster.Stop() + + cli := cluster.Client() + + // Tests will rely on a namespace existing, so let's create a namespace with + // some labels. + testNs := v1.Namespace{ + TypeMeta: metav1.TypeMeta{ + APIVersion: "v1", + Kind: "Namespace", + }, + ObjectMeta: metav1.ObjectMeta{ + Name: "enqueue-test", + Labels: map[string]string{"foo": "bar"}, + }, + } + err = cli.Create(ctx, &testNs) + require.NoError(t, err) + + testPod := &v1.Pod{ObjectMeta: metav1.ObjectMeta{ + Name: "test-pod", + Namespace: "enqueue-test", + Labels: map[string]string{"fizz": "buzz"}, + }} + + tt := []struct { + name string + sel Selector + expectEnqueue bool + }{ + { + name: "no watchers", + sel: nil, + expectEnqueue: false, + }, + { + name: "matches watcher", + sel: &LabelsSelector{ + NamespaceName: "enqueue-test", + NamespaceLabels: parseSelector(t, "foo in (bar)"), + Labels: parseSelector(t, "fizz in (buzz)"), + }, + expectEnqueue: true, + }, + { + name: "matches watcher with explicit namespace", + sel: &LabelsSelector{ + NamespaceName: "enqueue-test", + Labels: parseSelector(t, "fizz in (buzz)"), + }, + expectEnqueue: true, + }, + { + name: "bad namespace name selector", + sel: &LabelsSelector{ + NamespaceName: "default", + Labels: labels.Everything(), + }, + expectEnqueue: false, + }, + { + name: "bad namespace label selector", + sel: &LabelsSelector{ + NamespaceName: "enqueue-test", + NamespaceLabels: parseSelector(t, "foo notin (bar)"), + Labels: labels.Everything(), + }, + expectEnqueue: false, + }, + { + name: "bad label selector", + sel: &LabelsSelector{ + NamespaceName: "default", + NamespaceLabels: labels.Everything(), + Labels: parseSelector(t, "fizz notin (buzz)"), + }, + expectEnqueue: false, + }, + } + + for _, tc := range tt { + t.Run(tc.name, func(t *testing.T) { + limiter := workqueue.DefaultControllerRateLimiter() + q := workqueue.NewRateLimitingQueue(limiter) + + notifier := NewNotifier(l, cli) + + if tc.sel != nil { + err := notifier.Notify(Watcher{ + Object: &v1.Pod{}, + Owner: types.NamespacedName{Name: "watcher", Namespace: "enqueue-test"}, + Selector: tc.sel, + }) + require.NoError(t, err) + } + + e := notifier.EventHandler() + e.Create(event.CreateEvent{Object: testPod}, q) + if tc.expectEnqueue { + require.Equal(t, 1, q.Len(), "expected change enqueue") + } else { + require.Equal(t, 0, q.Len(), "no changes should have been enqueued") + } + }) + } +} + +func parseSelector(t *testing.T, selector string) labels.Selector { + t.Helper() + s, err := labels.Parse(selector) + require.NoError(t, err) + return s +} diff --git a/pkg/operator/hierarchy/list.go b/pkg/operator/hierarchy/list.go new file mode 100644 index 000000000000..016fdc33da6c --- /dev/null +++ b/pkg/operator/hierarchy/list.go @@ -0,0 +1,50 @@ +package hierarchy + +import ( + "context" + "fmt" + + "k8s.io/apimachinery/pkg/api/meta" + "k8s.io/apimachinery/pkg/runtime" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// List will populate list with elements that match sel. +func List(ctx context.Context, cli client.Client, list client.ObjectList, sel Selector) error { + if err := cli.List(ctx, list, sel); err != nil { + return fmt.Errorf("list failed: %w", err) + } + if err := filterList(ctx, cli, list, sel); err != nil { + return fmt.Errorf("filter failed: %w", err) + } + return nil +} + +// filterList updates the provided list to only elements which match sel. +func filterList(ctx context.Context, cli client.Client, list client.ObjectList, sel Selector) error { + allElements, err := meta.ExtractList(list) + if err != nil { + return fmt.Errorf("failed to get list: %w", err) + } + + filtered := make([]runtime.Object, 0, len(allElements)) + for _, element := range allElements { + obj, ok := element.(client.Object) + if !ok { + return fmt.Errorf("unexpected object of type %T in list", element) + } + + matches, err := sel.Matches(ctx, cli, obj) + if err != nil { + return fmt.Errorf("failed to validate object: %w", err) + } + if matches { + filtered = append(filtered, obj) + } + } + + if err := meta.SetList(list, filtered); err != nil { + return fmt.Errorf("failed to update list: %w", err) + } + return nil +} diff --git a/pkg/operator/hierarchy/selector.go b/pkg/operator/hierarchy/selector.go new file mode 100644 index 000000000000..06454e3a8f1e --- /dev/null +++ b/pkg/operator/hierarchy/selector.go @@ -0,0 +1,88 @@ +package hierarchy + +import ( + "context" + "fmt" + + corev1 "k8s.io/api/core/v1" + "k8s.io/apimachinery/pkg/fields" + "k8s.io/apimachinery/pkg/labels" + "sigs.k8s.io/controller-runtime/pkg/client" +) + +// Selector finding objects within the resource hierarchy. +type Selector interface { + // ListOption can be passed to List to perform initial filtering of returned + // objects. + client.ListOption + + // Matches returns true if the Selector matches the provided Object. The + // provided Client may be used to perform extra searches. + Matches(context.Context, client.Client, client.Object) (bool, error) +} + +// LabelsSelector is used for discovering a set of objects in a hierarchy based +// on labels. +type LabelsSelector struct { + // NamespaceName is the default namespace to search for objects in when + // NamespaceSelector is nil. + NamespaceName string + + // NamespaceLabels causes all namespaces whose labels match NamespaceLabels + // to be searched. When nil, only the namespace specified by NamespaceName + // will be searched. + NamespaceLabels labels.Selector + + // Labels discovers all objects whose labels match the selector. If nil, + // no objects will be discovered. + Labels labels.Selector +} + +var _ Selector = (*LabelsSelector)(nil) + +// ApplyToList implements Selector. +func (ls *LabelsSelector) ApplyToList(lo *client.ListOptions) { + if ls.NamespaceLabels == nil { + lo.Namespace = ls.NamespaceName + } + lo.LabelSelector = ls.Labels +} + +// Matches implements Selector. +func (ls *LabelsSelector) Matches(ctx context.Context, cli client.Client, o client.Object) (bool, error) { + if !ls.Labels.Matches(labels.Set(o.GetLabels())) { + return false, nil + } + + // Fast path: we don't need to retrieve the labels of the namespace. + if ls.NamespaceLabels == nil { + return o.GetNamespace() == ls.NamespaceName, nil + } + + // Slow path: we need to look up the namespace to see if its labels match. As + // long as cli implements caching, this won't be too bad. + var ns corev1.Namespace + if err := cli.Get(ctx, client.ObjectKey{Name: o.GetNamespace()}, &ns); err != nil { + return false, fmt.Errorf("error looking up namespace %q: %w", o.GetNamespace(), err) + } + return ls.NamespaceLabels.Matches(labels.Set(ns.GetLabels())), nil +} + +// KeySelector is used for discovering a single object based on namespace and +// name. +type KeySelector struct { + Namespace, Name string +} + +var _ Selector = (*KeySelector)(nil) + +// ApplyToList implements Selector. +func (ks *KeySelector) ApplyToList(lo *client.ListOptions) { + lo.Namespace = ks.Namespace + lo.FieldSelector = fields.OneTermEqualSelector("metadata.name", ks.Name) +} + +// Matches implements Selector. +func (ks *KeySelector) Matches(ctx context.Context, cli client.Client, o client.Object) (bool, error) { + return ks.Name == o.GetName() && ks.Namespace == o.GetNamespace(), nil +} diff --git a/pkg/operator/operator.go b/pkg/operator/operator.go index edb66c3535bf..f78b93c6efc6 100644 --- a/pkg/operator/operator.go +++ b/pkg/operator/operator.go @@ -13,8 +13,6 @@ import ( "k8s.io/apimachinery/pkg/runtime" controller "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/builder" - "sigs.k8s.io/controller-runtime/pkg/client" - "sigs.k8s.io/controller-runtime/pkg/client/apiutil" "sigs.k8s.io/controller-runtime/pkg/healthz" "sigs.k8s.io/controller-runtime/pkg/manager" "sigs.k8s.io/controller-runtime/pkg/predicate" @@ -22,6 +20,7 @@ import ( "sigs.k8s.io/controller-runtime/pkg/source" grafana_v1alpha1 "github.com/grafana/agent/pkg/operator/apis/monitoring/v1alpha1" + "github.com/grafana/agent/pkg/operator/hierarchy" promop_v1 "github.com/prometheus-operator/prometheus-operator/pkg/apis/monitoring/v1" promop "github.com/prometheus-operator/prometheus-operator/pkg/operator" apps_v1 "k8s.io/api/apps/v1" @@ -138,12 +137,10 @@ func New(l log.Logger, c *Config) (*Operator, error) { } var ( - events = newResourceEventHandlers(manager.GetClient(), l) - - applyGVK = func(obj client.Object) client.Object { return applyGVK(obj, manager) } - watchType = func(obj client.Object) source.Source { return watchType(obj, manager) } - agentPredicates []predicate.Predicate + + notifier = hierarchy.NewNotifier(log.With(l, "component", "hierarchy_notifier"), manager.GetClient()) + notifierHandler = notifier.EventHandler() ) // Initialize agentPredicates if an GrafanaAgent selector is configured. @@ -168,9 +165,9 @@ func New(l log.Logger, c *Config) (*Operator, error) { kubeletName := parts[1] err := controller.NewControllerManagedBy(manager). - For(applyGVK(&core_v1.Node{})). - Owns(applyGVK(&core_v1.Service{})). - Owns(applyGVK(&core_v1.Endpoints{})). + For(&core_v1.Node{}). + Owns(&core_v1.Service{}). + Owns(&core_v1.Endpoints{}). Complete(&lazyKubeletReconciler) if err != nil { return nil, fmt.Errorf("failed to create kubelet controller: %w", err) @@ -185,30 +182,30 @@ func New(l log.Logger, c *Config) (*Operator, error) { } err = controller.NewControllerManagedBy(manager). - For(applyGVK(&grafana_v1alpha1.GrafanaAgent{}), builder.WithPredicates(agentPredicates...)). - Owns(applyGVK(&apps_v1.StatefulSet{})). - Owns(applyGVK(&apps_v1.DaemonSet{})). - Owns(applyGVK(&core_v1.Secret{})). - Owns(applyGVK(&core_v1.Service{})). - Watches(watchType(&core_v1.Secret{}), events[resourceSecret]). - Watches(watchType(&grafana_v1alpha1.LogsInstance{}), events[resourceLogsInstance]). - Watches(watchType(&grafana_v1alpha1.PodLogs{}), events[resourcePodLogs]). - Watches(watchType(&grafana_v1alpha1.MetricsInstance{}), events[resourcePromInstance]). - Watches(watchType(&promop_v1.PodMonitor{}), events[resourcePodMonitor]). - Watches(watchType(&promop_v1.Probe{}), events[resourceProbe]). - Watches(watchType(&promop_v1.ServiceMonitor{}), events[resourceServiceMonitor]). - Watches(watchType(&core_v1.Secret{}), events[resourceSecret]). - Watches(watchType(&core_v1.ConfigMap{}), events[resourceConfigMap]). + For(&grafana_v1alpha1.GrafanaAgent{}, builder.WithPredicates(agentPredicates...)). + Owns(&apps_v1.StatefulSet{}). + Owns(&apps_v1.DaemonSet{}). + Owns(&core_v1.Secret{}). + Owns(&core_v1.Service{}). + Watches(&source.Kind{Type: &core_v1.Secret{}}, notifierHandler). + Watches(&source.Kind{Type: &grafana_v1alpha1.LogsInstance{}}, notifierHandler). + Watches(&source.Kind{Type: &grafana_v1alpha1.PodLogs{}}, notifierHandler). + Watches(&source.Kind{Type: &grafana_v1alpha1.MetricsInstance{}}, notifierHandler). + Watches(&source.Kind{Type: &promop_v1.PodMonitor{}}, notifierHandler). + Watches(&source.Kind{Type: &promop_v1.Probe{}}, notifierHandler). + Watches(&source.Kind{Type: &promop_v1.ServiceMonitor{}}, notifierHandler). + Watches(&source.Kind{Type: &core_v1.Secret{}}, notifierHandler). + Watches(&source.Kind{Type: &core_v1.ConfigMap{}}, notifierHandler). Complete(&lazyAgentReconciler) if err != nil { return nil, fmt.Errorf("failed to create GrafanaAgent controller: %w", err) } lazyAgentReconciler.Set(&reconciler{ - Client: manager.GetClient(), - scheme: manager.GetScheme(), - eventHandlers: events, - config: c, + Client: manager.GetClient(), + scheme: manager.GetScheme(), + notifier: notifier, + config: c, }) return &Operator{ @@ -225,26 +222,6 @@ func (o *Operator) Start(ctx context.Context) error { return o.manager.Start(ctx) } -// watchType applies the GVK to an object and returns a source to watch it. -// watchType is a convenience function; without it, the GVK won't show up in -// logs. -func watchType(obj client.Object, m manager.Manager) source.Source { - applyGVK(obj, m) - return &source.Kind{Type: obj} -} - -// applyGVK applies a GVK to an object based on the scheme. applyGVK is a -// convenience function; without it, the GVK won't show up in logs. -// nolint: interfacer -func applyGVK(obj client.Object, m manager.Manager) client.Object { - gvk, err := apiutil.GVKForObject(obj, m.GetScheme()) - if err != nil { - panic(err) - } - obj.GetObjectKind().SetGroupVersionKind(gvk) - return obj -} - type lazyReconciler struct { mut sync.RWMutex inner reconcile.Reconciler diff --git a/pkg/operator/reconciler.go b/pkg/operator/reconciler.go index c1476e7ef872..a1c5208be195 100644 --- a/pkg/operator/reconciler.go +++ b/pkg/operator/reconciler.go @@ -10,12 +10,12 @@ import ( "github.com/grafana/agent/pkg/operator/assets" "github.com/grafana/agent/pkg/operator/clientutil" "github.com/grafana/agent/pkg/operator/config" + "github.com/grafana/agent/pkg/operator/hierarchy" "github.com/grafana/agent/pkg/operator/logutil" core_v1 "k8s.io/api/core/v1" k8s_errors "k8s.io/apimachinery/pkg/api/errors" v1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/runtime" - "k8s.io/apimachinery/pkg/types" controller "sigs.k8s.io/controller-runtime" "sigs.k8s.io/controller-runtime/pkg/client" ) @@ -25,7 +25,7 @@ type reconciler struct { scheme *runtime.Scheme config *Config - eventHandlers eventHandlers + notifier *hierarchy.Notifier } func (r *reconciler) Reconcile(ctx context.Context, req controller.Request) (controller.Result, error) { @@ -33,11 +33,12 @@ func (r *reconciler) Reconcile(ctx context.Context, req controller.Request) (con level.Info(l).Log("msg", "reconciling grafana-agent") defer level.Debug(l).Log("msg", "done reconciling grafana-agent") + // Reset our notifications while we re-handle the reconcile. + r.notifier.StopNotify(req.NamespacedName) + var agent grafana_v1alpha1.GrafanaAgent if err := r.Get(ctx, req.NamespacedName, &agent); k8s_errors.IsNotFound(err) { - level.Debug(l).Log("msg", "detected deleted agent, cleaning up watchers") - r.eventHandlers.Clear(req.NamespacedName) - + level.Debug(l).Log("msg", "detected deleted agent") return controller.Result{}, nil } else if err != nil { level.Error(l).Log("msg", "unable to get grafana-agent", "err", err) @@ -48,37 +49,13 @@ func (r *reconciler) Reconcile(ctx context.Context, req controller.Request) (con return controller.Result{}, nil } - secrets := make(assets.SecretStore) - builder := deploymentBuilder{ - Logger: l, - Config: r.config, - Client: r.Client, - Agent: &agent, - Secrets: secrets, - ResourceSelectors: make(map[secondaryResource][]resourceSelector), - } - - deployment, err := builder.Build(ctx) + deployment, watchers, err := buildHierarchy(ctx, l, r.Client, &agent) if err != nil { - level.Error(l).Log("msg", "unable to collect resources", "err", err) + level.Error(l).Log("msg", "unable to build hierarchy", "err", err) return controller.Result{}, nil } - - // Update our notifiers with the objects we discovered from building the - // deployment. This allows us to re-reconcile when any of the objects that - // composed our final deployment changes. - for _, secondary := range secondaryResources { - r.eventHandlers[secondary].Notify(req.NamespacedName, builder.ResourceSelectors[secondary]) - } - - assetRefs := deployment.AssetReferences() - - // Update our notifiers with asset references discovered through the deployment. - r.watchSecrets(req, assetRefs) - - // Fill secrets in store - if err := r.fillStore(ctx, assetRefs, secrets); err != nil { - level.Error(l).Log("msg", "unable to cache secrets for building config", "err", err) + if err := r.notifier.Notify(watchers...); err != nil { + level.Error(l).Log("msg", "unable to update notifier", "err", err) return controller.Result{}, nil } @@ -97,7 +74,7 @@ func (r *reconciler) Reconcile(ctx context.Context, req controller.Request) (con r.createLogsDaemonSet, } for _, actor := range actors { - err := actor(ctx, l, deployment, secrets) + err := actor(ctx, l, deployment, deployment.Secrets) if err != nil { level.Error(l).Log("msg", "error during reconciling", "err", err) return controller.Result{Requeue: true}, nil @@ -107,75 +84,6 @@ func (r *reconciler) Reconcile(ctx context.Context, req controller.Request) (con return controller.Result{}, nil } -// watchSecrets will go iterate over asset references and configure them to be -// watched for updates. This allows reconciles to trigger when a referenced -// Secret or ConfigMap changes. -func (r *reconciler) watchSecrets(req controller.Request, refs []config.AssetReference) { - var ( - configMapSelectors []resourceSelector - secretSelectors []resourceSelector - ) - - for _, ref := range refs { - switch { - case ref.Reference.ConfigMap != nil: - configMapSelectors = append(configMapSelectors, &assetReferenceSelector{Reference: ref}) - case ref.Reference.Secret != nil: - secretSelectors = append(secretSelectors, &assetReferenceSelector{Reference: ref}) - default: - panic("unknown AssetReference") - } - } - - r.eventHandlers[resourceConfigMap].Notify(req.NamespacedName, configMapSelectors) - r.eventHandlers[resourceSecret].Notify(req.NamespacedName, secretSelectors) -} - -// fillStore retrieves all the values from refs and caches them in the provided store. -func (r *reconciler) fillStore(ctx context.Context, refs []config.AssetReference, store assets.SecretStore) error { - for _, ref := range refs { - var value string - - if ref.Reference.ConfigMap != nil { - var cm core_v1.ConfigMap - name := types.NamespacedName{ - Namespace: ref.Namespace, - Name: ref.Reference.ConfigMap.Name, - } - - if err := r.Get(ctx, name, &cm); err != nil { - return err - } - - rawValue, ok := cm.BinaryData[ref.Reference.ConfigMap.Key] - if !ok { - return fmt.Errorf("no key %s in ConfigMap %s", ref.Reference.ConfigMap.Key, name) - } - value = string(rawValue) - } else if ref.Reference.Secret != nil { - var secret core_v1.Secret - name := types.NamespacedName{ - Namespace: ref.Namespace, - Name: ref.Reference.Secret.Name, - } - - if err := r.Get(ctx, name, &secret); err != nil { - return err - } - - rawValue, ok := secret.Data[ref.Reference.Secret.Key] - if !ok { - return fmt.Errorf("no key %s in Secret %s", ref.Reference.ConfigMap.Key, name) - } - value = string(rawValue) - } - - store[assets.KeyForSelector(ref.Namespace, &ref.Reference)] = value - } - - return nil -} - // createSecrets creates secrets from the secret store. func (r *reconciler) createSecrets( ctx context.Context, diff --git a/pkg/operator/reconciler_logs.go b/pkg/operator/reconciler_logs.go index 8bdeb301abd0..a12e752a5db8 100644 --- a/pkg/operator/reconciler_logs.go +++ b/pkg/operator/reconciler_logs.go @@ -10,7 +10,6 @@ import ( "github.com/grafana/agent/pkg/operator/clientutil" "github.com/grafana/agent/pkg/operator/config" apps_v1 "k8s.io/api/apps/v1" - k8s_errors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/types" ) @@ -40,20 +39,8 @@ func (r *reconciler) createLogsDaemonSet( key := types.NamespacedName{Namespace: ds.Namespace, Name: ds.Name} if len(d.Logs) == 0 { - var ds apps_v1.DaemonSet - err := r.Client.Get(ctx, key, &ds) - if k8s_errors.IsNotFound(err) || !isManagedResource(&ds) { - return nil - } else if err != nil { - return fmt.Errorf("failed to find stale DaemonSet %s: %w", key, err) - } - - err = r.Client.Delete(ctx, &ds) - if err != nil { - return fmt.Errorf("failed to delete stale DaemonSet %s: %w", key, err) - } - return nil + return deleteManagedResource(ctx, r.Client, key, &ds) } level.Info(l).Log("msg", "reconciling logs daemonset", "ds", key) diff --git a/pkg/operator/reconciler_metrics.go b/pkg/operator/reconciler_metrics.go index 334f94adcd4a..6a52948133b5 100644 --- a/pkg/operator/reconciler_metrics.go +++ b/pkg/operator/reconciler_metrics.go @@ -14,10 +14,10 @@ import ( "github.com/grafana/agent/pkg/operator/config" apps_v1 "k8s.io/api/apps/v1" core_v1 "k8s.io/api/core/v1" - k8s_errors "k8s.io/apimachinery/pkg/api/errors" v1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/labels" "k8s.io/apimachinery/pkg/types" + "k8s.io/utils/pointer" "sigs.k8s.io/controller-runtime/pkg/client" ) @@ -57,18 +57,7 @@ func (r *reconciler) createTelemetryConfigurationSecret( // Delete the old Secret if one exists and we have nothing to create. if !shouldCreate { var secret core_v1.Secret - err := r.Client.Get(ctx, key, &secret) - if k8s_errors.IsNotFound(err) || !isManagedResource(&secret) { - return nil - } else if err != nil { - return fmt.Errorf("failed to find stale secret %s: %w", key, err) - } - - err = r.Client.Delete(ctx, &secret) - if err != nil { - return fmt.Errorf("failed to delete stale secret %s: %w", key, err) - } - return nil + return deleteManagedResource(ctx, r.Client, key, &secret) } rawConfig, err := d.BuildConfig(s, ty) @@ -83,8 +72,6 @@ func (r *reconciler) createTelemetryConfigurationSecret( return fmt.Errorf("unable to build config: %w", err) } - blockOwnerDeletion := true - secret := core_v1.Secret{ ObjectMeta: v1.ObjectMeta{ Namespace: key.Namespace, @@ -92,7 +79,7 @@ func (r *reconciler) createTelemetryConfigurationSecret( Labels: r.config.Labels.Merge(managedByOperatorLabels), OwnerReferences: []v1.OwnerReference{{ APIVersion: d.Agent.APIVersion, - BlockOwnerDeletion: &blockOwnerDeletion, + BlockOwnerDeletion: pointer.Bool(true), Kind: d.Agent.Kind, Name: d.Agent.Name, UID: d.Agent.UID, @@ -121,21 +108,9 @@ func (r *reconciler) createMetricsGoverningService( // Delete the old Secret if one exists and we have no prometheus instances. if len(d.Metrics) == 0 { - key := types.NamespacedName{Namespace: svc.Namespace, Name: svc.Name} - var service core_v1.Service - err := r.Client.Get(ctx, key, &service) - if k8s_errors.IsNotFound(err) || !isManagedResource(&service) { - return nil - } else if err != nil { - return fmt.Errorf("failed to find stale Service %s: %w", key, err) - } - - err = r.Client.Delete(ctx, &service) - if err != nil { - return fmt.Errorf("failed to delete stale Service %s: %w", key, err) - } - return nil + key := types.NamespacedName{Namespace: svc.Namespace, Name: svc.Name} + return deleteManagedResource(ctx, r.Client, key, &service) } level.Info(l).Log("msg", "reconciling statefulset service", "service", svc.Name) @@ -199,11 +174,11 @@ func (r *reconciler) createMetricsStatefulSets( return fmt.Errorf("failed to list statefulsets: %w", err) } for _, ss := range statefulSets.Items { - if _, keep := generated[ss.Name]; keep { + if _, keep := generated[ss.Name]; keep || !isManagedResource(&ss) { continue } level.Info(l).Log("msg", "deleting stale statefulset", "name", ss.Name) - if err := r.Client.Delete(ctx, &ss); err != nil { + if err := r.Delete(ctx, &ss); err != nil { return fmt.Errorf("failed to delete stale statefulset %s: %w", ss.Name, err) } } diff --git a/pkg/operator/resources_logs.go b/pkg/operator/resources_logs.go index 81d7119b2477..85dc07b6cce9 100644 --- a/pkg/operator/resources_logs.go +++ b/pkg/operator/resources_logs.go @@ -11,6 +11,7 @@ import ( v1 "k8s.io/api/core/v1" meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/util/intstr" + "k8s.io/utils/pointer" ) func generateLogsDaemonSet( @@ -32,7 +33,7 @@ func generateLogsDaemonSet( // Don't transfer any kubectl annotations to the DaemonSet so it doesn't get // pruned by kubectl. annotations := make(map[string]string) - for k, v := range d.Agent.ObjectMeta.Annotations { + for k, v := range d.Agent.Annotations { if !strings.HasPrefix(k, "kubectl.kubernetes.io/") { annotations[k] = v } @@ -46,7 +47,6 @@ func generateLogsDaemonSet( labels[agentTypeLabel] = "logs" labels[managedByOperatorLabel] = managedByOperatorLabelValue - boolTrue := true ds := &apps_v1.DaemonSet{ ObjectMeta: meta_v1.ObjectMeta{ Name: name, @@ -56,8 +56,8 @@ func generateLogsDaemonSet( OwnerReferences: []meta_v1.OwnerReference{{ APIVersion: d.Agent.APIVersion, Kind: d.Agent.Kind, - BlockOwnerDeletion: &boolTrue, - Controller: &boolTrue, + BlockOwnerDeletion: pointer.Bool(true), + Controller: pointer.Bool(true), Name: d.Agent.Name, UID: d.Agent.UID, }}, @@ -260,13 +260,6 @@ func generateLogsDaemonSetSpec( Value: "0", }} - var ( - privileged bool = true - runAsUser int64 = 0 - - terminationGracePeriodSeconds = int64(4800) - ) - operatorContainers := []v1.Container{ { Name: "config-reloader", @@ -274,8 +267,8 @@ func generateLogsDaemonSetSpec( VolumeMounts: volumeMounts, Env: envVars, SecurityContext: &v1.SecurityContext{ - Privileged: &privileged, - RunAsUser: &runAsUser, + Privileged: pointer.Bool(true), + RunAsUser: pointer.Int64(0), }, Args: []string{ "--config-file=/var/lib/grafana-agent/config-in/agent.yml", @@ -336,7 +329,7 @@ func generateLogsDaemonSetSpec( ServiceAccountName: d.Agent.Spec.ServiceAccountName, NodeSelector: d.Agent.Spec.NodeSelector, PriorityClassName: d.Agent.Spec.PriorityClassName, - TerminationGracePeriodSeconds: &terminationGracePeriodSeconds, + TerminationGracePeriodSeconds: pointer.Int64(4800), Volumes: volumes, Tolerations: d.Agent.Spec.Tolerations, Affinity: d.Agent.Spec.Affinity, diff --git a/pkg/operator/resources_metrics.go b/pkg/operator/resources_metrics.go index 69fd583966f0..4834fb8e9793 100644 --- a/pkg/operator/resources_metrics.go +++ b/pkg/operator/resources_metrics.go @@ -1,6 +1,7 @@ package operator import ( + "context" "fmt" "strings" @@ -10,8 +11,10 @@ import ( prom_operator "github.com/prometheus-operator/prometheus-operator/pkg/operator" apps_v1 "k8s.io/api/apps/v1" v1 "k8s.io/api/core/v1" + k8s_errors "k8s.io/apimachinery/pkg/api/errors" meta_v1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/util/intstr" + "k8s.io/utils/pointer" "sigs.k8s.io/controller-runtime/pkg/client" ) @@ -33,6 +36,22 @@ var ( probeTimeoutSeconds int32 = 3 ) +// deleteManagedResource deletes a managed resource. Ignores resources that are +// not managed. +func deleteManagedResource(ctx context.Context, cli client.Client, key client.ObjectKey, o client.Object) error { + err := cli.Get(ctx, key, o) + if k8s_errors.IsNotFound(err) || !isManagedResource(o) { + return nil + } else if err != nil { + return fmt.Errorf("failed to find stale resource %s: %w", key, err) + } + err = cli.Delete(ctx, o) + if err != nil { + return fmt.Errorf("failed to delete stale resource %s: %w", key, err) + } + return nil +} + // isManagedResource returns true if the given object has a managed-by // grafana-agent-operator label. func isManagedResource(obj client.Object) bool { @@ -51,18 +70,16 @@ func generateMetricsStatefulSetService(cfg *Config, d config.Deployment) *v1.Ser d.Agent.Spec.PortName = defaultPortName } - boolTrue := true - return &v1.Service{ ObjectMeta: meta_v1.ObjectMeta{ Name: governingServiceName(d.Agent.Name), - Namespace: d.Agent.ObjectMeta.Namespace, + Namespace: d.Agent.Namespace, OwnerReferences: []meta_v1.OwnerReference{{ APIVersion: d.Agent.APIVersion, Kind: d.Agent.Kind, Name: d.Agent.Name, - BlockOwnerDeletion: &boolTrue, - Controller: &boolTrue, + BlockOwnerDeletion: pointer.Bool(true), + Controller: pointer.Bool(true), UID: d.Agent.UID, }}, Labels: cfg.Labels.Merge(map[string]string{ @@ -122,7 +139,7 @@ func generateMetricsStatefulSet( // Don't transfer any kubectl annotations to the statefulset so it doesn't // get pruned by kubectl. annotations := make(map[string]string) - for k, v := range d.Agent.ObjectMeta.Annotations { + for k, v := range d.Agent.Annotations { if !strings.HasPrefix(k, "kubectl.kubernetes.io/") { annotations[k] = v } @@ -136,8 +153,6 @@ func generateMetricsStatefulSet( labels[agentTypeLabel] = "metrics" labels[managedByOperatorLabel] = managedByOperatorLabelValue - boolTrue := true - ss := &apps_v1.StatefulSet{ ObjectMeta: meta_v1.ObjectMeta{ Name: name, @@ -147,8 +162,8 @@ func generateMetricsStatefulSet( OwnerReferences: []meta_v1.OwnerReference{{ APIVersion: d.Agent.APIVersion, Kind: d.Agent.Kind, - BlockOwnerDeletion: &boolTrue, - Controller: &boolTrue, + BlockOwnerDeletion: pointer.Bool(true), + Controller: pointer.Bool(true), Name: d.Agent.Name, UID: d.Agent.UID, }}, @@ -214,8 +229,6 @@ func generateMetricsStatefulSetSpec( shards = *reqShards } - terminationGracePeriodSeconds := int64(4800) - useVersion := d.Agent.Spec.Version if useVersion == "" { useVersion = DefaultAgentVersion @@ -446,7 +459,7 @@ func generateMetricsStatefulSetSpec( ServiceAccountName: d.Agent.Spec.ServiceAccountName, NodeSelector: d.Agent.Spec.NodeSelector, PriorityClassName: d.Agent.Spec.PriorityClassName, - TerminationGracePeriodSeconds: &terminationGracePeriodSeconds, + TerminationGracePeriodSeconds: pointer.Int64(4800), Volumes: volumes, Tolerations: d.Agent.Spec.Tolerations, Affinity: d.Agent.Spec.Affinity, diff --git a/pkg/operator/secondary_resource.go b/pkg/operator/secondary_resource.go deleted file mode 100644 index 72d0b6d2674f..000000000000 --- a/pkg/operator/secondary_resource.go +++ /dev/null @@ -1,58 +0,0 @@ -package operator - -import ( - "github.com/go-kit/log" - "k8s.io/apimachinery/pkg/types" - "sigs.k8s.io/controller-runtime/pkg/client" -) - -// secondaryResource is a secondary resource that is consumed by the primary -// resource (the GrafanaAgent CR) that should trigger a reconcile when it -// changes. -type secondaryResource int - -// List of secondary resources that the reconciler will watch. -const ( - resourcePromInstance secondaryResource = iota - resourceServiceMonitor - resourcePodMonitor - resourceProbe - resourceSecret - resourceConfigMap - resourceLogsInstance - resourcePodLogs -) - -// secondaryResources is the list of valid secondaryResources. -var secondaryResources = []secondaryResource{ - resourcePromInstance, - resourceServiceMonitor, - resourcePodMonitor, - resourceProbe, - resourceSecret, - resourceConfigMap, - resourceLogsInstance, - resourcePodLogs, -} - -// eventHandlers is a set of EnqueueRequestForSelector event handlers, one per -// secondary resource. -type eventHandlers map[secondaryResource]*enqueueRequestForSelector - -// newResourceEventHandlers creates a new eventHandlers for all secondary -// resources using the given client and logger. -func newResourceEventHandlers(c client.Reader, l log.Logger) eventHandlers { - m := make(eventHandlers) - for _, r := range secondaryResources { - m[r] = &enqueueRequestForSelector{Client: c, Log: l} - } - return m -} - -// Clear informs all event handlers to stop sending events for the given -// namespaced name. -func (ev eventHandlers) Clear(name types.NamespacedName) { - for _, v := range ev { - v.Notify(name, nil) - } -} diff --git a/pkg/operator/selector_eventhandler.go b/pkg/operator/selector_eventhandler.go deleted file mode 100644 index e799d052239f..000000000000 --- a/pkg/operator/selector_eventhandler.go +++ /dev/null @@ -1,202 +0,0 @@ -package operator - -import ( - "context" - "fmt" - "sync" - "time" - - "github.com/go-kit/log" - "github.com/go-kit/log/level" - "github.com/grafana/agent/pkg/operator/config" - v1 "k8s.io/api/core/v1" - "k8s.io/apimachinery/pkg/labels" - "k8s.io/apimachinery/pkg/types" - "k8s.io/client-go/util/workqueue" - "sigs.k8s.io/controller-runtime/pkg/client" - "sigs.k8s.io/controller-runtime/pkg/event" - "sigs.k8s.io/controller-runtime/pkg/reconcile" -) - -// enqueueRequestForSelector allows for requesting that specific -// reconciliations occur whenever an object that matches a selector -// comes in. -// -// Implements handler.EventHandler. -type enqueueRequestForSelector struct { - Client client.Reader - Log log.Logger - - mut sync.RWMutex - watchers map[types.NamespacedName][]resourceSelector -} - -// Create implements handler.EventHandler. -func (e *enqueueRequestForSelector) Create(ev event.CreateEvent, q workqueue.RateLimitingInterface) { - level.Debug(e.logger(ev.Object)).Log("msg", "got create for object") - e.handleEvent(ev.Object, q) -} - -// Update implements handler.EventHandler. -func (e *enqueueRequestForSelector) Update(ev event.UpdateEvent, q workqueue.RateLimitingInterface) { - level.Debug(e.logger(ev.ObjectNew)).Log("msg", "got update for object") - e.handleEvent(ev.ObjectOld, q) - e.handleEvent(ev.ObjectNew, q) -} - -// Delete implements handler.EventHandler. -func (e *enqueueRequestForSelector) Delete(ev event.DeleteEvent, q workqueue.RateLimitingInterface) { - level.Debug(e.logger(ev.Object)).Log("msg", "got delete for object") - e.handleEvent(ev.Object, q) -} - -// Generic implements handler.EventHandler. -func (e *enqueueRequestForSelector) Generic(ev event.GenericEvent, q workqueue.RateLimitingInterface) { - level.Debug(e.logger(ev.Object)).Log("msg", "got generic event for object") - e.handleEvent(ev.Object, q) -} - -func (e *enqueueRequestForSelector) logger(obj client.Object) log.Logger { - gvk := obj.GetObjectKind().GroupVersionKind() - return log.With( - e.Log, - "kind", fmt.Sprintf("%s/%s.%s", gvk.Group, gvk.Version, gvk.Kind), - "key", client.ObjectKeyFromObject(obj), - ) -} - -func (e *enqueueRequestForSelector) handleEvent(obj client.Object, q workqueue.RateLimitingInterface) { - e.mut.RLock() - defer e.mut.RUnlock() - - if e.watchers == nil { - return - } - - // Go through our watchers. If any of their selectors match this object, - // enqueue a reconcile request for the watcher. - for watcher, selectors := range e.watchers { - var performReconcile bool - for _, selector := range selectors { - if selector.Matches(e.Log, e.Client, obj) { - performReconcile = true - break - } - } - if performReconcile { - q.Add(reconcile.Request{NamespacedName: watcher}) - } - } -} - -// Notify will notify obj to reconcile if an event was received that matches -// any selector in ss. -// -// To stop being notified for changes, call Notify again with nil for ss. -func (e *enqueueRequestForSelector) Notify(obj types.NamespacedName, ss []resourceSelector) { - e.mut.Lock() - defer e.mut.Unlock() - - if e.watchers == nil { - e.watchers = make(map[types.NamespacedName][]resourceSelector) - } - - if ss == nil { - delete(e.watchers, obj) - } else { - e.watchers[obj] = ss - } -} - -type resourceSelector interface { - Matches(l log.Logger, c client.Reader, o client.Object) bool - SetListOptions(lo *client.ListOptions) -} - -// multiSelector returns true if all inner selectors match. -type multiSelector struct { - Selectors []resourceSelector -} - -func (s *multiSelector) Matches(l log.Logger, c client.Reader, o client.Object) bool { - for _, inner := range s.Selectors { - if !inner.Matches(l, c, o) { - return false - } - } - return true -} - -func (s *multiSelector) SetListOptions(lo *client.ListOptions) { - for _, inner := range s.Selectors { - inner.SetListOptions(lo) - } -} - -type namespaceSelector struct { - Namespace string -} - -func (s *namespaceSelector) Matches(l log.Logger, c client.Reader, o client.Object) bool { - return o.GetNamespace() == s.Namespace -} - -func (s *namespaceSelector) SetListOptions(lo *client.ListOptions) { - lo.Namespace = s.Namespace -} - -type namespaceLabelSelector struct { - Selector labels.Selector -} - -func (s *namespaceLabelSelector) Matches(l log.Logger, c client.Reader, o client.Object) bool { - ctx, cancel := context.WithTimeout(context.Background(), 500*time.Millisecond) - defer cancel() - - in := o.GetNamespace() - var ns v1.Namespace - if err := c.Get(ctx, types.NamespacedName{Name: in}, &ns); err != nil { - level.Error(l).Log("msg", "failed to look up namespace", "namespace", in, "err", err) - return false - } - - return s.Selector.Matches(labels.Set(ns.Labels)) -} - -func (s *namespaceLabelSelector) SetListOptions(lo *client.ListOptions) { - // no-op -} - -type labelSelector struct { - Selector labels.Selector -} - -func (s *labelSelector) Matches(l log.Logger, c client.Reader, o client.Object) bool { - return s.Selector.Matches(labels.Set(o.GetLabels())) -} - -func (s *labelSelector) SetListOptions(lo *client.ListOptions) { - lo.LabelSelector = s.Selector -} - -type assetReferenceSelector struct { - Reference config.AssetReference -} - -func (s *assetReferenceSelector) Matches(l log.Logger, c client.Reader, o client.Object) bool { - if o.GetNamespace() != s.Reference.Namespace { - return false - } - - if sc, ok := o.(*v1.Secret); ok { - return s.Reference.Reference.Secret != nil && sc.Name == s.Reference.Reference.Secret.Name - } else if cm, ok := o.(*v1.ConfigMap); ok { - return s.Reference.Reference.ConfigMap != nil && cm.Name == s.Reference.Reference.ConfigMap.Name - } - - return false -} - -func (s *assetReferenceSelector) SetListOptions(lo *client.ListOptions) { - // no-op -} diff --git a/pkg/operator/selector_eventhandler_test.go b/pkg/operator/selector_eventhandler_test.go deleted file mode 100644 index 6f6b3d8bfc6a..000000000000 --- a/pkg/operator/selector_eventhandler_test.go +++ /dev/null @@ -1,151 +0,0 @@ -//go:build !nonetwork && !nodocker && !race -// +build !nonetwork,!nodocker,!race - -package operator - -import ( - "context" - "testing" - "time" - - "github.com/go-kit/log" - "github.com/grafana/agent/pkg/util/k8s" - "github.com/stretchr/testify/require" - v1 "k8s.io/api/core/v1" - metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" - "k8s.io/apimachinery/pkg/labels" - "k8s.io/apimachinery/pkg/types" - "k8s.io/client-go/util/workqueue" - "sigs.k8s.io/controller-runtime/pkg/event" -) - -// TestEnqueueRequestForSelector creates an example Kubenretes cluster and runs -// EnqueueRequestForSelector to validate it works. -func TestEnqueueRequestForSelector(t *testing.T) { - l := log.NewNopLogger() - - ctx, cancel := context.WithTimeout(context.Background(), 5*time.Minute) - defer cancel() - - cluster, err := k8s.NewCluster(ctx, k8s.Options{}) - require.NoError(t, err) - defer cluster.Stop() - - cli := cluster.Client() - - // Tests will rely on a namespace existing, so let's create a namespace with - // some labels. - testNs := v1.Namespace{ - TypeMeta: metav1.TypeMeta{ - APIVersion: "v1", - Kind: "Namespace", - }, - ObjectMeta: metav1.ObjectMeta{ - Name: "enqueue-test", - Labels: map[string]string{"foo": "bar"}, - }, - } - err = cli.Create(ctx, &testNs) - require.NoError(t, err) - - testPod := &v1.Pod{ObjectMeta: metav1.ObjectMeta{ - Name: "test-pod", - Namespace: "enqueue-test", - Labels: map[string]string{"fizz": "buzz"}, - }} - - t.Run("no watchers", func(t *testing.T) { - limiter := workqueue.DefaultControllerRateLimiter() - q := workqueue.NewRateLimitingQueue(limiter) - - e := enqueueRequestForSelector{Client: cli, Log: l} - e.Create(event.CreateEvent{Object: testPod}, q) - - require.Equal(t, 0, q.Len(), "no changes should have been enqueued") - }) - - t.Run("matches watcher", func(t *testing.T) { - limiter := workqueue.DefaultControllerRateLimiter() - q := workqueue.NewRateLimitingQueue(limiter) - - e := enqueueRequestForSelector{Client: cli, Log: l} - - e.Notify(types.NamespacedName{Name: "watcher"}, buildSelectorSet( - &namespaceLabelSelector{Selector: parseSelector(t, "foo in (bar)")}, - &labelSelector{Selector: parseSelector(t, "fizz in (buzz)")}, - )) - - e.Create(event.CreateEvent{Object: testPod}, q) - require.Equal(t, 1, q.Len(), "expected one enqueue") - }) - - t.Run("matches watcher with explicit namespace", func(t *testing.T) { - limiter := workqueue.DefaultControllerRateLimiter() - q := workqueue.NewRateLimitingQueue(limiter) - - e := enqueueRequestForSelector{Client: cli, Log: l} - e.Notify(types.NamespacedName{Name: "watcher"}, buildSelectorSet( - &namespaceSelector{Namespace: "enqueue-test"}, - &namespaceLabelSelector{Selector: labels.Everything()}, - &labelSelector{Selector: labels.Everything()}, - )) - - e.Create(event.CreateEvent{Object: testPod}, q) - require.Equal(t, 1, q.Len(), "expected one enqueue") - }) - - t.Run("bad namespace name selector", func(t *testing.T) { - limiter := workqueue.DefaultControllerRateLimiter() - q := workqueue.NewRateLimitingQueue(limiter) - - e := enqueueRequestForSelector{Client: cli, Log: l} - e.Notify(types.NamespacedName{Name: "watcher"}, buildSelectorSet( - &namespaceSelector{Namespace: "default"}, - &namespaceLabelSelector{Selector: labels.Everything()}, - &labelSelector{Selector: labels.Everything()}, - )) - - e.Create(event.CreateEvent{Object: testPod}, q) - require.Equal(t, 0, q.Len(), "expected no enqueues") - }) - - t.Run("bad namespace label selector", func(t *testing.T) { - limiter := workqueue.DefaultControllerRateLimiter() - q := workqueue.NewRateLimitingQueue(limiter) - - e := enqueueRequestForSelector{Client: cli, Log: l} - e.Notify(types.NamespacedName{Name: "watcher"}, buildSelectorSet( - &namespaceLabelSelector{Selector: parseSelector(t, "foo notin (bar)")}, - &labelSelector{Selector: parseSelector(t, "fizz in (buzz)")}, - )) - - e.Create(event.CreateEvent{Object: testPod}, q) - require.Equal(t, 0, q.Len(), "no changes should have been enqueued") - }) - - t.Run("bad label selector", func(t *testing.T) { - limiter := workqueue.DefaultControllerRateLimiter() - q := workqueue.NewRateLimitingQueue(limiter) - - e := enqueueRequestForSelector{Client: cli, Log: l} - e.Notify(types.NamespacedName{Name: "watcher"}, buildSelectorSet( - &namespaceLabelSelector{Selector: parseSelector(t, "foo in (bar)")}, - &labelSelector{Selector: parseSelector(t, "fizz notin (buzz)")}, - )) - - e.Create(event.CreateEvent{Object: testPod}, q) - require.Equal(t, 0, q.Len(), "no changes should have been enqueued") - }) -} - -// buildSelectorSet returns a single multiSelector composed of ss. -func buildSelectorSet(ss ...resourceSelector) []resourceSelector { - return []resourceSelector{&multiSelector{Selectors: ss}} -} - -func parseSelector(t *testing.T, selector string) labels.Selector { - t.Helper() - s, err := labels.Parse(selector) - require.NoError(t, err) - return s -} diff --git a/pkg/operator/testdata/test-resource-hierarchy.yaml b/pkg/operator/testdata/test-resource-hierarchy.yaml new file mode 100644 index 000000000000..c3e412550f7c --- /dev/null +++ b/pkg/operator/testdata/test-resource-hierarchy.yaml @@ -0,0 +1,165 @@ +apiVersion: monitoring.grafana.com/v1alpha1 +kind: GrafanaAgent +metadata: + name: grafana-agent-example + namespace: default + labels: + app: grafana-agent-example +spec: + image: grafana/agent:latest + serviceAccountName: grafana-agent + logs: + instanceSelector: + matchLabels: + agent: grafana-agent-example + metrics: + instanceSelector: + matchLabels: + agent: grafana-agent-example + +--- + +apiVersion: monitoring.grafana.com/v1alpha1 +kind: MetricsInstance +metadata: + name: primary + namespace: default + labels: + agent: grafana-agent-example +spec: + remoteWrite: + - url: http://prometheus:9090/api/v1/write + basicAuth: + username: + name: prometheus-fake-credentials + key: fakeUsername + password: + name: prometheus-fake-credentials + key: fakePassword + # Supply an empty namespace selector to look in all namespaces. + podMonitorNamespaceSelector: {} + podMonitorSelector: + matchLabels: + instance: primary + # Supply an empty namespace selector to look in all namespaces. + serviceMonitorNamespaceSelector: {} + serviceMonitorSelector: + matchLabels: + instance: primary + +--- + +apiVersion: monitoring.grafana.com/v1alpha1 +kind: LogsInstance +metadata: + name: primary + namespace: default + labels: + agent: grafana-agent-example +spec: + clients: + - url: http://loki:8080/loki/api/v1/push + + # Supply an empty namespace selector to look in all namespaces. + podLogsNamespaceSelector: {} + podLogsSelector: + matchLabels: + instance: primary + +--- + +# Have the Agent monitor itself. +apiVersion: monitoring.coreos.com/v1 +kind: PodMonitor +metadata: + name: grafana-agents + namespace: default + labels: + instance: primary +spec: + selector: + matchLabels: + app.kubernetes.io/name: grafana-agent + podMetricsEndpoints: + - port: http-metrics + +--- + +# Have the Agent get logs from itself. +apiVersion: monitoring.grafana.com/v1alpha1 +kind: PodLogs +metadata: + name: grafana-agents + namespace: default + labels: + instance: primary +spec: + selector: + matchLabels: + app.kubernetes.io/name: grafana-agent + pipelineStages: + - cri: {} + +# +# Pretend credentials +# + +--- +apiVersion: v1 +kind: Secret +metadata: + name: prometheus-fake-credentials + namespace: default +data: + # "user" + fakeUsername: "dXNlcg==" + # "password" + fakePassword: "cGFzc3dvcmQ=" + +# +# Extra resources +# + +--- +apiVersion: v1 +kind: ServiceAccount +metadata: + name: grafana-agent + namespace: default +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRole +metadata: + name: grafana-agent +rules: +- apiGroups: + - "" + resources: + - nodes + - nodes/proxy + - nodes/metrics + - services + - endpoints + - pods + verbs: + - get + - list + - watch +- nonResourceURLs: + - /metrics + - /metrics/cadvisor + verbs: + - get +--- +apiVersion: rbac.authorization.k8s.io/v1 +kind: ClusterRoleBinding +metadata: + name: grafana-agent +roleRef: + apiGroup: rbac.authorization.k8s.io + kind: ClusterRole + name: grafana-agent +subjects: +- kind: ServiceAccount + name: grafana-agent + namespace: default