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