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})