Use common controller/store for httpTriggerSet and functionReferenceResolver (#390)
This fixes a bug where functionReferenceResolver returned out-of-date function metadata and caused the router to proxy requests to old function pods. It also uses the go context package to shutdown the controller when the router is shutting down.
This commit is contained in:
committed by
Soam Vasani
parent
9694ba4fd1
commit
da07b35a96
@@ -36,8 +36,6 @@ type (
|
||||
// functionReferenceResolver provides a resolver to turn a function
|
||||
// reference into a resolveResult
|
||||
functionReferenceResolver struct {
|
||||
fissionClient *crd.FissionClient
|
||||
|
||||
// FunctionReference -> function metadata
|
||||
refCache *cache.Cache
|
||||
|
||||
@@ -68,28 +66,14 @@ const (
|
||||
resolveResultSingleFunction = iota
|
||||
)
|
||||
|
||||
func makeFunctionReferenceResolver(fissionClient *crd.FissionClient) *functionReferenceResolver {
|
||||
func makeFunctionReferenceResolver(store k8sCache.Store) *functionReferenceResolver {
|
||||
frr := &functionReferenceResolver{
|
||||
fissionClient: fissionClient,
|
||||
refCache: cache.MakeCache(time.Minute, 0),
|
||||
refCache: cache.MakeCache(time.Minute, 0),
|
||||
store: store,
|
||||
}
|
||||
return frr
|
||||
}
|
||||
|
||||
// Sync starts syncing crd function resources from k8s api server
|
||||
func (frr *functionReferenceResolver) Sync(crdClient *rest.RESTClient) {
|
||||
stopCh := make(chan struct{})
|
||||
store, controller := makeK8SCache(crdClient)
|
||||
frr.stopCh = stopCh
|
||||
frr.store = store
|
||||
go controller.Run(stopCh)
|
||||
}
|
||||
|
||||
// Stop stops crd resources syncing
|
||||
func (frr *functionReferenceResolver) Stop() {
|
||||
frr.stopCh <- struct{}{}
|
||||
}
|
||||
|
||||
func makeK8SCache(crdClient *rest.RESTClient) (k8sCache.Store, k8sCache.Controller) {
|
||||
watchlist := k8sCache.NewListWatchFromClient(crdClient, "functions", metav1.NamespaceDefault, fields.Everything())
|
||||
listWatch := &k8sCache.ListWatch{
|
||||
|
||||
+40
-49
@@ -17,6 +17,7 @@ limitations under the License.
|
||||
package router
|
||||
|
||||
import (
|
||||
"context"
|
||||
"log"
|
||||
"net/http"
|
||||
"time"
|
||||
@@ -24,8 +25,6 @@ import (
|
||||
"github.com/gorilla/mux"
|
||||
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
|
||||
"k8s.io/apimachinery/pkg/fields"
|
||||
"k8s.io/apimachinery/pkg/runtime"
|
||||
"k8s.io/apimachinery/pkg/watch"
|
||||
"k8s.io/client-go/rest"
|
||||
k8sCache "k8s.io/client-go/tools/cache"
|
||||
|
||||
@@ -37,30 +36,42 @@ import (
|
||||
type HTTPTriggerSet struct {
|
||||
*functionServiceMap
|
||||
*mutableRouter
|
||||
fissionClient *crd.FissionClient
|
||||
poolmgr *poolmgrClient.Client
|
||||
resolver *functionReferenceResolver
|
||||
triggers []crd.HTTPTrigger
|
||||
triggerStore k8sCache.Store
|
||||
functions []crd.Function
|
||||
funcStore k8sCache.Store
|
||||
crdClient *rest.RESTClient
|
||||
fissionClient *crd.FissionClient
|
||||
poolmgr *poolmgrClient.Client
|
||||
resolver *functionReferenceResolver
|
||||
crdClient *rest.RESTClient
|
||||
triggers []crd.HTTPTrigger
|
||||
triggerStore k8sCache.Store
|
||||
triggerController k8sCache.Controller
|
||||
functions []crd.Function
|
||||
funcStore k8sCache.Store
|
||||
funcController k8sCache.Controller
|
||||
}
|
||||
|
||||
func makeHTTPTriggerSet(fmap *functionServiceMap, fissionClient *crd.FissionClient,
|
||||
poolmgr *poolmgrClient.Client, resolver *functionReferenceResolver, crdClient *rest.RESTClient) *HTTPTriggerSet {
|
||||
triggers := make([]crd.HTTPTrigger, 1)
|
||||
return &HTTPTriggerSet{
|
||||
poolmgr *poolmgrClient.Client, crdClient *rest.RESTClient) (*HTTPTriggerSet, k8sCache.Store, k8sCache.Store) {
|
||||
httpTriggerSet := &HTTPTriggerSet{
|
||||
functionServiceMap: fmap,
|
||||
triggers: triggers,
|
||||
triggers: []crd.HTTPTrigger{},
|
||||
fissionClient: fissionClient,
|
||||
poolmgr: poolmgr,
|
||||
resolver: resolver,
|
||||
crdClient: crdClient,
|
||||
}
|
||||
var tStore, fnStore k8sCache.Store
|
||||
var tController, fnController k8sCache.Controller
|
||||
if httpTriggerSet.crdClient != nil {
|
||||
tStore, tController = httpTriggerSet.initTriggerController()
|
||||
httpTriggerSet.triggerStore = tStore
|
||||
httpTriggerSet.triggerController = tController
|
||||
fnStore, fnController = httpTriggerSet.initFunctionController()
|
||||
httpTriggerSet.funcStore = fnStore
|
||||
httpTriggerSet.funcController = fnController
|
||||
}
|
||||
return httpTriggerSet, tStore, fnStore
|
||||
}
|
||||
|
||||
func (ts *HTTPTriggerSet) subscribeRouter(mr *mutableRouter) {
|
||||
func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter, resolver *functionReferenceResolver) {
|
||||
ts.resolver = resolver
|
||||
ts.mutableRouter = mr
|
||||
mr.updateRouter(ts.getRouter())
|
||||
|
||||
@@ -69,8 +80,8 @@ func (ts *HTTPTriggerSet) subscribeRouter(mr *mutableRouter) {
|
||||
log.Printf("Skipping continuous trigger updates")
|
||||
return
|
||||
}
|
||||
go ts.watchTriggers()
|
||||
go ts.watchFunctions()
|
||||
go ts.runWatcher(ctx, ts.funcController)
|
||||
go ts.runWatcher(ctx, ts.triggerController)
|
||||
}
|
||||
|
||||
func defaultHomeHandler(w http.ResponseWriter, r *http.Request) {
|
||||
@@ -140,17 +151,9 @@ func (ts *HTTPTriggerSet) updateTriggerStatusFailed(ht *crd.HTTPTrigger, err err
|
||||
// TODO
|
||||
}
|
||||
|
||||
func (ts *HTTPTriggerSet) watchTriggers() {
|
||||
watchlist := k8sCache.NewListWatchFromClient(ts.crdClient, "httptriggers", metav1.NamespaceDefault, fields.Everything())
|
||||
listWatch := &k8sCache.ListWatch{
|
||||
ListFunc: func(options metav1.ListOptions) (runtime.Object, error) {
|
||||
return watchlist.List(options)
|
||||
},
|
||||
WatchFunc: func(options metav1.ListOptions) (watch.Interface, error) {
|
||||
return watchlist.Watch(options)
|
||||
},
|
||||
}
|
||||
func (ts *HTTPTriggerSet) initTriggerController() (k8sCache.Store, k8sCache.Controller) {
|
||||
resyncPeriod := 30 * time.Second
|
||||
listWatch := k8sCache.NewListWatchFromClient(ts.crdClient, "httptriggers", metav1.NamespaceDefault, fields.Everything())
|
||||
store, controller := k8sCache.NewInformer(listWatch, &crd.HTTPTrigger{}, resyncPeriod,
|
||||
k8sCache.ResourceEventHandlerFuncs{
|
||||
AddFunc: func(obj interface{}) {
|
||||
@@ -163,25 +166,12 @@ func (ts *HTTPTriggerSet) watchTriggers() {
|
||||
ts.syncTriggers()
|
||||
},
|
||||
})
|
||||
ts.triggerStore = store
|
||||
stop := make(chan struct{})
|
||||
defer func() {
|
||||
stop <- struct{}{}
|
||||
}()
|
||||
controller.Run(stop)
|
||||
return store, controller
|
||||
}
|
||||
|
||||
func (ts *HTTPTriggerSet) watchFunctions() {
|
||||
watchlist := k8sCache.NewListWatchFromClient(ts.crdClient, "functions", metav1.NamespaceDefault, fields.Everything())
|
||||
listWatch := &k8sCache.ListWatch{
|
||||
ListFunc: func(options metav1.ListOptions) (runtime.Object, error) {
|
||||
return watchlist.List(options)
|
||||
},
|
||||
WatchFunc: func(options metav1.ListOptions) (watch.Interface, error) {
|
||||
return watchlist.Watch(options)
|
||||
},
|
||||
}
|
||||
func (ts *HTTPTriggerSet) initFunctionController() (k8sCache.Store, k8sCache.Controller) {
|
||||
resyncPeriod := 30 * time.Second
|
||||
listWatch := k8sCache.NewListWatchFromClient(ts.crdClient, "functions", metav1.NamespaceDefault, fields.Everything())
|
||||
store, controller := k8sCache.NewInformer(listWatch, &crd.Function{}, resyncPeriod,
|
||||
k8sCache.ResourceEventHandlerFuncs{
|
||||
AddFunc: func(obj interface{}) {
|
||||
@@ -206,12 +196,13 @@ func (ts *HTTPTriggerSet) watchFunctions() {
|
||||
ts.syncTriggers()
|
||||
},
|
||||
})
|
||||
ts.funcStore = store
|
||||
stop := make(chan struct{})
|
||||
defer func() {
|
||||
stop <- struct{}{}
|
||||
return store, controller
|
||||
}
|
||||
|
||||
func (ts *HTTPTriggerSet) runWatcher(ctx context.Context, controller k8sCache.Controller) {
|
||||
go func() {
|
||||
controller.Run(ctx.Done())
|
||||
}()
|
||||
controller.Run(stop)
|
||||
}
|
||||
|
||||
func (ts *HTTPTriggerSet) syncTriggers() {
|
||||
|
||||
+10
-11
@@ -40,6 +40,7 @@ Its job is to:
|
||||
package router
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log"
|
||||
"net/http"
|
||||
@@ -57,15 +58,15 @@ import (
|
||||
|
||||
// request url ---[trigger]---> Function(name, deployment) ----[deployment]----> Function(name, uid) ----[pool mgr]---> k8s service url
|
||||
|
||||
func router(httpTriggerSet *HTTPTriggerSet) *mutableRouter {
|
||||
func router(ctx context.Context, httpTriggerSet *HTTPTriggerSet, resolver *functionReferenceResolver) *mutableRouter {
|
||||
muxRouter := mux.NewRouter()
|
||||
mr := NewMutableRouter(muxRouter)
|
||||
httpTriggerSet.subscribeRouter(mr)
|
||||
httpTriggerSet.subscribeRouter(ctx, mr, resolver)
|
||||
return mr
|
||||
}
|
||||
|
||||
func serve(port int, httpTriggerSet *HTTPTriggerSet) {
|
||||
mr := router(httpTriggerSet)
|
||||
func serve(ctx context.Context, port int, httpTriggerSet *HTTPTriggerSet, resolver *functionReferenceResolver) {
|
||||
mr := router(ctx, httpTriggerSet, resolver)
|
||||
url := fmt.Sprintf(":%v", port)
|
||||
http.ListenAndServe(url, handlers.LoggingHandler(os.Stdout, mr))
|
||||
}
|
||||
@@ -79,12 +80,10 @@ func Start(port int, poolmgrUrl string) {
|
||||
}
|
||||
restClient := fissionClient.GetCrdClient()
|
||||
poolmgr := poolmgrClient.MakeClient(poolmgrUrl)
|
||||
resolver := makeFunctionReferenceResolver(fissionClient)
|
||||
resolver.Sync(restClient)
|
||||
defer func() {
|
||||
resolver.Stop()
|
||||
}()
|
||||
triggers := makeHTTPTriggerSet(fmap, fissionClient, poolmgr, resolver, restClient)
|
||||
triggers, _, fnStore := makeHTTPTriggerSet(fmap, fissionClient, poolmgr, restClient)
|
||||
resolver := makeFunctionReferenceResolver(fnStore)
|
||||
log.Printf("Starting router at port %v\n", port)
|
||||
serve(port, triggers)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
serve(ctx, port, triggers, resolver)
|
||||
}
|
||||
|
||||
@@ -17,6 +17,7 @@ limitations under the License.
|
||||
package router
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"testing"
|
||||
"time"
|
||||
@@ -58,7 +59,7 @@ func TestRouter(t *testing.T) {
|
||||
frr.refCache.Set(nfr, rr)
|
||||
|
||||
// HTTP trigger set with a trigger for this function
|
||||
triggers := makeHTTPTriggerSet(fmap, nil, nil, frr, nil)
|
||||
triggers, _, _ := makeHTTPTriggerSet(fmap, nil, nil, nil)
|
||||
triggerUrl := "/foo"
|
||||
triggers.triggers = append(triggers.triggers,
|
||||
crd.HTTPTrigger{
|
||||
@@ -75,7 +76,9 @@ func TestRouter(t *testing.T) {
|
||||
|
||||
// run the router
|
||||
port := 4242
|
||||
go serve(port, triggers)
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
go serve(ctx, port, triggers, frr)
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
|
||||
// hit the router
|
||||
|
||||
Reference in New Issue
Block a user