Skip to content
Closed
Show file tree
Hide file tree
Changes from all 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
Original file line number Diff line number Diff line change
@@ -0,0 +1,60 @@
package openshiftmanager

import (
"k8s.io/apimachinery/pkg/runtime/schema"

"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/event"
)

type inputResourceEventFilter func(obj client.Object) bool

// inputResourceDispatcher is a simple dispatcher that applies GVK scoped filters
// and forwards matching operator.
//
// Each GVK has its own set of filters. Today these
// may include name/namespace checks, and in the future label selectors.
//
// Longer term, this dispatcher is expected to track which input resources are
// associated with which operator.
type inputResourceDispatcher struct {
// eventsCh channel on which an operator name
// to reconcile will be sent
eventsCh chan event.TypedGenericEvent[string]
filters map[schema.GroupVersionKind][]inputResourceEventFilter
}

func newInputResourceDispatcher() *inputResourceDispatcher {
return &inputResourceDispatcher{
eventsCh: make(chan event.TypedGenericEvent[string]),
}
}

func (d *inputResourceDispatcher) SetFilters(filters map[schema.GroupVersionKind][]inputResourceEventFilter) {
d.filters = filters
}

func (d *inputResourceDispatcher) Handle(gvk schema.GroupVersionKind, cObj client.Object) {
filters := d.filters[gvk]
if len(filters) == 0 {
// for the POC we always return cao
// TODO: implement proper operator discovery
d.eventsCh <- event.TypedGenericEvent[string]{Object: "cluster-authentication-operator"}
return
}

for _, filter := range filters {
if filter(cObj) {
// for the POC we always return cao
// TODO: implement proper operator discovery
d.eventsCh <- event.TypedGenericEvent[string]{Object: "cluster-authentication-operator"}
return
}
}
}

// ResultChan returns a channel on which
// an operator name to reconcile will be sent
func (d *inputResourceDispatcher) ResultChan() <-chan event.TypedGenericEvent[string] {
return d.eventsCh
}
Original file line number Diff line number Diff line change
@@ -0,0 +1,127 @@
package openshiftmanager

import (
"testing"
"time"

corev1 "k8s.io/api/core/v1"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/runtime/schema"

"sigs.k8s.io/controller-runtime/pkg/client"
"sigs.k8s.io/controller-runtime/pkg/event"

"github.com/stretchr/testify/require"
)

func TestInputResourceDispatcherHandle(t *testing.T) {
wellKnownGVK := schema.GroupVersionKind{Group: "example.io", Version: "v1", Kind: "Widget"}
obj := &corev1.ConfigMap{
ObjectMeta: metav1.ObjectMeta{
Name: "widget-a",
Namespace: "default",
},
}

scenarios := []struct {
name string
filters map[schema.GroupVersionKind][]inputResourceEventFilter
inputGVK schema.GroupVersionKind
inputObj client.Object
expectedOperatorNames []event.TypedGenericEvent[string]
}{
{
name: "dispatches matching filter",
filters: map[schema.GroupVersionKind][]inputResourceEventFilter{
wellKnownGVK: {
func(cObj client.Object) bool {
return cObj.GetName() == "widget-a"
},
},
},
inputGVK: wellKnownGVK,
inputObj: obj,
expectedOperatorNames: []event.TypedGenericEvent[string]{
{Object: "cluster-authentication-operator"},
},
},
{
name: "does not dispatch when filters do not match",
filters: map[schema.GroupVersionKind][]inputResourceEventFilter{
wellKnownGVK: {
func(cObj client.Object) bool {
return cObj.GetName() == "widget-b"
},
},
},
inputGVK: wellKnownGVK,
inputObj: obj,
},
{
name: "dispatches when gvk has no filters",
filters: map[schema.GroupVersionKind][]inputResourceEventFilter{},
inputGVK: wellKnownGVK,
inputObj: obj,
expectedOperatorNames: []event.TypedGenericEvent[string]{
{Object: "cluster-authentication-operator"},
},
},
{
name: "dispatches when any filter matches",
filters: map[schema.GroupVersionKind][]inputResourceEventFilter{
wellKnownGVK: {
func(cObj client.Object) bool {
return cObj.GetName() == "widget-b"
},
func(cObj client.Object) bool {
return cObj.GetNamespace() == "default"
},
},
},
inputGVK: wellKnownGVK,
inputObj: obj,
expectedOperatorNames: []event.TypedGenericEvent[string]{
{Object: "cluster-authentication-operator"},
},
},
}

for _, scenario := range scenarios {
t.Run(scenario.name, func(t *testing.T) {
dispatcher := newInputResourceDispatcher()
dispatcher.SetFilters(scenario.filters)
// dispatch in a goroutine for simplicity with an unbuffered channel
go dispatcher.Handle(scenario.inputGVK, scenario.inputObj)

events := readEvents(t, dispatcher.ResultChan(), len(scenario.expectedOperatorNames))
require.Equal(t, scenario.expectedOperatorNames, events)
ensureNoMoreEvents(t, dispatcher.ResultChan())
})
}
}

func readEvents(t *testing.T, ch <-chan event.TypedGenericEvent[string], expected int) []event.TypedGenericEvent[string] {
if expected == 0 {
return nil
}

events := make([]event.TypedGenericEvent[string], 0, expected)
for i := 0; i < expected; i++ {
select {
case evt := <-ch:
events = append(events, evt)
case <-time.After(100 * time.Millisecond):
require.Failf(t, "expected event not received", "received %d/%d events", len(events), expected)
}
}

return events
}

func ensureNoMoreEvents(t *testing.T, ch <-chan event.TypedGenericEvent[string]) {
select {
case ev := <-ch:
require.Failf(t, "unexpected event received", "got %+v", ev)
case <-time.After(100 * time.Millisecond):
}
}
Loading