diff --git a/router/httpTriggers.go b/router/httpTriggers.go index 51bba411..0999ed50 100644 --- a/router/httpTriggers.go +++ b/router/httpTriggers.go @@ -38,17 +38,18 @@ type HTTPTriggerSet struct { *functionServiceMap *mutableRouter - fissionClient *crd.FissionClient - kubeClient *kubernetes.Clientset - executor *executorClient.Client - resolver *functionReferenceResolver - crdClient *rest.RESTClient - triggers []crd.HTTPTrigger - triggerStore k8sCache.Store - triggerController k8sCache.Controller - functions []crd.Function - funcStore k8sCache.Store - funcController k8sCache.Controller + fissionClient *crd.FissionClient + kubeClient *kubernetes.Clientset + executor *executorClient.Client + resolver *functionReferenceResolver + crdClient *rest.RESTClient + triggers []crd.HTTPTrigger + triggerStore k8sCache.Store + triggerController k8sCache.Controller + functions []crd.Function + funcStore k8sCache.Store + funcController k8sCache.Controller + updateRouterRequestChannel chan struct{} tsRoundTripperParams *tsRoundTripperParams } @@ -56,13 +57,14 @@ type HTTPTriggerSet struct { func makeHTTPTriggerSet(fmap *functionServiceMap, fissionClient *crd.FissionClient, kubeClient *kubernetes.Clientset, executor *executorClient.Client, crdClient *rest.RESTClient, params *tsRoundTripperParams) (*HTTPTriggerSet, k8sCache.Store, k8sCache.Store) { httpTriggerSet := &HTTPTriggerSet{ - functionServiceMap: fmap, - triggers: []crd.HTTPTrigger{}, - fissionClient: fissionClient, - kubeClient: kubeClient, - executor: executor, - crdClient: crdClient, - tsRoundTripperParams: params, + functionServiceMap: fmap, + triggers: []crd.HTTPTrigger{}, + fissionClient: fissionClient, + kubeClient: kubeClient, + executor: executor, + crdClient: crdClient, + updateRouterRequestChannel: make(chan struct{}), + tsRoundTripperParams: params, } var tStore, fnStore k8sCache.Store var tController, fnController k8sCache.Controller @@ -87,6 +89,7 @@ func (ts *HTTPTriggerSet) subscribeRouter(ctx context.Context, mr *mutableRouter log.Printf("Skipping continuous trigger updates") return } + go ts.updateRouter() go ts.runWatcher(ctx, ts.funcController) go ts.runWatcher(ctx, ts.triggerController) } @@ -192,6 +195,11 @@ func (ts *HTTPTriggerSet) initTriggerController() (k8sCache.Store, k8sCache.Cont UpdateFunc: func(oldObj interface{}, newObj interface{}) { oldTrigger := oldObj.(*crd.HTTPTrigger) newTrigger := newObj.(*crd.HTTPTrigger) + + if oldTrigger.Metadata.ResourceVersion == newTrigger.Metadata.ResourceVersion { + return + } + go updateIngress(oldTrigger, newTrigger, ts.kubeClient) ts.syncTriggers() }, @@ -211,7 +219,13 @@ func (ts *HTTPTriggerSet) initFunctionController() (k8sCache.Store, k8sCache.Con ts.syncTriggers() }, UpdateFunc: func(oldObj interface{}, newObj interface{}) { + oldFn := oldObj.(*crd.Function) fn := newObj.(*crd.Function) + + if oldFn.Metadata.ResourceVersion == fn.Metadata.ResourceVersion { + return + } + // update resolver function reference cache for key, rr := range ts.resolver.copy() { if key.functionReference.Name == fn.Metadata.Name && @@ -236,22 +250,28 @@ func (ts *HTTPTriggerSet) runWatcher(ctx context.Context, controller k8sCache.Co } func (ts *HTTPTriggerSet) syncTriggers() { - // get triggers - latestTriggers := ts.triggerStore.List() - triggers := make([]crd.HTTPTrigger, len(latestTriggers)) - for _, t := range latestTriggers { - triggers = append(triggers, *t.(*crd.HTTPTrigger)) - } - ts.triggers = triggers - - // get functions - latestFunctions := ts.funcStore.List() - functions := make([]crd.Function, len(latestFunctions)) - for _, f := range latestFunctions { - functions = append(functions, *f.(*crd.Function)) - } - ts.functions = functions - - // make a new router and use it - ts.mutableRouter.updateRouter(ts.getRouter()) + ts.updateRouterRequestChannel <- struct{}{} +} + +func (ts *HTTPTriggerSet) updateRouter() { + for range ts.updateRouterRequestChannel { + // get triggers + latestTriggers := ts.triggerStore.List() + triggers := make([]crd.HTTPTrigger, len(latestTriggers)) + for _, t := range latestTriggers { + triggers = append(triggers, *t.(*crd.HTTPTrigger)) + } + ts.triggers = triggers + + // get functions + latestFunctions := ts.funcStore.List() + functions := make([]crd.Function, len(latestFunctions)) + for _, f := range latestFunctions { + functions = append(functions, *f.(*crd.Function)) + } + ts.functions = functions + + // make a new router and use it + ts.mutableRouter.updateRouter(ts.getRouter()) + } }