diff --git a/charts/fission-all/templates/deployment.yaml b/charts/fission-all/templates/deployment.yaml index 9836eae2..f1d57f42 100644 --- a/charts/fission-all/templates/deployment.yaml +++ b/charts/fission-all/templates/deployment.yaml @@ -167,6 +167,11 @@ spec: imagePullPolicy: {{ .Values.pullPolicy }} command: ["/fission-bundle"] args: ["--routerPort", "8888", "--executorUrl", "http://executor.{{ .Release.Namespace }}"] + env: + - name: POD_NAMESPACE + valueFrom: + fieldRef: + fieldPath: metadata.namespace readinessProbe: httpGet: path: "/router-healthz" diff --git a/charts/fission-core/templates/deployment.yaml b/charts/fission-core/templates/deployment.yaml index 5b0a70ee..b66b72f7 100644 --- a/charts/fission-core/templates/deployment.yaml +++ b/charts/fission-core/templates/deployment.yaml @@ -163,6 +163,11 @@ spec: imagePullPolicy: {{ .Values.pullPolicy }} command: ["/fission-bundle"] args: ["--routerPort", "8888", "--executorUrl", "http://executor.{{ .Release.Namespace }}"] + env: + - name: POD_NAMESPACE + valueFrom: + fieldRef: + fieldPath: metadata.namespace readinessProbe: httpGet: path: "/router-healthz" diff --git a/fission/httptrigger.go b/fission/httptrigger.go index 3e56c0ce..820d9a30 100644 --- a/fission/httptrigger.go +++ b/fission/httptrigger.go @@ -93,6 +93,12 @@ func htCreate(c *cli.Context) error { } checkFunctionExistence(client, fnName, fnNamespace) + createIngress := false + if c.IsSet("createingress") { + createIngress = c.Bool("createingress") + } + + host := c.String("host") // just name triggers by uuid. triggerName := uuid.NewV4().String() @@ -103,12 +109,14 @@ func htCreate(c *cli.Context) error { Namespace: fnNamespace, }, Spec: fission.HTTPTriggerSpec{ + Host: host, RelativeURL: triggerUrl, Method: getMethod(method), FunctionReference: fission.FunctionReference{ Type: fission.FunctionReferenceTypeFunctionName, Name: fnName, }, + CreateIngress: createIngress, }, } @@ -157,6 +165,14 @@ func htUpdate(c *cli.Context) error { ht.Spec.FunctionReference.Name = newFn } + if c.IsSet("createingress") { + ht.Spec.CreateIngress = c.Bool("createingress") + } + + if c.IsSet("host") { + ht.Spec.Host = c.String("host") + } + _, err = client.HTTPTriggerUpdate(ht) checkErr(err, "update HTTP trigger") @@ -191,10 +207,10 @@ func htList(c *cli.Context) error { w := tabwriter.NewWriter(os.Stdout, 0, 0, 1, ' ', 0) - fmt.Fprintf(w, "%v\t%v\t%v\t%v\t%v\n", "NAME", "METHOD", "HOST", "URL", "FUNCTION_NAME") + fmt.Fprintf(w, "%v\t%v\t%v\t%v\t%v\t%v\n", "NAME", "METHOD", "HOST", "URL", "INGRESS", "FUNCTION_NAME") for _, ht := range hts { - fmt.Fprintf(w, "%v\t%v\t%v\t%v\t%v\n", - ht.Metadata.Name, ht.Spec.Method, ht.Spec.Host, ht.Spec.RelativeURL, ht.Spec.FunctionReference.Name) + fmt.Fprintf(w, "%v\t%v\t%v\t%v\t%v\t%v\n", + ht.Metadata.Name, ht.Spec.Method, ht.Spec.Host, ht.Spec.RelativeURL, ht.Spec.CreateIngress, ht.Spec.FunctionReference.Name) } w.Flush() diff --git a/fission/main.go b/fission/main.go index 1e1d3524..7fca5401 100644 --- a/fission/main.go +++ b/fission/main.go @@ -164,10 +164,13 @@ func main() { // httptriggers htNameFlag := cli.StringFlag{Name: "name", Usage: "HTTP Trigger name"} htFnNameFlag := cli.StringFlag{Name: "function", Usage: "Function name"} + htHostFlag := cli.StringFlag{Name: "host", Usage: "FQDN of the network host for route"} + htIngressFlag := cli.BoolFlag{Name: "createingress", Usage: "Creates ingress with same URL, defaults to false"} htSubcommands := []cli.Command{ - {Name: "create", Aliases: []string{"add"}, Usage: "Create HTTP trigger", Flags: []cli.Flag{htMethodFlag, htUrlFlag, htFnNameFlag, fnNameFlag, fnNamespaceFlag, specSaveFlag}, Action: htCreate}, + + {Name: "create", Aliases: []string{"add"}, Usage: "Create HTTP trigger", Flags: []cli.Flag{htMethodFlag, htUrlFlag, htFnNameFlag, htHostFlag, htIngressFlag, fnNamespaceFlag, specSaveFlag}, Action: htCreate}, {Name: "get", Usage: "Get HTTP trigger", Flags: []cli.Flag{htMethodFlag, htUrlFlag}, Action: htGet}, - {Name: "update", Usage: "Update HTTP trigger", Flags: []cli.Flag{htNameFlag, triggerNamespaceFlag, htFnNameFlag}, Action: htUpdate}, + {Name: "update", Usage: "Update HTTP trigger", Flags: []cli.Flag{htNameFlag, triggerNamespaceFlag, htFnNameFlag, htHostFlag, htIngressFlag}, Action: htUpdate}, {Name: "delete", Usage: "Delete HTTP trigger", Flags: []cli.Flag{htNameFlag, triggerNamespaceFlag}, Action: htDelete}, {Name: "list", Usage: "List HTTP triggers", Flags: []cli.Flag{triggerNamespaceFlag}, Action: htList}, } diff --git a/router/httpTriggers.go b/router/httpTriggers.go index 04b1cd25..460f084e 100644 --- a/router/httpTriggers.go +++ b/router/httpTriggers.go @@ -25,6 +25,7 @@ import ( "github.com/gorilla/mux" metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/apimachinery/pkg/fields" + "k8s.io/client-go/kubernetes" "k8s.io/client-go/rest" k8sCache "k8s.io/client-go/tools/cache" @@ -38,6 +39,7 @@ type HTTPTriggerSet struct { *mutableRouter fissionClient *crd.FissionClient + kubeClient *kubernetes.Clientset executor *executorClient.Client resolver *functionReferenceResolver crdClient *rest.RESTClient @@ -50,11 +52,12 @@ type HTTPTriggerSet struct { } func makeHTTPTriggerSet(fmap *functionServiceMap, fissionClient *crd.FissionClient, - executor *executorClient.Client, crdClient *rest.RESTClient) (*HTTPTriggerSet, k8sCache.Store, k8sCache.Store) { + kubeClient *kubernetes.Clientset, executor *executorClient.Client, crdClient *rest.RESTClient) (*HTTPTriggerSet, k8sCache.Store, k8sCache.Store) { httpTriggerSet := &HTTPTriggerSet{ functionServiceMap: fmap, triggers: []crd.HTTPTrigger{}, fissionClient: fissionClient, + kubeClient: kubeClient, executor: executor, crdClient: crdClient, } @@ -171,12 +174,19 @@ func (ts *HTTPTriggerSet) initTriggerController() (k8sCache.Store, k8sCache.Cont store, controller := k8sCache.NewInformer(listWatch, &crd.HTTPTrigger{}, resyncPeriod, k8sCache.ResourceEventHandlerFuncs{ AddFunc: func(obj interface{}) { + trigger := obj.(*crd.HTTPTrigger) + go createIngress(trigger, ts.kubeClient) ts.syncTriggers() }, DeleteFunc: func(obj interface{}) { ts.syncTriggers() + trigger := obj.(*crd.HTTPTrigger) + go deleteIngress(trigger, ts.kubeClient) }, UpdateFunc: func(oldObj interface{}, newObj interface{}) { + oldTrigger := oldObj.(*crd.HTTPTrigger) + newTrigger := newObj.(*crd.HTTPTrigger) + go updateIngress(oldTrigger, newTrigger, ts.kubeClient) ts.syncTriggers() }, }) diff --git a/router/ingress.go b/router/ingress.go new file mode 100644 index 00000000..cb650f0d --- /dev/null +++ b/router/ingress.go @@ -0,0 +1,155 @@ +/* +Copyright 2016 The Fission Authors. + +Licensed under the Apache License, Version 2.0 (the "License"); +you may not use this file except in compliance with the License. +You may obtain a copy of the License at + + http://www.apache.org/licenses/LICENSE-2.0 + +Unless required by applicable law or agreed to in writing, software +distributed under the License is distributed on an "AS IS" BASIS, +WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. +See the License for the specific language governing permissions and +limitations under the License. +*/ + +package router + +import ( + "log" + "os" + + metav1 "k8s.io/apimachinery/pkg/apis/meta/v1" + v1 "k8s.io/apimachinery/pkg/apis/meta/v1" + "k8s.io/apimachinery/pkg/util/intstr" + "k8s.io/client-go/kubernetes" + "k8s.io/client-go/pkg/apis/extensions/v1beta1" + + "github.com/fission/fission/crd" +) + +var podNamespace string + +func init() { + podNamespace = os.Getenv("POD_NAMESPACE") + if podNamespace == "" { + podNamespace = "fission" + } +} + +func createIngress(trigger *crd.HTTPTrigger, kubeClient *kubernetes.Clientset) { + + if !trigger.Spec.CreateIngress { + log.Printf("Skipping creation of ingress for trigger: %v", trigger.Metadata.Name) + return + } + + _, err := kubeClient.ExtensionsV1beta1().Ingresses(podNamespace).Get(trigger.Metadata.Name, v1.GetOptions{}) + if err == nil { + log.Printf("Ingress for trigger exists already %v", trigger.Metadata.Name) + return + } + + ing := &v1beta1.Ingress{ + ObjectMeta: metav1.ObjectMeta{ + Labels: getDeployLabels(trigger), + Name: trigger.Metadata.Name, + // The Ingress NS MUST be same as Router NS, check long discussion: + // https://github.com/kubernetes/kubernetes/issues/17088 + // We need to revisit this in future, once Kubernetes supports cross namespace ingress + Namespace: podNamespace, + }, + Spec: v1beta1.IngressSpec{ + Rules: []v1beta1.IngressRule{ + { + Host: trigger.Spec.Host, + IngressRuleValue: v1beta1.IngressRuleValue{ + HTTP: &v1beta1.HTTPIngressRuleValue{ + Paths: []v1beta1.HTTPIngressPath{ + { + Backend: v1beta1.IngressBackend{ + ServiceName: "router", + ServicePort: intstr.IntOrString{ + Type: intstr.Int, + IntVal: 80, + }, + }, + Path: trigger.Spec.RelativeURL, + }, + }, + }, + }, + }, + }, + }, + } + + _, err = kubeClient.ExtensionsV1beta1().Ingresses(podNamespace).Create(ing) + if err != nil { + log.Printf("Failed to create ingress: %v", err) + return + } + log.Printf("Created ingress successfully for trigger %v", trigger.Metadata.Name) +} + +func getDeployLabels(trigger *crd.HTTPTrigger) map[string]string { + return map[string]string{ + "triggerName": trigger.Metadata.Name, + "functionName": trigger.Spec.FunctionReference.Name, + "triggerNamespace": trigger.Metadata.Namespace, + } +} + +func deleteIngress(trigger *crd.HTTPTrigger, kubeClient *kubernetes.Clientset) { + if !trigger.Spec.CreateIngress { + return + } + + ingress, err := kubeClient.ExtensionsV1beta1().Ingresses(podNamespace).Get(trigger.Metadata.Name, v1.GetOptions{}) + if err != nil { + log.Printf("Failed to get ingress when deleting trigger: %v, %v", err, trigger) + } + + err = kubeClient.ExtensionsV1beta1().Ingresses(podNamespace).Delete(ingress.Name, &v1.DeleteOptions{}) + + if err != nil { + log.Printf("Failed to delete ingress %v error: %v", ingress, err) + } + +} + +func updateIngress(oldT *crd.HTTPTrigger, newT *crd.HTTPTrigger, kubeClient *kubernetes.Clientset) { + + if oldT.Spec.CreateIngress == false && newT.Spec.CreateIngress == true { + createIngress(newT, kubeClient) + return + } + + if newT.Spec.CreateIngress == false && oldT.Spec.CreateIngress == true { + deleteIngress(oldT, kubeClient) + return + } + + if newT.Spec.Host != oldT.Spec.Host || newT.Spec.RelativeURL != oldT.Spec.RelativeURL { + log.Printf("Updating ingress for trigger %v", oldT.Metadata.Name) + ingress, err := kubeClient.ExtensionsV1beta1().Ingresses(podNamespace).Get(oldT.Metadata.Name, v1.GetOptions{}) + if err != nil { + log.Printf("Failed to get ingress when updating trigger: %v", err) + } + + if newT.Spec.Host != oldT.Spec.Host { + ingress.Spec.Rules[0].Host = newT.Spec.Host + } + + if newT.Spec.RelativeURL != oldT.Spec.RelativeURL { + ingress.Spec.Rules[0].HTTP.Paths[0].Path = newT.Spec.RelativeURL + } + + _, err = kubeClient.ExtensionsV1beta1().Ingresses(podNamespace).Update(ingress) + if err != nil { + log.Printf("Failed to update ingress for trigger: %v", err) + } + } + +} diff --git a/router/router.go b/router/router.go index fe579ae3..9ea146b7 100644 --- a/router/router.go +++ b/router/router.go @@ -84,7 +84,7 @@ func Start(port int, executorUrl string) { fmap := makeFunctionServiceMap(time.Minute) - fissionClient, _, _, err := crd.MakeFissionClient() + fissionClient, kubeClient, _, err := crd.MakeFissionClient() if err != nil { log.Fatalf("Error connecting to kubernetes API: %v", err) } @@ -97,7 +97,7 @@ func Start(port int, executorUrl string) { restClient := fissionClient.GetCrdClient() executor := executorClient.MakeClient(executorUrl) - triggers, _, fnStore := makeHTTPTriggerSet(fmap, fissionClient, executor, restClient) + triggers, _, fnStore := makeHTTPTriggerSet(fmap, fissionClient, kubeClient, executor, restClient) resolver := makeFunctionReferenceResolver(fnStore) go serveMetric() diff --git a/router/router_test.go b/router/router_test.go index 9e631400..9a6c61d9 100644 --- a/router/router_test.go +++ b/router/router_test.go @@ -59,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, nil) + triggers, _, _ := makeHTTPTriggerSet(fmap, nil, nil, nil, nil) triggerUrl := "/foo" triggers.triggers = append(triggers.triggers, crd.HTTPTrigger{ diff --git a/test/tests/test_ingress.sh b/test/tests/test_ingress.sh new file mode 100755 index 00000000..c96c9333 --- /dev/null +++ b/test/tests/test_ingress.sh @@ -0,0 +1,45 @@ +#!/bin/bash +set -euo pipefail + +ROOT=$(dirname $0)/../.. + +relativeUrl="/itest" +functionName="hellotest" +hostName="test.com" + +cleanup() { + fission route delete --name $1 +} + +log "Creating route for URL $relativeUrl" +route_name=$(fission route create --url $relativeUrl --function $functionName --createingress| grep trigger| cut -d" " -f 2|cut -d"'" -f 2) +trap "cleanup $route_name" EXIT + +log "Route $route_name created" + +sleep 5 + +log "Ingresses matching this trigger:" +kubectl get ing -l 'functionName='$functionName',triggerName='$route_name --all-namespaces -o=json + +log "Verifying to route value in ingress" +actual_route=$(kubectl -n fission get ing -l 'functionName='$functionName',triggerName='$route_name --all-namespaces -o=jsonpath='{.items[0].spec.rules[0].http.paths[0].path}') + +if [ $actual_route != $relativeUrl ] +then + log "Provided route and route in ingress don't match" + exit 1 +fi + +log "Modifying the route by adding host" +fission route update --name $route_name --host $hostName --function $functionName + +sleep 2 + +actual_host=$(kubectl get ing -l 'functionName='$functionName',triggerName='$route_name --all-namespaces -o=jsonpath='{.items[0].spec.rules[0].host}') + +if [ $hostName != $actual_host ] +then + log "Provided host and host in ingress don't match" + exit 1 +fi \ No newline at end of file diff --git a/types.go b/types.go index 512e6c63..00467ea5 100644 --- a/types.go +++ b/types.go @@ -279,6 +279,7 @@ type ( HTTPTriggerSpec struct { Host string `json:"host"` RelativeURL string `json:"relativeurl"` + CreateIngress bool `json:"createingress"` Method string `json:"method"` FunctionReference FunctionReference `json:"functionref"` }