Merge pull request #15 from platform9/poolmgr
Poolmgr -- manage generic containers and their specialization
This commit is contained in:
@@ -54,7 +54,7 @@ func TestFunctionApi(t *testing.T) {
|
||||
m, err := g.client.FunctionCreate(testFunc)
|
||||
panicIf(err)
|
||||
uid1 := m.Uid
|
||||
log.Printf("Created function %v: %v", m.Name, m.Uid)
|
||||
//log.Printf("Created function %v: %v", m.Name, m.Uid)
|
||||
|
||||
code, err := g.client.FunctionGetRaw(m)
|
||||
panicIf(err)
|
||||
@@ -64,7 +64,7 @@ func TestFunctionApi(t *testing.T) {
|
||||
m, err = g.client.FunctionUpdate(testFunc)
|
||||
panicIf(err)
|
||||
uid2 := m.Uid
|
||||
log.Printf("Updated function %v: %v", m.Name, m.Uid)
|
||||
//log.Printf("Updated function %v: %v", m.Name, m.Uid)
|
||||
|
||||
m.Uid = uid1
|
||||
testFunc.Code = "code1"
|
||||
@@ -72,8 +72,8 @@ func TestFunctionApi(t *testing.T) {
|
||||
panicIf(err)
|
||||
|
||||
testFunc.Metadata.Uid = m.Uid
|
||||
log.Printf("f = %#v", f)
|
||||
log.Printf("testFunc = %#v", testFunc)
|
||||
//log.Printf("f = %#v", f)
|
||||
//log.Printf("testFunc = %#v", testFunc)
|
||||
assert(*f == *testFunc, "first version should match when read by uid")
|
||||
|
||||
m.Uid = uid2
|
||||
@@ -212,7 +212,7 @@ func TestMain(m *testing.M) {
|
||||
HTTPTriggerStore: HTTPTriggerStore{ResourceStore: *rs},
|
||||
EnvironmentStore: EnvironmentStore{ResourceStore: *rs},
|
||||
}
|
||||
g.client = client.New("http://localhost:8888")
|
||||
g.client = client.MakeClient("http://localhost:8888")
|
||||
|
||||
ks.Delete(context.Background(), "Function", &etcdClient.DeleteOptions{Recursive: true})
|
||||
ks.Delete(context.Background(), "HTTPTrigger", &etcdClient.DeleteOptions{Recursive: true})
|
||||
|
||||
@@ -36,7 +36,7 @@ type (
|
||||
}
|
||||
)
|
||||
|
||||
func New(serverUrl string) *Client {
|
||||
func MakeClient(serverUrl string) *Client {
|
||||
return &Client{Url: serverUrl}
|
||||
}
|
||||
|
||||
|
||||
+161
@@ -0,0 +1,161 @@
|
||||
/*
|
||||
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 poolmgr
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"io/ioutil"
|
||||
"log"
|
||||
"net/http"
|
||||
"os"
|
||||
"time"
|
||||
|
||||
"github.com/gorilla/handlers"
|
||||
"github.com/gorilla/mux"
|
||||
|
||||
"github.com/platform9/fission"
|
||||
"github.com/platform9/fission/cache"
|
||||
controllerclient "github.com/platform9/fission/controller/client"
|
||||
)
|
||||
|
||||
type funcSvc struct {
|
||||
function *fission.Metadata // function this thing is for
|
||||
environment *fission.Environment // env it was obtained from
|
||||
serviceName string // name of k8s svc
|
||||
|
||||
ctime time.Time
|
||||
atime time.Time
|
||||
}
|
||||
|
||||
type API struct {
|
||||
poolMgr *GenericPoolManager
|
||||
functionEnv *cache.Cache // map[fission.Metadata]fission.Environment
|
||||
functionService *cache.Cache // map[fission.Metadata]funcSvc
|
||||
controller *controllerclient.Client
|
||||
}
|
||||
|
||||
func MakeAPI(gpm *GenericPoolManager, controller *controllerclient.Client) *API {
|
||||
return &API{
|
||||
poolMgr: gpm,
|
||||
functionEnv: cache.MakeCache(),
|
||||
functionService: cache.MakeCache(),
|
||||
controller: controller,
|
||||
}
|
||||
}
|
||||
|
||||
func (api *API) lookupApi(w http.ResponseWriter, r *http.Request) {
|
||||
body, err := ioutil.ReadAll(r.Body)
|
||||
if err != nil {
|
||||
http.Error(w, "Failed to read request", 500)
|
||||
return
|
||||
}
|
||||
|
||||
// get function metadata
|
||||
m := fission.Metadata{}
|
||||
err = json.Unmarshal(body, &m)
|
||||
if err != nil {
|
||||
http.Error(w, "Failed to parse request", 400)
|
||||
return
|
||||
}
|
||||
|
||||
serviceUrl, err := api.lookup(&m)
|
||||
if err != nil {
|
||||
code, msg := fission.GetHTTPError(err)
|
||||
log.Printf("Error: %v: %v", code, msg)
|
||||
http.Error(w, msg, code)
|
||||
}
|
||||
|
||||
// return serviceUrl
|
||||
w.Write([]byte(serviceUrl))
|
||||
}
|
||||
|
||||
func (api *API) getFunctionEnv(m *fission.Metadata) (*fission.Environment, error) {
|
||||
var env *fission.Environment
|
||||
|
||||
// Cached ?
|
||||
result, err := api.functionEnv.Get(m)
|
||||
if err == nil {
|
||||
env = result.(*fission.Environment)
|
||||
return env, nil
|
||||
}
|
||||
|
||||
// Cache miss -- get func from controller
|
||||
f, err := api.controller.FunctionGet(m)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// Get env from metadata
|
||||
env, err = api.controller.EnvironmentGet(&f.Environment)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// cache for future
|
||||
api.functionEnv.Set(m, env)
|
||||
|
||||
return env, nil
|
||||
}
|
||||
|
||||
func (api *API) lookup(m *fission.Metadata) (string, error) {
|
||||
// Check function -> svc map
|
||||
result, err := api.functionService.Get(m)
|
||||
if err == nil {
|
||||
// Ok: return svc name
|
||||
svc := result.(*funcSvc)
|
||||
return svc.serviceName, nil
|
||||
}
|
||||
|
||||
// None exists, so create a new funcSvc:
|
||||
|
||||
// from Func -> get Env
|
||||
env, err := api.getFunctionEnv(m)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
// from Env -> get GenericPool
|
||||
pool, err := api.poolMgr.GetPool(env)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
// from GenericPool -> get one function container
|
||||
funcSvc, err := pool.GetFuncSvc(m)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
|
||||
// add to cache
|
||||
err = api.functionService.Set(m, funcSvc)
|
||||
if err != nil {
|
||||
// log and ignore error
|
||||
log.Printf("Error saving function service: %v", err)
|
||||
}
|
||||
|
||||
return funcSvc.serviceName, nil
|
||||
}
|
||||
|
||||
func (api *API) Serve(port int) {
|
||||
r := mux.NewRouter()
|
||||
r.HandleFunc("/v1/lookup", api.lookupApi).Methods("GET")
|
||||
|
||||
address := fmt.Sprintf(":%v", port)
|
||||
log.Printf("starting poolmgr at port %v", port)
|
||||
log.Fatal(http.ListenAndServe(address, handlers.LoggingHandler(os.Stdout, r)))
|
||||
}
|
||||
+373
@@ -0,0 +1,373 @@
|
||||
/*
|
||||
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 poolmgr
|
||||
|
||||
import (
|
||||
"bytes"
|
||||
"errors"
|
||||
"fmt"
|
||||
"log"
|
||||
"math/rand"
|
||||
"net/http"
|
||||
"time"
|
||||
|
||||
"k8s.io/client-go/1.4/kubernetes"
|
||||
"k8s.io/client-go/1.4/pkg/api"
|
||||
"k8s.io/client-go/1.4/pkg/api/v1"
|
||||
"k8s.io/client-go/1.4/pkg/apis/extensions/v1beta1"
|
||||
"k8s.io/client-go/1.4/pkg/labels"
|
||||
|
||||
"github.com/platform9/fission"
|
||||
)
|
||||
|
||||
type (
|
||||
GenericPool struct {
|
||||
env *fission.Environment
|
||||
replicas int32 // num containers
|
||||
deployment *v1beta1.Deployment // kubernetes deployment
|
||||
namespace string // namespace to keep our resources
|
||||
podReadyTimeout time.Duration // timeout for generic pods to become ready
|
||||
controllerUrl string
|
||||
|
||||
kubernetesClient *kubernetes.Clientset
|
||||
requestChannel chan *choosePodRequest
|
||||
}
|
||||
|
||||
// serialize the choosing of pods so that choices don't conflict
|
||||
choosePodRequest struct {
|
||||
newLabels map[string]string
|
||||
responseChannel chan *choosePodResponse
|
||||
}
|
||||
choosePodResponse struct {
|
||||
pod *v1.Pod
|
||||
error
|
||||
}
|
||||
)
|
||||
|
||||
func MakeGenericPool(
|
||||
controllerUrl string,
|
||||
kubernetesClient *kubernetes.Clientset,
|
||||
env *fission.Environment,
|
||||
initialReplicas int32,
|
||||
namespace string) (*GenericPool, error) {
|
||||
|
||||
gp := &GenericPool{
|
||||
env: env,
|
||||
replicas: initialReplicas,
|
||||
requestChannel: make(chan *choosePodRequest),
|
||||
kubernetesClient: kubernetesClient,
|
||||
namespace: namespace,
|
||||
podReadyTimeout: 5 * time.Minute,
|
||||
controllerUrl: controllerUrl,
|
||||
}
|
||||
|
||||
// create the pool
|
||||
err := gp.createPool()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// wait for at least one pod to be ready
|
||||
err = gp.waitForReadyPod()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
go gp.choosePodService()
|
||||
return gp, nil
|
||||
}
|
||||
|
||||
// choosePodService serializes the choosing of pods
|
||||
func (gp *GenericPool) choosePodService() {
|
||||
for {
|
||||
select {
|
||||
case req := <-gp.requestChannel:
|
||||
pod, err := gp._choosePod(req.newLabels)
|
||||
if err != nil {
|
||||
req.responseChannel <- &choosePodResponse{error: err}
|
||||
continue
|
||||
}
|
||||
req.responseChannel <- &choosePodResponse{pod: pod}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// choosePod picks a ready pod from the pool and relabels it, waiting if necessary.
|
||||
// returns the pod API object.
|
||||
func (gp *GenericPool) choosePod(newLabels map[string]string) (*v1.Pod, error) {
|
||||
req := &choosePodRequest{
|
||||
newLabels: newLabels,
|
||||
responseChannel: make(chan *choosePodResponse),
|
||||
}
|
||||
gp.requestChannel <- req
|
||||
resp := <-req.responseChannel
|
||||
return resp.pod, resp.error
|
||||
}
|
||||
|
||||
// _choosePod is called serially by choosePodService
|
||||
func (gp *GenericPool) _choosePod(newLabels map[string]string) (*v1.Pod, error) {
|
||||
startTime := time.Now()
|
||||
for {
|
||||
// Retries took too long, error out.
|
||||
if time.Now().Sub(startTime) > gp.podReadyTimeout {
|
||||
return nil, errors.New("timeout: waited too long to get a ready pod")
|
||||
}
|
||||
|
||||
// Get pods; filter the ones that are ready
|
||||
podList, err := gp.kubernetesClient.Core().Pods(gp.namespace).List(
|
||||
api.ListOptions{
|
||||
LabelSelector: labels.Set(
|
||||
gp.deployment.Spec.Selector.MatchLabels).AsSelector(),
|
||||
})
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
readyPods := make([]v1.Pod, len(podList.Items))
|
||||
for _, pod := range podList.Items {
|
||||
podReady := true
|
||||
for _, cs := range pod.Status.ContainerStatuses {
|
||||
podReady = podReady && cs.Ready
|
||||
}
|
||||
if podReady {
|
||||
readyPods = append(readyPods, pod)
|
||||
}
|
||||
}
|
||||
|
||||
// If there are no ready pods, wait and retry.
|
||||
if len(readyPods) == 0 {
|
||||
err = gp.waitForReadyPod()
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
continue
|
||||
}
|
||||
|
||||
// Pick a ready pod. For now just choose randomly;
|
||||
// ideally we'd care about which node it's running on,
|
||||
// and make a good scheduling decision.
|
||||
chosenPod := readyPods[rand.Intn(len(readyPods))]
|
||||
|
||||
// Relabel. If the pod already got picked and
|
||||
// modified, this should fail; in that case just
|
||||
// retry.
|
||||
chosenPod.ObjectMeta.Labels = newLabels
|
||||
_, err = gp.kubernetesClient.Core().Pods(gp.namespace).Update(&chosenPod)
|
||||
if err != nil {
|
||||
log.Printf("failed to relabel pod: %v", err)
|
||||
continue
|
||||
}
|
||||
log.Printf("Chose a pod: %v", chosenPod.ObjectMeta.Name)
|
||||
return &chosenPod, nil
|
||||
}
|
||||
}
|
||||
|
||||
func labelsForMetadata(metadata *fission.Metadata) map[string]string {
|
||||
return map[string]string{
|
||||
"functionName": metadata.Name,
|
||||
"functionUid": metadata.Uid,
|
||||
}
|
||||
}
|
||||
|
||||
// specializePod chooses a pod, copies the required user-defined function to that pod
|
||||
// (via fetcher), and calls the function-run container to load it, resulting in a
|
||||
// specialized pod.
|
||||
func (gp *GenericPool) specializePod(metadata *fission.Metadata) (*v1.Pod, error) {
|
||||
newLabels := labelsForMetadata(metadata)
|
||||
|
||||
pod, err := gp.choosePod(newLabels)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
// for fetcher we don't need to create a service, just talk to the pod directly
|
||||
podIP := pod.Status.PodIP
|
||||
if len(podIP) == 0 {
|
||||
return nil, errors.New("Pod has no IP")
|
||||
}
|
||||
|
||||
// tell fetcher to get the function
|
||||
fetcherUrl := fmt.Sprintf("http://%v:8000/", podIP)
|
||||
functionUrl := fmt.Sprintf("%v/v1/functions/%v?uid=%v&raw=1",
|
||||
gp.controllerUrl, metadata.Name, metadata.Uid)
|
||||
fetcherRequest := fmt.Sprintf("{\"url\": \"%v\", \"filename\": \"user\"}", functionUrl)
|
||||
|
||||
resp, err := http.Post(fetcherUrl, "application/json", bytes.NewReader([]byte(fetcherRequest)))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
defer resp.Body.Close()
|
||||
if resp.StatusCode != 200 {
|
||||
return nil, errors.New(fmt.Sprintf("Error from fetcher: %v", resp.Status))
|
||||
}
|
||||
|
||||
// get function run container to specialize
|
||||
specializeUrl := fmt.Sprintf("http://%v:8888/specialize", podIP)
|
||||
resp2, err := http.Post(specializeUrl, "", bytes.NewReader([]byte{}))
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
resp2.Body.Close()
|
||||
return pod, nil
|
||||
}
|
||||
|
||||
// A pool is a deployment of generic containers for an env. This
|
||||
// creates the pool but doesn't wait for any pods to be ready.
|
||||
func (gp *GenericPool) createPool() error {
|
||||
poolDeploymentName := fmt.Sprintf("deployment-%v-%v-0",
|
||||
gp.env.Metadata.Name, gp.env.Metadata.Uid)
|
||||
|
||||
podLabels := map[string]string{
|
||||
"pool": poolDeploymentName,
|
||||
}
|
||||
|
||||
sharedMountPath := "/userfunc"
|
||||
deployment := &v1beta1.Deployment{
|
||||
ObjectMeta: v1.ObjectMeta{
|
||||
Name: poolDeploymentName,
|
||||
Labels: map[string]string{
|
||||
"environmentName": gp.env.Metadata.Name,
|
||||
"environmentUid": gp.env.Metadata.Uid,
|
||||
},
|
||||
},
|
||||
Spec: v1beta1.DeploymentSpec{
|
||||
Replicas: &gp.replicas,
|
||||
Selector: &v1beta1.LabelSelector{
|
||||
MatchLabels: podLabels,
|
||||
},
|
||||
Template: v1.PodTemplateSpec{
|
||||
ObjectMeta: v1.ObjectMeta{
|
||||
Labels: podLabels,
|
||||
},
|
||||
Spec: v1.PodSpec{
|
||||
Volumes: []v1.Volume{
|
||||
v1.Volume{
|
||||
Name: "userfunc",
|
||||
VolumeSource: v1.VolumeSource{
|
||||
EmptyDir: &v1.EmptyDirVolumeSource{},
|
||||
},
|
||||
},
|
||||
},
|
||||
Containers: []v1.Container{
|
||||
v1.Container{
|
||||
Name: gp.env.Metadata.Name,
|
||||
Image: gp.env.RunContainerImageUrl,
|
||||
ImagePullPolicy: v1.PullIfNotPresent,
|
||||
TerminationMessagePath: "/dev/termination-log",
|
||||
VolumeMounts: []v1.VolumeMount{
|
||||
v1.VolumeMount{
|
||||
Name: "userfunc",
|
||||
MountPath: sharedMountPath,
|
||||
},
|
||||
},
|
||||
},
|
||||
v1.Container{
|
||||
Name: "fetcher",
|
||||
Image: "fission/fetcher",
|
||||
ImagePullPolicy: v1.PullIfNotPresent,
|
||||
TerminationMessagePath: "/dev/termination-log",
|
||||
VolumeMounts: []v1.VolumeMount{
|
||||
v1.VolumeMount{
|
||||
Name: "userfunc",
|
||||
MountPath: sharedMountPath,
|
||||
},
|
||||
},
|
||||
Command: []string{"/fetcher", sharedMountPath},
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
},
|
||||
}
|
||||
depl, err := gp.kubernetesClient.Extensions().Deployments(gp.namespace).Create(deployment)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
gp.deployment = depl
|
||||
return nil
|
||||
}
|
||||
|
||||
func (gp *GenericPool) waitForReadyPod() error {
|
||||
startTime := time.Now()
|
||||
for {
|
||||
// TODO: for now we just poll; use a watch instead
|
||||
depl, err := gp.kubernetesClient.Extensions().Deployments(gp.namespace).Get(gp.deployment.ObjectMeta.Name)
|
||||
if err != nil {
|
||||
log.Printf("err: %v", err)
|
||||
return err
|
||||
}
|
||||
gp.deployment = depl
|
||||
if gp.deployment.Status.AvailableReplicas > 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
if time.Now().Sub(startTime) > gp.podReadyTimeout {
|
||||
return errors.New("timeout: waited too long for pod to be ready")
|
||||
}
|
||||
time.Sleep(1000 * time.Millisecond)
|
||||
}
|
||||
}
|
||||
|
||||
func (gp *GenericPool) createSvc(name string, labels map[string]string) (*v1.Service, error) {
|
||||
service := v1.Service{
|
||||
ObjectMeta: v1.ObjectMeta{
|
||||
Name: name,
|
||||
},
|
||||
Spec: v1.ServiceSpec{
|
||||
Type: v1.ServiceTypeClusterIP,
|
||||
Ports: []v1.ServicePort{
|
||||
v1.ServicePort{
|
||||
Protocol: v1.ProtocolTCP,
|
||||
Port: 8888,
|
||||
},
|
||||
},
|
||||
Selector: labels,
|
||||
},
|
||||
}
|
||||
svc, err := gp.kubernetesClient.Core().Services(gp.namespace).Create(&service)
|
||||
return svc, err
|
||||
}
|
||||
|
||||
func (gp *GenericPool) GetFuncSvc(m *fission.Metadata) (*funcSvc, error) {
|
||||
pod, err := gp.specializePod(m)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
log.Printf("Specialized pod: %v", pod.ObjectMeta.Name)
|
||||
|
||||
svcName := fmt.Sprintf("svc-%v", m.Name)
|
||||
if len(m.Uid) > 0 {
|
||||
svcName += ("-" + m.Uid)
|
||||
}
|
||||
|
||||
labels := labelsForMetadata(m)
|
||||
svc, err := gp.createSvc(svcName, labels)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
if svc.ObjectMeta.Name != svcName {
|
||||
return nil, errors.New(fmt.Sprintf("sanity check failed for svc %v", svc.ObjectMeta.Name))
|
||||
}
|
||||
|
||||
fsvc := &funcSvc{
|
||||
function: m,
|
||||
environment: gp.env,
|
||||
serviceName: svcName,
|
||||
ctime: time.Now(),
|
||||
atime: time.Now(),
|
||||
}
|
||||
return fsvc, nil
|
||||
}
|
||||
@@ -0,0 +1,93 @@
|
||||
package poolmgr
|
||||
|
||||
import (
|
||||
"github.com/platform9/fission"
|
||||
|
||||
"fmt"
|
||||
"k8s.io/kubernetes/pkg/api"
|
||||
"k8s.io/kubernetes/pkg/client/unversioned"
|
||||
"k8s.io/kubernetes/pkg/client/unversioned/clientcmd"
|
||||
"log"
|
||||
"net/http"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func getKubeClient() *unversioned.Client {
|
||||
loadingRules := clientcmd.NewDefaultClientConfigLoadingRules()
|
||||
configOverrides := &clientcmd.ConfigOverrides{}
|
||||
kubeConfig := clientcmd.NewNonInteractiveDeferredLoadingClientConfig(loadingRules, configOverrides)
|
||||
config, err := kubeConfig.ClientConfig()
|
||||
if err != nil {
|
||||
panic("failed loading client config")
|
||||
}
|
||||
client := unversioned.NewOrDie(config)
|
||||
return client
|
||||
}
|
||||
|
||||
type staticHandler struct {
|
||||
resp string
|
||||
}
|
||||
|
||||
func (s *staticHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
w.Write([]byte(s.resp))
|
||||
}
|
||||
|
||||
// staticHttpServer starts an http server at port and responds to any
|
||||
// request with the given response. Use this to mock the controller
|
||||
// raw function fetch HTTP endpoint.
|
||||
func staticHttpServer(port int, response string) {
|
||||
s := &staticHandler{resp: response}
|
||||
log.Fatal(http.ListenAndServe(fmt.Sprintf(":%v", port), s))
|
||||
}
|
||||
|
||||
func TestGenericPool(t *testing.T) {
|
||||
namespace := "fission-test"
|
||||
|
||||
client := getKubeClient()
|
||||
|
||||
_, err := client.Namespaces().Create(&api.Namespace{
|
||||
ObjectMeta: api.ObjectMeta{
|
||||
Name: namespace,
|
||||
Labels: map[string]string{},
|
||||
},
|
||||
})
|
||||
if err != nil {
|
||||
log.Panicf("failed to create namespace: %v", err)
|
||||
}
|
||||
|
||||
// destroys everything in the namespace
|
||||
defer client.Namespaces().Delete(namespace)
|
||||
|
||||
env := &fission.Environment{
|
||||
Metadata: fission.Metadata{
|
||||
Name: "test-env",
|
||||
Uid: "",
|
||||
},
|
||||
RunContainerImageUrl: "fission/testing",
|
||||
}
|
||||
|
||||
gp, err := MakeGenericPool(client, env, 3, namespace)
|
||||
if err != nil {
|
||||
log.Panicf("failed to make generic pool: %v", err)
|
||||
}
|
||||
log.Printf("Pool created")
|
||||
|
||||
// test specialization
|
||||
|
||||
testFunc := `
|
||||
module.exports = function (context, callback) {
|
||||
callback(200, "Hello, world!");
|
||||
}
|
||||
`
|
||||
go staticHttpServer(2222, testFunc)
|
||||
|
||||
m := fission.Metadata{
|
||||
Name: "foo",
|
||||
Uid: "xxx-yyy",
|
||||
}
|
||||
fsvc, err := gp.GetFuncSvc(&m)
|
||||
if err != nil {
|
||||
log.Fatalf("Error getting function svc: %v", err)
|
||||
}
|
||||
log.Printf("fsvc: %v", fsvc)
|
||||
}
|
||||
@@ -0,0 +1,79 @@
|
||||
/*
|
||||
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 poolmgr
|
||||
|
||||
import (
|
||||
"github.com/platform9/fission"
|
||||
"k8s.io/client-go/1.4/kubernetes"
|
||||
)
|
||||
|
||||
type (
|
||||
GenericPoolManager struct {
|
||||
pools map[fission.Environment]*GenericPool
|
||||
kubernetesClient *kubernetes.Clientset
|
||||
namespace string
|
||||
controllerUrl string
|
||||
|
||||
requestChannel chan *request
|
||||
}
|
||||
request struct {
|
||||
env *fission.Environment
|
||||
responseChannel chan *response
|
||||
}
|
||||
response struct {
|
||||
error
|
||||
pool *GenericPool
|
||||
}
|
||||
)
|
||||
|
||||
func MakeGenericPoolManager(controllerUrl string, kubernetesClient *kubernetes.Clientset, namespace string) *GenericPoolManager {
|
||||
gpm := &GenericPoolManager{
|
||||
pools: make(map[fission.Environment]*GenericPool),
|
||||
kubernetesClient: kubernetesClient,
|
||||
namespace: namespace,
|
||||
controllerUrl: controllerUrl,
|
||||
requestChannel: make(chan *request),
|
||||
}
|
||||
go gpm.service()
|
||||
return gpm
|
||||
}
|
||||
|
||||
func (gpm *GenericPoolManager) service() {
|
||||
for {
|
||||
select {
|
||||
case req := <-gpm.requestChannel:
|
||||
pool, ok := gpm.pools[*req.env]
|
||||
if !ok {
|
||||
pool, err := MakeGenericPool(gpm.controllerUrl, gpm.kubernetesClient, req.env, 3, gpm.namespace)
|
||||
if err != nil {
|
||||
req.responseChannel <- &response{error: err}
|
||||
continue
|
||||
}
|
||||
gpm.pools[*req.env] = pool
|
||||
req.responseChannel <- &response{pool: pool}
|
||||
}
|
||||
req.responseChannel <- &response{pool: pool}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (gpm *GenericPoolManager) GetPool(env *fission.Environment) (*GenericPool, error) {
|
||||
c := make(chan *response)
|
||||
gpm.requestChannel <- &request{env: env, responseChannel: c}
|
||||
resp := <-c
|
||||
return resp.pool, resp.error
|
||||
}
|
||||
Reference in New Issue
Block a user