Skip to content

Commit

Permalink
Broker eventtype autocreate fixes (#7161)
Browse files Browse the repository at this point in the history
* Fixed undefined typemeta on broker in eventtype autocreate

Signed-off-by: Calum Murray <cmurray@redhat.com>

* Fixed autocreate so that only one eventtype is created when events are sent to mt channel broker

Signed-off-by: Calum Murray <cmurray@redhat.com>

* Clean up

Signed-off-by: Calum Murray <cmurray@redhat.com>

* Fixed unit tests

Signed-off-by: Calum Murray <cmurray@redhat.com>

* channel only creates eventtypes if not owned by broker

Signed-off-by: Calum Murray <cmurray@redhat.com>

---------

Signed-off-by: Calum Murray <cmurray@redhat.com>
  • Loading branch information
Cali0707 authored Aug 24, 2023
1 parent 7749771 commit 0045fa9
Show file tree
Hide file tree
Showing 3 changed files with 34 additions and 8 deletions.
18 changes: 16 additions & 2 deletions cmd/broker/ingress/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -131,12 +131,26 @@ func main() {
logger.Fatal("Error setting up trace publishing", zap.Error(err))
}

featureStore := feature.NewStore(logging.FromContext(ctx).Named("feature-config-store"))
var featureStore *feature.Store
var handler *ingress.Handler

featureStore = feature.NewStore(logging.FromContext(ctx).Named("feature-config-store"), func(name string, value interface{}) {
featureFlags := value.(feature.Flags)
if featureFlags.IsEnabled(feature.EvenTypeAutoCreate) && featureStore != nil && handler != nil {
autoCreate := &eventtype.EventTypeAutoHandler{
EventTypeLister: eventtypeinformer.Get(ctx).Lister(),
EventingClient: eventingclient.Get(ctx).EventingV1beta2(),
FeatureStore: featureStore,
Logger: logger,
}
handler.EvenTypeHandler = autoCreate
}
})
featureStore.WatchConfigs(configMapWatcher)

reporter := ingress.NewStatsReporter(env.ContainerName, kmeta.ChildName(env.PodName, uuid.New().String()))

handler, err := ingress.NewHandler(logger, reporter, broker.TTLDefaulter(logger, int32(env.MaxTTL)), brokerInformer)
handler, err = ingress.NewHandler(logger, reporter, broker.TTLDefaulter(logger, int32(env.MaxTTL)), brokerInformer)
if err != nil {
logger.Fatal("Error creating Handler", zap.Error(err))
}
Expand Down
9 changes: 8 additions & 1 deletion pkg/broker/ingress/ingress_handler.go
Original file line number Diff line number Diff line change
Expand Up @@ -247,13 +247,20 @@ func (h *Handler) ServeHTTP(writer http.ResponseWriter, request *http.Request) {
}

func toKReference(broker *eventingv1.Broker) *duckv1.KReference {
return &duckv1.KReference{
kref := &duckv1.KReference{
Kind: broker.Kind,
Namespace: broker.Namespace,
Name: broker.Name,
APIVersion: broker.APIVersion,
Address: broker.Status.Address.Name,
}
if kref.Kind == "" {
kref.Kind = "Broker"
}
if kref.APIVersion == "" {
kref.APIVersion = "eventing.knative.dev/v1"
}
return kref
}

func (h *Handler) receive(ctx context.Context, headers http.Header, event *cloudevents.Event, brokerNamespace, brokerName string) (int, time.Duration) {
Expand Down
15 changes: 10 additions & 5 deletions pkg/reconciler/inmemorychannel/dispatcher/inmemorychannel.go
Original file line number Diff line number Diff line change
Expand Up @@ -93,17 +93,22 @@ func (r *Reconciler) reconcile(ctx context.Context, imc *v1.InMemoryChannel) rec
return err
}
var eventTypeAutoHandler *eventtype.EventTypeAutoHandler
var channelReference *duckv1.KReference
var channelRef *duckv1.KReference
var UID *types.UID
if r.featureStore.IsEnabled(feature.EvenTypeAutoCreate) {
if ownerReferences := imc.GetOwnerReferences(); r.featureStore.IsEnabled(feature.EvenTypeAutoCreate) &&
(len(ownerReferences) == 0 ||
ownerReferences[0].Kind != "Broker") {
logging.FromContext(ctx).Info("EventType autocreate is enabled, creating handler")
eventTypeAutoHandler = &eventtype.EventTypeAutoHandler{
EventTypeLister: r.eventTypeLister,
EventingClient: r.eventingClient,
FeatureStore: r.featureStore,
Logger: logging.FromContext(ctx).Desugar(),
}
channelReference = toKReference(imc)

channelRef = toKReference(imc)
UID = &imc.UID

}

// First grab the host based MultiChannelFanoutMessage httpHandler
Expand All @@ -115,7 +120,7 @@ func (r *Reconciler) reconcile(ctx context.Context, imc *v1.InMemoryChannel) rec
config.FanoutConfig,
r.reporter,
eventTypeAutoHandler,
channelReference,
channelRef,
UID,
)
if err != nil {
Expand Down Expand Up @@ -143,7 +148,7 @@ func (r *Reconciler) reconcile(ctx context.Context, imc *v1.InMemoryChannel) rec
config.FanoutConfig,
r.reporter,
eventTypeAutoHandler,
channelReference,
channelRef,
UID,
channel.ResolveChannelFromPath(channel.ParseChannelFromPath),
)
Expand Down

0 comments on commit 0045fa9

Please sign in to comment.