Ingress integration (#688)
Ingress integration to allow the optional creation of ingress for a given route. The ingress controller needs to be set up by the user separately so that ingress path is accessible outside the cluster.
This commit is contained in:
+11
-1
@@ -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()
|
||||
},
|
||||
})
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
+2
-2
@@ -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()
|
||||
|
||||
@@ -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{
|
||||
|
||||
Reference in New Issue
Block a user