From da07b35a968a310418566aca5d7fae2262f374ca Mon Sep 17 00:00:00 2001 From: Ta-Ching Chen Date: Thu, 9 Nov 2017 09:07:14 -0600 Subject: [PATCH] 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. --- router/functionReferenceResolver.go | 22 +------ router/httpTriggers.go | 89 +++++++++++++---------------- router/router.go | 21 ++++--- router/router_test.go | 7 ++- 4 files changed, 58 insertions(+), 81 deletions(-) diff --git a/router/functionReferenceResolver.go b/router/functionReferenceResolver.go index feff53c9..e98a11c1 100644 --- a/router/functionReferenceResolver.go +++ b/router/functionReferenceResolver.go @@ -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{ diff --git a/router/httpTriggers.go b/router/httpTriggers.go index f7c379ae..9fc72c92 100644 --- a/router/httpTriggers.go +++ b/router/httpTriggers.go @@ -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() { diff --git a/router/router.go b/router/router.go index 735ac42d..4496f4b9 100644 --- a/router/router.go +++ b/router/router.go @@ -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) } diff --git a/router/router_test.go b/router/router_test.go index 01e37b0f..9e631400 100644 --- a/router/router_test.go +++ b/router/router_test.go @@ -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