Skip to content
Merged
Show file tree
Hide file tree
Changes from 12 commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
30 changes: 15 additions & 15 deletions pkg/controller/config/config_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -29,6 +29,8 @@ import (
"github.com/open-policy-agent/gatekeeper/v3/pkg/metrics"
"github.com/open-policy-agent/gatekeeper/v3/pkg/mutation"
"github.com/open-policy-agent/gatekeeper/v3/pkg/readiness"
syncutil "github.com/open-policy-agent/gatekeeper/v3/pkg/syncutil"
cm "github.com/open-policy-agent/gatekeeper/v3/pkg/syncutil/cachemanager"
"github.com/open-policy-agent/gatekeeper/v3/pkg/target"
"github.com/open-policy-agent/gatekeeper/v3/pkg/watch"
"k8s.io/apimachinery/pkg/api/errors"
Expand Down Expand Up @@ -109,16 +111,14 @@ func (a *Adder) InjectWatchSet(watchSet *watch.Set) {
// events is the channel from which sync controller will receive the events
// regEvents is the channel registered by Registrar to put the events in
// events and regEvents point to same event channel except for testing.
func newReconciler(mgr manager.Manager, opa syncc.OpaDataClient, wm *watch.Manager, cs *watch.ControllerSwitch, tracker *readiness.Tracker, processExcluder *process.Excluder, events <-chan event.GenericEvent, watchSet *watch.Set, regEvents chan<- event.GenericEvent) (*ReconcileConfig, error) {
filteredOpa := syncc.NewFilteredOpaDataClient(opa, watchSet)
syncMetricsCache := syncc.NewMetricsCache()
func newReconciler(mgr manager.Manager, opa syncutil.OpaDataClient, wm *watch.Manager, cs *watch.ControllerSwitch, tracker *readiness.Tracker, processExcluder *process.Excluder, events <-chan event.GenericEvent, watchSet *watch.Set, regEvents chan<- event.GenericEvent) (*ReconcileConfig, error) {
filteredOpa := syncutil.NewFilteredOpaDataClient(opa, watchSet)
syncMetricsCache := syncutil.NewMetricsCache()
cm := cm.NewCacheManager(opa, syncMetricsCache, tracker, processExcluder)

syncAdder := syncc.Adder{
Opa: filteredOpa,
Events: events,
MetricsCache: syncMetricsCache,
Tracker: tracker,
ProcessExcluder: processExcluder,
Events: events,
CacheManager: cm,
}
// Create subordinate controller - we will feed it events dynamically via watch
if err := syncAdder.Add(mgr); err != nil {
Expand Down Expand Up @@ -176,8 +176,8 @@ type ReconcileConfig struct {
statusClient client.StatusClient

scheme *runtime.Scheme
opa syncc.OpaDataClient
syncMetricsCache *syncc.MetricsCache
opa syncutil.OpaDataClient
syncMetricsCache *syncutil.MetricsCache
cs *watch.ControllerSwitch
watcher *watch.Registrar

Expand Down Expand Up @@ -327,7 +327,7 @@ func (r *ReconcileConfig) wipeCacheIfNeeded(ctx context.Context) error {

// reset sync cache before sending the metric
r.syncMetricsCache.ResetCache()
r.syncMetricsCache.ReportSync(&syncc.Reporter{})
r.syncMetricsCache.ReportSync(&syncutil.Reporter{}, log)

r.needsWipe = false
}
Expand All @@ -352,10 +352,10 @@ func (r *ReconcileConfig) replayData(ctx context.Context) error {
return fmt.Errorf("replaying data for %+v: %w", gvk, err)
}

defer r.syncMetricsCache.ReportSync(&syncc.Reporter{})
defer r.syncMetricsCache.ReportSync(&syncutil.Reporter{}, log)

for i := range u.Items {
syncKey := r.syncMetricsCache.GetSyncKey(u.Items[i].GetNamespace(), u.Items[i].GetName())
syncKey := syncutil.GetKeyForSyncMetrics(u.Items[i].GetNamespace(), u.Items[i].GetName())

isExcludedNamespace, err := r.skipExcludedNamespace(&u.Items[i])
if err != nil {
Expand All @@ -367,14 +367,14 @@ func (r *ReconcileConfig) replayData(ctx context.Context) error {
}

if _, err := r.opa.AddData(ctx, &u.Items[i]); err != nil {
r.syncMetricsCache.AddObject(syncKey, syncc.Tags{
r.syncMetricsCache.AddObject(syncKey, syncutil.Tags{
Kind: u.Items[i].GetKind(),
Status: metrics.ErrorStatus,
})
return fmt.Errorf("adding data for %+v: %w", gvk, err)
}

r.syncMetricsCache.AddObject(syncKey, syncc.Tags{
r.syncMetricsCache.AddObject(syncKey, syncutil.Tags{
Kind: u.Items[i].GetKind(),
Status: metrics.ActiveStatus,
})
Expand Down
206 changes: 34 additions & 172 deletions pkg/controller/sync/sync_controller.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,16 +17,13 @@ package sync

import (
"context"
"strings"
"sync"
"time"

"github.com/go-logr/logr"
"github.com/open-policy-agent/gatekeeper/v3/pkg/controller/config/process"
"github.com/open-policy-agent/gatekeeper/v3/pkg/logging"
"github.com/open-policy-agent/gatekeeper/v3/pkg/metrics"
"github.com/open-policy-agent/gatekeeper/v3/pkg/operations"
"github.com/open-policy-agent/gatekeeper/v3/pkg/readiness"
"github.com/open-policy-agent/gatekeeper/v3/pkg/syncutil"
cm "github.com/open-policy-agent/gatekeeper/v3/pkg/syncutil/cachemanager"
"github.com/open-policy-agent/gatekeeper/v3/pkg/util"
"k8s.io/apimachinery/pkg/api/errors"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
Expand All @@ -44,11 +41,8 @@ import (
var log = logf.Log.WithName("controller").WithValues("metaKind", "Sync")

type Adder struct {
Opa OpaDataClient
Events <-chan event.GenericEvent
MetricsCache *MetricsCache
Tracker *readiness.Tracker
ProcessExcluder *process.Excluder
CacheManager *cm.CacheManager
Events <-chan event.GenericEvent
}

// Add creates a new Sync Controller and adds it to the Manager with default RBAC. The Manager will set fields on the Controller
Expand All @@ -57,34 +51,28 @@ func (a *Adder) Add(mgr manager.Manager) error {
if !operations.HasValidationOperations() {
return nil
}
reporter, err := NewStatsReporter()
reporter, err := syncutil.NewStatsReporter()
if err != nil {
log.Error(err, "Sync metrics reporter could not start")
return err
}

r := newReconciler(mgr, a.Opa, *reporter, a.MetricsCache, a.Tracker, a.ProcessExcluder)
r := newReconciler(mgr, *reporter, a.CacheManager)
return add(mgr, r, a.Events)
}

// newReconciler returns a new reconcile.Reconciler.
func newReconciler(
mgr manager.Manager,
opa OpaDataClient,
reporter Reporter,
metricsCache *MetricsCache,
tracker *readiness.Tracker,
processExcluder *process.Excluder,
reporter syncutil.Reporter,
cmt *cm.CacheManager,
) reconcile.Reconciler {
return &ReconcileSync{
reader: mgr.GetCache(),
scheme: mgr.GetScheme(),
opa: opa,
log: log,
reporter: reporter,
metricsCache: metricsCache,
tracker: tracker,
processExcluder: processExcluder,
reader: mgr.GetCache(),
scheme: mgr.GetScheme(),
log: log,
reporter: reporter,
cm: cmt,
}
}

Expand All @@ -108,28 +96,14 @@ func add(mgr manager.Manager, r reconcile.Reconciler, events <-chan event.Generi

var _ reconcile.Reconciler = &ReconcileSync{}

type MetricsCache struct {
mux sync.RWMutex
Cache map[string]Tags
KnownKinds map[string]bool
}

type Tags struct {
Kind string
Status metrics.Status
}

// ReconcileSync reconciles an arbitrary object described by Kind.
type ReconcileSync struct {
reader client.Reader

scheme *runtime.Scheme
opa OpaDataClient
log logr.Logger
reporter Reporter
metricsCache *MetricsCache
tracker *readiness.Tracker
processExcluder *process.Excluder
scheme *runtime.Scheme
log logr.Logger
reporter syncutil.Reporter
cm *cm.CacheManager
}

// +kubebuilder:rbac:groups=constraints.gatekeeper.sh,resources=*,verbs=get;list;watch;create;update;patch;delete
Expand All @@ -147,17 +121,19 @@ func (r *ReconcileSync) Reconcile(ctx context.Context, request reconcile.Request
return reconcile.Result{}, nil
}

syncKey := r.metricsCache.GetSyncKey(unpackedRequest.Namespace, unpackedRequest.Name)
reportMetrics := false
// todo acpana -- double check that request namespace & name match instance namespace & name
// syncKey := syncutil.GetKeyForSyncMetrics(unpackedRequest.Namespace, unpackedRequest.Name)
Comment thread
acpana marked this conversation as resolved.
Outdated

reportMetricsForRenconcileRun := false
Comment thread
acpana marked this conversation as resolved.
Outdated
defer func() {
if reportMetrics {
if err := r.reporter.reportSyncDuration(time.Since(timeStart)); err != nil {
if reportMetricsForRenconcileRun {
if err := r.reporter.ReportSyncDuration(time.Since(timeStart)); err != nil {
log.Error(err, "failed to report sync duration")
}

r.metricsCache.ReportSync(&r.reporter)
r.cm.ReportSyncMetrics(&r.reporter, log)

if err := r.reporter.reportLastSync(); err != nil {
if err := r.reporter.ReportLastSync(); err != nil {
log.Error(err, "failed to report last sync timestamp")
}
}
Expand All @@ -171,46 +147,25 @@ func (r *ReconcileSync) Reconcile(ctx context.Context, request reconcile.Request
// This is a deletion; remove the data
instance.SetNamespace(unpackedRequest.Namespace)
instance.SetName(unpackedRequest.Name)
if _, err := r.opa.RemoveData(ctx, instance); err != nil {
if err := r.cm.RemoveObject(ctx, instance); err != nil {
return reconcile.Result{}, err
}

// cancel expectations
t := r.tracker.ForData(instance.GroupVersionKind())
t.CancelExpect(instance)

r.metricsCache.DeleteObject(syncKey)
reportMetrics = true
reportMetricsForRenconcileRun = true
return reconcile.Result{}, nil
}
// Error reading the object - requeue the request.
return reconcile.Result{}, err
}

// namespace is excluded from sync
isExcludedNamespace, err := r.skipExcludedNamespace(instance)
if err != nil {
log.Error(err, "error while excluding namespaces")
}

if isExcludedNamespace {
// cancel expectations
t := r.tracker.ForData(instance.GroupVersionKind())
t.CancelExpect(instance)
return reconcile.Result{}, nil
}

// todo acpana -- double check that it is okay to remove what has been
// namespace excluded now (but was not namesapced excluded before)
Comment thread
acpana marked this conversation as resolved.
Outdated
if !instance.GetDeletionTimestamp().IsZero() {
if _, err := r.opa.RemoveData(ctx, instance); err != nil {
if err := r.cm.RemoveObject(ctx, instance); err != nil {
return reconcile.Result{}, err
}

// cancel expectations
t := r.tracker.ForData(instance.GroupVersionKind())
t.CancelExpect(instance)

r.metricsCache.DeleteObject(syncKey)
reportMetrics = true
reportMetricsForRenconcileRun = true
return reconcile.Result{}, nil
}

Expand All @@ -222,106 +177,13 @@ func (r *ReconcileSync) Reconcile(ctx context.Context, request reconcile.Request
logging.ResourceName, instance.GetName(),
)

if _, err := r.opa.AddData(ctx, instance); err != nil {
r.metricsCache.AddObject(syncKey, Tags{
Kind: instance.GetKind(),
Status: metrics.ErrorStatus,
})
reportMetrics = true
if err := r.cm.AddObject(ctx, instance); err != nil {
reportMetricsForRenconcileRun = true

return reconcile.Result{}, err
}
r.tracker.ForData(gvk).Observe(instance)
log.V(1).Info("[readiness] observed data", "gvk", gvk, "namespace", instance.GetNamespace(), "name", instance.GetName())

r.metricsCache.AddObject(syncKey, Tags{
Kind: instance.GetKind(),
Status: metrics.ActiveStatus,
})

r.metricsCache.addKind(instance.GetKind())

reportMetrics = true
reportMetricsForRenconcileRun = true

return reconcile.Result{}, nil
}

func (r *ReconcileSync) skipExcludedNamespace(obj *unstructured.Unstructured) (bool, error) {
isNamespaceExcluded, err := r.processExcluder.IsNamespaceExcluded(process.Sync, obj)
if err != nil {
return false, err
}

return isNamespaceExcluded, err
}

func NewMetricsCache() *MetricsCache {
return &MetricsCache{
Cache: make(map[string]Tags),
KnownKinds: make(map[string]bool),
}
}

func (c *MetricsCache) GetSyncKey(namespace string, name string) string {
return strings.Join([]string{namespace, name}, "/")
}

// need to know encountered kinds to reset metrics for that kind
// this is a known memory leak
// footprint should naturally reset on Pod upgrade b/c the container restarts.
func (c *MetricsCache) addKind(key string) {
c.mux.Lock()
defer c.mux.Unlock()

c.KnownKinds[key] = true
}

func (c *MetricsCache) ResetCache() {
c.mux.Lock()
defer c.mux.Unlock()

c.Cache = make(map[string]Tags)
}

func (c *MetricsCache) AddObject(key string, t Tags) {
c.mux.Lock()
defer c.mux.Unlock()

c.Cache[key] = Tags{
Kind: t.Kind,
Status: t.Status,
}
}

func (c *MetricsCache) DeleteObject(key string) {
c.mux.Lock()
defer c.mux.Unlock()

delete(c.Cache, key)
}

func (c *MetricsCache) ReportSync(reporter *Reporter) {
c.mux.RLock()
defer c.mux.RUnlock()

totals := make(map[Tags]int)
for _, v := range c.Cache {
totals[v]++
}

for kind := range c.KnownKinds {
for _, status := range metrics.AllStatuses {
if err := reporter.reportSync(
Tags{
Kind: kind,
Status: status,
},
int64(totals[Tags{
Kind: kind,
Status: status,
}])); err != nil {
log.Error(err, "failed to report sync")
}
}
}
}
Loading