diff --git a/poolmgr/api.go b/poolmgr/api.go index e7c0e4f1..d457650d 100644 --- a/poolmgr/api.go +++ b/poolmgr/api.go @@ -18,14 +18,19 @@ 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" - "io/ioutil" - "log" ) type funcSvc struct { @@ -44,6 +49,15 @@ type API struct { 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 { @@ -136,3 +150,12 @@ func (api *API) lookup(m *fission.Metadata) (string, error) { 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))) +} diff --git a/poolmgr/gp.go b/poolmgr/gp.go index 59bfa1c4..cb1919ee 100644 --- a/poolmgr/gp.go +++ b/poolmgr/gp.go @@ -25,25 +25,25 @@ import ( "net/http" "time" - "github.com/platform9/fission" + "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" - "k8s.io/kubernetes/pkg/api" - apiUnversioned "k8s.io/kubernetes/pkg/api/unversioned" - "k8s.io/kubernetes/pkg/apis/extensions" - clientUnversioned "k8s.io/kubernetes/pkg/client/unversioned" - "k8s.io/kubernetes/pkg/labels" + "github.com/platform9/fission" ) type ( GenericPool struct { - env *fission.Environment - replicas int // num containers - deployment *extensions.Deployment // kubernetes deployment - namespace string // namespace to keep our resources - podReadyTimeout time.Duration // timeout for generic pods to become ready - controllerHostName string + 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 *clientUnversioned.Client + kubernetesClient *kubernetes.Clientset requestChannel chan *choosePodRequest } @@ -53,25 +53,26 @@ type ( responseChannel chan *choosePodResponse } choosePodResponse struct { - pod *api.Pod + pod *v1.Pod error } ) func MakeGenericPool( - kubernetesClient *clientUnversioned.Client, + controllerUrl string, + kubernetesClient *kubernetes.Clientset, env *fission.Environment, - initialReplicas int, + initialReplicas int32, namespace string) (*GenericPool, error) { gp := &GenericPool{ - env: env, - replicas: initialReplicas, - requestChannel: make(chan *choosePodRequest), - kubernetesClient: kubernetesClient, - namespace: namespace, - podReadyTimeout: 5 * time.Minute, - controllerHostName: "controller", + env: env, + replicas: initialReplicas, + requestChannel: make(chan *choosePodRequest), + kubernetesClient: kubernetesClient, + namespace: namespace, + podReadyTimeout: 5 * time.Minute, + controllerUrl: controllerUrl, } // create the pool @@ -107,7 +108,7 @@ func (gp *GenericPool) choosePodService() { // 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) (*api.Pod, error) { +func (gp *GenericPool) choosePod(newLabels map[string]string) (*v1.Pod, error) { req := &choosePodRequest{ newLabels: newLabels, responseChannel: make(chan *choosePodResponse), @@ -118,7 +119,7 @@ func (gp *GenericPool) choosePod(newLabels map[string]string) (*api.Pod, error) } // _choosePod is called serially by choosePodService -func (gp *GenericPool) _choosePod(newLabels map[string]string) (*api.Pod, error) { +func (gp *GenericPool) _choosePod(newLabels map[string]string) (*v1.Pod, error) { startTime := time.Now() for { // Retries took too long, error out. @@ -127,7 +128,7 @@ func (gp *GenericPool) _choosePod(newLabels map[string]string) (*api.Pod, error) } // Get pods; filter the ones that are ready - podList, err := gp.kubernetesClient.Pods(gp.namespace).List( + podList, err := gp.kubernetesClient.Core().Pods(gp.namespace).List( api.ListOptions{ LabelSelector: labels.Set( gp.deployment.Spec.Selector.MatchLabels).AsSelector(), @@ -135,7 +136,7 @@ func (gp *GenericPool) _choosePod(newLabels map[string]string) (*api.Pod, error) if err != nil { return nil, err } - readyPods := make([]api.Pod, len(podList.Items)) + readyPods := make([]v1.Pod, len(podList.Items)) for _, pod := range podList.Items { podReady := true for _, cs := range pod.Status.ContainerStatuses { @@ -164,7 +165,7 @@ func (gp *GenericPool) _choosePod(newLabels map[string]string) (*api.Pod, error) // modified, this should fail; in that case just // retry. chosenPod.ObjectMeta.Labels = newLabels - _, err = gp.kubernetesClient.Pods(gp.namespace).Update(&chosenPod) + _, err = gp.kubernetesClient.Core().Pods(gp.namespace).Update(&chosenPod) if err != nil { log.Printf("failed to relabel pod: %v", err) continue @@ -184,7 +185,7 @@ func labelsForMetadata(metadata *fission.Metadata) map[string]string { // 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) (*api.Pod, error) { +func (gp *GenericPool) specializePod(metadata *fission.Metadata) (*v1.Pod, error) { newLabels := labelsForMetadata(metadata) pod, err := gp.choosePod(newLabels) @@ -200,8 +201,8 @@ func (gp *GenericPool) specializePod(metadata *fission.Metadata) (*api.Pod, erro // tell fetcher to get the function fetcherUrl := fmt.Sprintf("http://%v:8000/", podIP) - functionUrl := fmt.Sprintf("http://%v/v1/functions/%v?uid=%v&raw=1", - gp.controllerHostName, metadata.Name, metadata.Uid) + 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))) @@ -234,52 +235,52 @@ func (gp *GenericPool) createPool() error { } sharedMountPath := "/userfunc" - deployment := &extensions.Deployment{ - ObjectMeta: api.ObjectMeta{ + deployment := &v1beta1.Deployment{ + ObjectMeta: v1.ObjectMeta{ Name: poolDeploymentName, Labels: map[string]string{ "environmentName": gp.env.Metadata.Name, "environmentUid": gp.env.Metadata.Uid, }, }, - Spec: extensions.DeploymentSpec{ - Replicas: int32(gp.replicas), - Selector: &apiUnversioned.LabelSelector{ + Spec: v1beta1.DeploymentSpec{ + Replicas: &gp.replicas, + Selector: &v1beta1.LabelSelector{ MatchLabels: podLabels, }, - Template: api.PodTemplateSpec{ - ObjectMeta: api.ObjectMeta{ + Template: v1.PodTemplateSpec{ + ObjectMeta: v1.ObjectMeta{ Labels: podLabels, }, - Spec: api.PodSpec{ - Volumes: []api.Volume{ - api.Volume{ + Spec: v1.PodSpec{ + Volumes: []v1.Volume{ + v1.Volume{ Name: "userfunc", - VolumeSource: api.VolumeSource{ - EmptyDir: &api.EmptyDirVolumeSource{}, + VolumeSource: v1.VolumeSource{ + EmptyDir: &v1.EmptyDirVolumeSource{}, }, }, }, - Containers: []api.Container{ - api.Container{ + Containers: []v1.Container{ + v1.Container{ Name: gp.env.Metadata.Name, Image: gp.env.RunContainerImageUrl, - ImagePullPolicy: api.PullIfNotPresent, + ImagePullPolicy: v1.PullIfNotPresent, TerminationMessagePath: "/dev/termination-log", - VolumeMounts: []api.VolumeMount{ - api.VolumeMount{ + VolumeMounts: []v1.VolumeMount{ + v1.VolumeMount{ Name: "userfunc", MountPath: sharedMountPath, }, }, }, - api.Container{ + v1.Container{ Name: "fetcher", Image: "fission/fetcher", - ImagePullPolicy: api.PullIfNotPresent, + ImagePullPolicy: v1.PullIfNotPresent, TerminationMessagePath: "/dev/termination-log", - VolumeMounts: []api.VolumeMount{ - api.VolumeMount{ + VolumeMounts: []v1.VolumeMount{ + v1.VolumeMount{ Name: "userfunc", MountPath: sharedMountPath, }, @@ -291,7 +292,7 @@ func (gp *GenericPool) createPool() error { }, }, } - depl, err := gp.kubernetesClient.ExtensionsClient.Deployments(gp.namespace).Create(deployment) + depl, err := gp.kubernetesClient.Extensions().Deployments(gp.namespace).Create(deployment) if err != nil { return err } @@ -303,7 +304,7 @@ func (gp *GenericPool) waitForReadyPod() error { startTime := time.Now() for { // TODO: for now we just poll; use a watch instead - depl, err := gp.kubernetesClient.ExtensionsClient.Deployments(gp.namespace).Get(gp.deployment.ObjectMeta.Name) + depl, err := gp.kubernetesClient.Extensions().Deployments(gp.namespace).Get(gp.deployment.ObjectMeta.Name) if err != nil { log.Printf("err: %v", err) return err @@ -320,23 +321,23 @@ func (gp *GenericPool) waitForReadyPod() error { } } -func (gp *GenericPool) createSvc(name string, labels map[string]string) (*api.Service, error) { - service := api.Service{ - ObjectMeta: api.ObjectMeta{ +func (gp *GenericPool) createSvc(name string, labels map[string]string) (*v1.Service, error) { + service := v1.Service{ + ObjectMeta: v1.ObjectMeta{ Name: name, }, - Spec: api.ServiceSpec{ - Type: api.ServiceTypeClusterIP, - Ports: []api.ServicePort{ - api.ServicePort{ - Protocol: api.ProtocolTCP, + Spec: v1.ServiceSpec{ + Type: v1.ServiceTypeClusterIP, + Ports: []v1.ServicePort{ + v1.ServicePort{ + Protocol: v1.ProtocolTCP, Port: 8888, }, }, Selector: labels, }, } - svc, err := gp.kubernetesClient.Services(gp.namespace).Create(&service) + svc, err := gp.kubernetesClient.Core().Services(gp.namespace).Create(&service) return svc, err } diff --git a/poolmgr/gpm.go b/poolmgr/gpm.go index 49799bcb..876a39f2 100644 --- a/poolmgr/gpm.go +++ b/poolmgr/gpm.go @@ -18,14 +18,15 @@ package poolmgr import ( "github.com/platform9/fission" - clientUnversioned "k8s.io/kubernetes/pkg/client/unversioned" + "k8s.io/client-go/1.4/kubernetes" ) type ( GenericPoolManager struct { pools map[fission.Environment]*GenericPool - kubernetesClient *clientUnversioned.Client + kubernetesClient *kubernetes.Clientset namespace string + controllerUrl string requestChannel chan *request } @@ -39,11 +40,12 @@ type ( } ) -func MakeGenericPoolManager(client *clientUnversioned.Client, namespace string) *GenericPoolManager { +func MakeGenericPoolManager(controllerUrl string, kubernetesClient *kubernetes.Clientset, namespace string) *GenericPoolManager { gpm := &GenericPoolManager{ pools: make(map[fission.Environment]*GenericPool), - kubernetesClient: client, + kubernetesClient: kubernetesClient, namespace: namespace, + controllerUrl: controllerUrl, requestChannel: make(chan *request), } go gpm.service() @@ -56,7 +58,7 @@ func (gpm *GenericPoolManager) service() { case req := <-gpm.requestChannel: pool, ok := gpm.pools[*req.env] if !ok { - pool, err := MakeGenericPool(gpm.kubernetesClient, req.env, 3, gpm.namespace) + pool, err := MakeGenericPool(gpm.controllerUrl, gpm.kubernetesClient, req.env, 3, gpm.namespace) if err != nil { req.responseChannel <- &response{error: err} continue