From 3d5e52dc3d05a23303412b2ecbdfecfdf429aac3 Mon Sep 17 00:00:00 2001 From: Soam Vasani Date: Tue, 1 Nov 2016 18:22:04 -0700 Subject: [PATCH] Router integration with poolmgr and controller Router uses poolmgr to specialize pods when necessary. Router uses controller to get the list of triggers to listen for. For now this integration is pretty crappy -- we just poll the controller every few seconds and cache the result. The right way would be to have some sort of watch API on the controller and use that. Or maybe share access to etcd directly. --- router/functionHandler.go | 27 ++++++++++++----- router/httpTriggers.go | 63 ++++++++++++++++++++++++++++----------- router/mutablemux.go | 2 +- router/router.go | 13 ++++++-- router/router_test.go | 2 +- 5 files changed, 79 insertions(+), 28 deletions(-) diff --git a/router/functionHandler.go b/router/functionHandler.go index 37f2cb6c..dd85a393 100644 --- a/router/functionHandler.go +++ b/router/functionHandler.go @@ -17,29 +17,40 @@ limitations under the License. package router import ( - "errors" + "fmt" "log" "net/http" "net/http/httputil" "net/url" "github.com/platform9/fission" + poolmgrClient "github.com/platform9/fission/poolmgr/client" ) type functionHandler struct { - fmap *functionServiceMap - poolManagerUrl string - Function fission.Metadata + fmap *functionServiceMap + poolmgr *poolmgrClient.Client + Function fission.Metadata } -func (*functionHandler) getServiceForFunction() (*url.URL, error) { - return nil, errors.New("not implemented") +func (fh *functionHandler) getServiceForFunction() (*url.URL, error) { + // call poolmgr, get a url for a function + svcName, err := fh.poolmgr.GetServiceForFunction(&fh.Function) + if err != nil { + return nil, err + } + svcUrl, err := url.Parse(fmt.Sprintf("http://%v", svcName)) + if err != nil { + return nil, err + } + return svcUrl, nil } func (fh *functionHandler) handler(responseWriter http.ResponseWriter, request *http.Request) { serviceUrl, err := fh.fmap.lookup(&fh.Function) if err != nil { // Cache miss: request the Pool Manager to make a new service. + log.Printf("Not cached, getting new service for %v", fh.Function) serviceUrl, poolErr := fh.getServiceForFunction() if poolErr != nil { log.Printf("Failed to get service for function (%v,%v): %v", @@ -55,9 +66,11 @@ func (fh *functionHandler) handler(responseWriter http.ResponseWriter, request * } // Proxy off our request to the serviceUrl, and send the response back. - // TODO: As an optimization we may want to cache proxies too -- this would get us + // TODO: As an optimization we may want to cache proxies too -- this might get us // connection reuse and possibly better performance director := func(req *http.Request) { + log.Printf("Proxying request for %v", req.URL) + // send this request to serviceurl req.URL.Scheme = serviceUrl.Scheme req.URL.Host = serviceUrl.Host diff --git a/router/httpTriggers.go b/router/httpTriggers.go index 5731ab5f..2ee95c9e 100644 --- a/router/httpTriggers.go +++ b/router/httpTriggers.go @@ -17,47 +17,76 @@ limitations under the License. package router import ( + "log" + "time" + "github.com/gorilla/mux" + "github.com/platform9/fission" + controllerClient "github.com/platform9/fission/controller/client" + poolmgrClient "github.com/platform9/fission/poolmgr/client" ) type HTTPTriggerSet struct { *functionServiceMap *mutableRouter - controllerUrl string - poolManagerUrl string - triggers []fission.HTTPTrigger + controller *controllerClient.Client + poolmgr *poolmgrClient.Client + triggers []fission.HTTPTrigger } -func makeHTTPTriggerSet(fmap *functionServiceMap, controllerUrl string, poolManagerUrl string) *HTTPTriggerSet { +func makeHTTPTriggerSet(fmap *functionServiceMap, controller *controllerClient.Client, poolmgr *poolmgrClient.Client) *HTTPTriggerSet { triggers := make([]fission.HTTPTrigger, 1) return &HTTPTriggerSet{ functionServiceMap: fmap, triggers: triggers, - controllerUrl: controllerUrl, - poolManagerUrl: poolManagerUrl, + controller: controller, + poolmgr: poolmgr, } } -func (triggers *HTTPTriggerSet) subscribeRouter(mr *mutableRouter) { - triggers.mutableRouter = mr - mr.updateRouter(triggers.getRouterFromTriggers()) - go triggers.watchTriggers() +func (ts *HTTPTriggerSet) subscribeRouter(mr *mutableRouter) { + ts.mutableRouter = mr + mr.updateRouter(ts.getRouterFromTriggers()) + go ts.watchTriggers() } -func (triggers *HTTPTriggerSet) getRouterFromTriggers() *mux.Router { +func (ts *HTTPTriggerSet) getRouterFromTriggers() *mux.Router { muxRouter := mux.NewRouter() - for _, trigger := range triggers.triggers { + for _, trigger := range ts.triggers { fh := &functionHandler{ - fmap: triggers.functionServiceMap, - Function: trigger.Function, - poolManagerUrl: triggers.poolManagerUrl, + fmap: ts.functionServiceMap, + Function: trigger.Function, + poolmgr: ts.poolmgr, } muxRouter.HandleFunc(trigger.UrlPattern, fh.handler) } return muxRouter } -func (triggers *HTTPTriggerSet) watchTriggers() { - // watch controller for updates to triggers and update the router accordingly +func (ts *HTTPTriggerSet) watchTriggers() { + if ts.controller == nil { + return + } + + failureCount := 0 + maxFailures := 5 + + // Watch controller for updates to triggers and update the router accordingly. + // TODO change this to use a watch API; or maybe even watch etcd directly. + for { + triggers, err := ts.controller.HTTPTriggerList() + if err != nil { + failureCount += 1 + if failureCount >= maxFailures { + log.Fatalf("Failed to connect to controller after %v retries: %v", failureCount, err) + } + } + log.Printf("Updating router, %v triggers", len(triggers)) + + ts.triggers = triggers + ts.mutableRouter.updateRouter(ts.getRouterFromTriggers()) + + time.Sleep(3 * time.Second) + } } diff --git a/router/mutablemux.go b/router/mutablemux.go index 0f1a4b5c..6e0354c2 100644 --- a/router/mutablemux.go +++ b/router/mutablemux.go @@ -49,6 +49,6 @@ func (mr *mutableRouter) ServeHTTP(responseWriter http.ResponseWriter, request * } func (mr *mutableRouter) updateRouter(newHandler *mux.Router) { - log.Print("Updating router") + log.Println("Updating router") mr.router.Store(newHandler) } diff --git a/router/router.go b/router/router.go index 41037ef0..90e8dd05 100644 --- a/router/router.go +++ b/router/router.go @@ -41,8 +41,13 @@ package router import ( "fmt" - "github.com/gorilla/mux" "net/http" + + "github.com/gorilla/mux" + + controllerClient "github.com/platform9/fission/controller/client" + poolmgrClient "github.com/platform9/fission/poolmgr/client" + "log" ) // request url ---[mux]---> Function(name,uid) ----[fmap]----> k8s service url @@ -64,6 +69,10 @@ func serve(port int, httpTriggerSet *HTTPTriggerSet) { func Start(port int, controllerUrl string, poolmgrUrl string) { fmap := makeFunctionServiceMap() - triggers := makeHTTPTriggerSet(fmap, controllerUrl, poolmgrUrl) + controller := controllerClient.MakeClient(controllerUrl) + poolmgr := poolmgrClient.MakeClient(poolmgrUrl) + + triggers := makeHTTPTriggerSet(fmap, controller, poolmgr) + log.Printf("Starting router at port %v\n", port) serve(port, triggers) } diff --git a/router/router_test.go b/router/router_test.go index db6c8c59..521c6cd4 100644 --- a/router/router_test.go +++ b/router/router_test.go @@ -33,7 +33,7 @@ func TestRouter(t *testing.T) { fmap.assign(fn, testServiceUrl) - triggers := makeHTTPTriggerSet(fmap, "", "") + triggers := makeHTTPTriggerSet(fmap, nil, nil) triggerUrl := "/foo" triggers.triggers = append(triggers.triggers, fission.HTTPTrigger{UrlPattern: triggerUrl, Function: *fn})