From 380b7d1bcd7dad9141e4a24682d876902eac2e59 Mon Sep 17 00:00:00 2001 From: Ta-Ching Chen Date: Mon, 6 Nov 2017 14:58:19 -0600 Subject: [PATCH] Use k8s client store to sync functions and triggers for fast synchronization (#382) --- router/httpTriggers.go | 31 ++++++++++++++++--------------- 1 file changed, 16 insertions(+), 15 deletions(-) diff --git a/router/httpTriggers.go b/router/httpTriggers.go index eb65fca3..f7c379ae 100644 --- a/router/httpTriggers.go +++ b/router/httpTriggers.go @@ -41,7 +41,9 @@ type HTTPTriggerSet struct { poolmgr *poolmgrClient.Client resolver *functionReferenceResolver triggers []crd.HTTPTrigger + triggerStore k8sCache.Store functions []crd.Function + funcStore k8sCache.Store crdClient *rest.RESTClient } @@ -139,9 +141,6 @@ func (ts *HTTPTriggerSet) updateTriggerStatusFailed(ht *crd.HTTPTrigger, err err } func (ts *HTTPTriggerSet) watchTriggers() { - // sync all http triggers - ts.syncTriggers() - watchlist := k8sCache.NewListWatchFromClient(ts.crdClient, "httptriggers", metav1.NamespaceDefault, fields.Everything()) listWatch := &k8sCache.ListWatch{ ListFunc: func(options metav1.ListOptions) (runtime.Object, error) { @@ -152,7 +151,7 @@ func (ts *HTTPTriggerSet) watchTriggers() { }, } resyncPeriod := 30 * time.Second - _, controller := k8sCache.NewInformer(listWatch, &crd.HTTPTrigger{}, resyncPeriod, + store, controller := k8sCache.NewInformer(listWatch, &crd.HTTPTrigger{}, resyncPeriod, k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { ts.syncTriggers() @@ -164,6 +163,7 @@ func (ts *HTTPTriggerSet) watchTriggers() { ts.syncTriggers() }, }) + ts.triggerStore = store stop := make(chan struct{}) defer func() { stop <- struct{}{} @@ -172,8 +172,6 @@ func (ts *HTTPTriggerSet) watchTriggers() { } func (ts *HTTPTriggerSet) watchFunctions() { - ts.syncTriggers() - watchlist := k8sCache.NewListWatchFromClient(ts.crdClient, "functions", metav1.NamespaceDefault, fields.Everything()) listWatch := &k8sCache.ListWatch{ ListFunc: func(options metav1.ListOptions) (runtime.Object, error) { @@ -184,7 +182,7 @@ func (ts *HTTPTriggerSet) watchFunctions() { }, } resyncPeriod := 30 * time.Second - _, controller := k8sCache.NewInformer(listWatch, &crd.Function{}, resyncPeriod, + store, controller := k8sCache.NewInformer(listWatch, &crd.Function{}, resyncPeriod, k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { ts.syncTriggers() @@ -208,6 +206,7 @@ func (ts *HTTPTriggerSet) watchFunctions() { ts.syncTriggers() }, }) + ts.funcStore = store stop := make(chan struct{}) defer func() { stop <- struct{}{} @@ -219,18 +218,20 @@ func (ts *HTTPTriggerSet) syncTriggers() { log.Printf("Syncing http triggers") // get triggers - triggers, err := ts.fissionClient.HTTPTriggers(metav1.NamespaceAll).List(metav1.ListOptions{}) - if err != nil { - log.Fatalf("Failed to get http trigger list: %v", err) + latestTriggers := ts.triggerStore.List() + triggers := make([]crd.HTTPTrigger, len(latestTriggers)) + for _, t := range latestTriggers { + triggers = append(triggers, *t.(*crd.HTTPTrigger)) } - ts.triggers = triggers.Items + ts.triggers = triggers // get functions - functions, err := ts.fissionClient.Functions(metav1.NamespaceAll).List(metav1.ListOptions{}) - if err != nil { - log.Fatalf("Failed to get function list: %v", err) + latestFunctions := ts.funcStore.List() + functions := make([]crd.Function, len(latestFunctions)) + for _, f := range latestFunctions { + functions = append(functions, *f.(*crd.Function)) } - ts.functions = functions.Items + ts.functions = functions // make a new router and use it ts.mutableRouter.updateRouter(ts.getRouter())