Switch to official Kubernetes Go client -- client-go/1.4
Switch to client-go package instead of pulling the from kubernetes. Use a versioned client with sensible compatiblity. This change breaks 'go get'. For now you have to manually checkout the 'release-1.4' branch of the client-go package after 'go get' fetches it. TODO: use one of the build tools to fix this.
This commit is contained in:
+25
-2
@@ -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)))
|
||||
}
|
||||
|
||||
+64
-63
@@ -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
|
||||
}
|
||||
|
||||
|
||||
+7
-5
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user