-
Notifications
You must be signed in to change notification settings - Fork 1.1k
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
move Source definitions to a single folder without breaking the API
Signed-off-by: Tim Ramlot <42113979+inteon@users.noreply.github.com>
- Loading branch information
Showing
9 changed files
with
292 additions
and
169 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,150 @@ | ||
/* | ||
Copyright 2018 The Kubernetes Authors. | ||
Licensed under the Apache License, Version 2.0 (the "License"); | ||
you may not use this file except in compliance with the License. | ||
You may obtain a copy of the License at | ||
http://www.apache.org/licenses/LICENSE-2.0 | ||
Unless required by applicable law or agreed to in writing, software | ||
distributed under the License is distributed on an "AS IS" BASIS, | ||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
See the License for the specific language governing permissions and | ||
limitations under the License. | ||
*/ | ||
|
||
package internal | ||
|
||
import ( | ||
"context" | ||
"fmt" | ||
"sync" | ||
|
||
"k8s.io/client-go/util/workqueue" | ||
"sigs.k8s.io/controller-runtime/pkg/event" | ||
"sigs.k8s.io/controller-runtime/pkg/handler" | ||
"sigs.k8s.io/controller-runtime/pkg/predicate" | ||
) | ||
|
||
const ( | ||
// defaultBufferSize is the default number of event notifications that can be buffered. | ||
defaultBufferSize = 1024 | ||
) | ||
|
||
// Channel is used to provide a source of events originating outside the cluster | ||
// (e.g. GitHub Webhook callback). Channel requires the user to wire the external | ||
// source (eh.g. http handler) to write GenericEvents to the underlying channel. | ||
type Channel struct { | ||
// once ensures the event distribution goroutine will be performed only once | ||
once sync.Once | ||
|
||
// Source is the source channel to fetch GenericEvents | ||
Source <-chan event.GenericEvent | ||
|
||
// dest is the destination channels of the added event handlers | ||
dest []chan event.GenericEvent | ||
|
||
// DestBufferSize is the specified buffer size of dest channels. | ||
// Default to 1024 if not specified. | ||
DestBufferSize int | ||
|
||
// destLock is to ensure the destination channels are safely added/removed | ||
destLock sync.Mutex | ||
} | ||
|
||
func (cs *Channel) String() string { | ||
return fmt.Sprintf("channel source: %p", cs) | ||
} | ||
|
||
// Start implements Source and should only be called by the Controller. | ||
func (cs *Channel) Start( | ||
ctx context.Context, | ||
handler handler.EventHandler, | ||
queue workqueue.RateLimitingInterface, | ||
prct ...predicate.Predicate) error { | ||
// Source should have been specified by the user. | ||
if cs.Source == nil { | ||
return fmt.Errorf("must specify Channel.Source") | ||
} | ||
|
||
// use default value if DestBufferSize not specified | ||
if cs.DestBufferSize == 0 { | ||
cs.DestBufferSize = defaultBufferSize | ||
} | ||
|
||
dst := make(chan event.GenericEvent, cs.DestBufferSize) | ||
|
||
cs.destLock.Lock() | ||
cs.dest = append(cs.dest, dst) | ||
cs.destLock.Unlock() | ||
|
||
cs.once.Do(func() { | ||
// Distribute GenericEvents to all EventHandler / Queue pairs Watching this source | ||
go cs.syncLoop(ctx) | ||
}) | ||
|
||
go func() { | ||
for evt := range dst { | ||
shouldHandle := true | ||
for _, p := range prct { | ||
if !p.Generic(evt) { | ||
shouldHandle = false | ||
break | ||
} | ||
} | ||
|
||
if shouldHandle { | ||
func() { | ||
ctx, cancel := context.WithCancel(ctx) | ||
defer cancel() | ||
handler.Generic(ctx, evt, queue) | ||
}() | ||
} | ||
} | ||
}() | ||
|
||
return nil | ||
} | ||
|
||
func (cs *Channel) doStop() { | ||
cs.destLock.Lock() | ||
defer cs.destLock.Unlock() | ||
|
||
for _, dst := range cs.dest { | ||
close(dst) | ||
} | ||
} | ||
|
||
func (cs *Channel) distribute(evt event.GenericEvent) { | ||
cs.destLock.Lock() | ||
defer cs.destLock.Unlock() | ||
|
||
for _, dst := range cs.dest { | ||
// We cannot make it under goroutine here, or we'll meet the | ||
// race condition of writing message to closed channels. | ||
// To avoid blocking, the dest channels are expected to be of | ||
// proper buffer size. If we still see it blocked, then | ||
// the controller is thought to be in an abnormal state. | ||
dst <- evt | ||
} | ||
} | ||
|
||
func (cs *Channel) syncLoop(ctx context.Context) { | ||
for { | ||
select { | ||
case <-ctx.Done(): | ||
// Close destination channels | ||
cs.doStop() | ||
return | ||
case evt, stillOpen := <-cs.Source: | ||
if !stillOpen { | ||
// if the source channel is closed, we're never gonna get | ||
// anything more on it, so stop & bail | ||
cs.doStop() | ||
return | ||
} | ||
cs.distribute(evt) | ||
} | ||
} | ||
} |
File renamed without changes.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,39 @@ | ||
/* | ||
Copyright 2018 The Kubernetes Authors. | ||
Licensed under the Apache License, Version 2.0 (the "License"); | ||
you may not use this file except in compliance with the License. | ||
You may obtain a copy of the License at | ||
http://www.apache.org/licenses/LICENSE-2.0 | ||
Unless required by applicable law or agreed to in writing, software | ||
distributed under the License is distributed on an "AS IS" BASIS, | ||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
See the License for the specific language governing permissions and | ||
limitations under the License. | ||
*/ | ||
|
||
package internal | ||
|
||
import ( | ||
"context" | ||
"fmt" | ||
|
||
"k8s.io/client-go/util/workqueue" | ||
"sigs.k8s.io/controller-runtime/pkg/handler" | ||
"sigs.k8s.io/controller-runtime/pkg/predicate" | ||
) | ||
|
||
// Func is a function that implements Source. | ||
type Func func(context.Context, handler.EventHandler, workqueue.RateLimitingInterface, ...predicate.Predicate) error | ||
|
||
// Start implements Source. | ||
func (f Func) Start(ctx context.Context, evt handler.EventHandler, queue workqueue.RateLimitingInterface, | ||
pr ...predicate.Predicate) error { | ||
return f(ctx, evt, queue, pr...) | ||
} | ||
|
||
func (f Func) String() string { | ||
return fmt.Sprintf("func source: %p", f) | ||
} |
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,53 @@ | ||
/* | ||
Copyright 2018 The Kubernetes Authors. | ||
Licensed under the Apache License, Version 2.0 (the "License"); | ||
you may not use this file except in compliance with the License. | ||
You may obtain a copy of the License at | ||
http://www.apache.org/licenses/LICENSE-2.0 | ||
Unless required by applicable law or agreed to in writing, software | ||
distributed under the License is distributed on an "AS IS" BASIS, | ||
WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. | ||
See the License for the specific language governing permissions and | ||
limitations under the License. | ||
*/ | ||
|
||
package internal | ||
|
||
import ( | ||
"context" | ||
"fmt" | ||
|
||
"k8s.io/client-go/util/workqueue" | ||
"sigs.k8s.io/controller-runtime/pkg/cache" | ||
"sigs.k8s.io/controller-runtime/pkg/handler" | ||
"sigs.k8s.io/controller-runtime/pkg/predicate" | ||
) | ||
|
||
// Informer is used to provide a source of events originating inside the cluster from Watches (e.g. Pod Create). | ||
type Informer struct { | ||
// Informer is the controller-runtime Informer | ||
Informer cache.Informer | ||
} | ||
|
||
// Start is internal and should be called only by the Controller to register an EventHandler with the Informer | ||
// to enqueue reconcile.Requests. | ||
func (is *Informer) Start(ctx context.Context, handler handler.EventHandler, queue workqueue.RateLimitingInterface, | ||
prct ...predicate.Predicate) error { | ||
// Informer should have been specified by the user. | ||
if is.Informer == nil { | ||
return fmt.Errorf("must specify Informer.Informer") | ||
} | ||
|
||
_, err := is.Informer.AddEventHandler(NewEventHandler(ctx, queue, handler, prct).HandlerFuncs()) | ||
if err != nil { | ||
return err | ||
} | ||
return nil | ||
} | ||
|
||
func (is *Informer) String() string { | ||
return fmt.Sprintf("informer source: %p", is.Informer) | ||
} |
File renamed without changes.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Oops, something went wrong.