refactor file tree structure

This commit is contained in:
Murat Kabilov
2017-05-08 12:10:25 +02:00
parent 77b01c67c9
commit 256ff37c19
7 changed files with 14 additions and 14 deletions
+27
View File
@@ -0,0 +1,27 @@
package controller
import (
"fmt"
"github.com/coreos/etcd/client"
"golang.org/x/net/context"
"log"
)
func (z *SpiloSupervisor) DeleteEtcdKey(clusterName string) error {
options := client.DeleteOptions{
Recursive: true,
}
keyName := fmt.Sprintf(etcdKeyTemplate, clusterName)
resp, err := z.etcdApiClient.Delete(context.Background(), keyName, &options)
if resp != nil {
log.Printf("Response: %+v", *resp)
} else {
log.Fatal("No response from etcd")
}
log.Printf("Deleting key %s from ETCD", clusterName)
return err
}
+234
View File
@@ -0,0 +1,234 @@
package controller
import (
"k8s.io/client-go/pkg/api/resource"
"k8s.io/client-go/pkg/api/v1"
"k8s.io/client-go/pkg/apis/apps/v1beta1"
"k8s.io/client-go/pkg/util/intstr"
"log"
)
func (z *SpiloSupervisor) CreateStatefulSet(spilo *Spilo) {
ns := (*spilo).Metadata.Namespace
statefulSet := z.createSetFromSpilo(spilo)
_, err := z.Clientset.StatefulSets(ns).Create(&statefulSet)
if err != nil {
log.Printf("Petset error: %+v", err)
} else {
log.Printf("Petset created: %+v", statefulSet)
}
}
func (z *SpiloSupervisor) createSetFromSpilo(spilo *Spilo) v1beta1.StatefulSet {
clusterName := (*spilo).Metadata.Name
envVars := []v1.EnvVar{
{
Name: "SCOPE",
Value: clusterName,
},
{
Name: "PGROOT",
Value: "/home/postgres/pgdata/pgroot",
},
{
Name: "ETCD_HOST",
Value: spilo.Spec.EtcdHost,
},
{
Name: "POD_IP",
ValueFrom: &v1.EnvVarSource{
FieldRef: &v1.ObjectFieldSelector{
APIVersion: "v1",
FieldPath: "status.podIP",
},
},
},
{
Name: "POD_NAMESPACE",
ValueFrom: &v1.EnvVarSource{
FieldRef: &v1.ObjectFieldSelector{
APIVersion: "v1",
FieldPath: "metadata.namespace",
},
},
},
{
Name: "PGPASSWORD_SUPERUSER",
ValueFrom: &v1.EnvVarSource{
SecretKeyRef: &v1.SecretKeySelector{
LocalObjectReference: v1.LocalObjectReference{
Name: clusterName,
},
Key: "superuser-password",
},
},
},
{
Name: "PGPASSWORD_ADMIN",
ValueFrom: &v1.EnvVarSource{
SecretKeyRef: &v1.SecretKeySelector{
LocalObjectReference: v1.LocalObjectReference{
Name: clusterName,
},
Key: "admin-password",
},
},
},
{
Name: "PGPASSWORD_STANDBY",
ValueFrom: &v1.EnvVarSource{
SecretKeyRef: &v1.SecretKeySelector{
LocalObjectReference: v1.LocalObjectReference{
Name: clusterName,
},
Key: "replication-password",
},
},
},
}
resourceList := v1.ResourceList{}
if (*spilo).Spec.ResourceCPU != "" {
resourceList[v1.ResourceCPU] = resource.MustParse((*spilo).Spec.ResourceCPU)
}
if (*spilo).Spec.ResourceMemory != "" {
resourceList[v1.ResourceMemory] = resource.MustParse((*spilo).Spec.ResourceMemory)
}
container := v1.Container{
Name: clusterName,
Image: spilo.Spec.DockerImage,
ImagePullPolicy: v1.PullAlways,
Resources: v1.ResourceRequirements{
Requests: resourceList,
},
Ports: []v1.ContainerPort{
{
ContainerPort: 8008,
Protocol: v1.ProtocolTCP,
},
{
ContainerPort: 5432,
Protocol: v1.ProtocolTCP,
},
},
VolumeMounts: []v1.VolumeMount{
{
Name: "pgdata",
MountPath: "/home/postgres/pgdata",
},
},
Env: envVars,
}
terminateGracePeriodSeconds := int64(30)
podSpec := v1.PodSpec{
TerminationGracePeriodSeconds: &terminateGracePeriodSeconds,
Volumes: []v1.Volume{
{
Name: "pgdata",
VolumeSource: v1.VolumeSource{EmptyDir: &v1.EmptyDirVolumeSource{}},
},
},
Containers: []v1.Container{container},
}
template := v1.PodTemplateSpec{
ObjectMeta: v1.ObjectMeta{
Labels: map[string]string{
"application": "spilo",
"spilo-cluster": clusterName,
},
Annotations: map[string]string{"pod.alpha.kubernetes.io/initialized": "true"},
},
Spec: podSpec,
}
return v1beta1.StatefulSet{
ObjectMeta: v1.ObjectMeta{
Name: clusterName,
Labels: map[string]string{
"application": "spilo",
"spilo-cluster": clusterName,
},
},
Spec: v1beta1.StatefulSetSpec{
Replicas: &spilo.Spec.NumberOfInstances,
ServiceName: clusterName,
Template: template,
},
}
}
func (z *SpiloSupervisor) CreateSecrets(ns, name string) {
secret := v1.Secret{
ObjectMeta: v1.ObjectMeta{
Name: name,
Labels: map[string]string{
"application": "spilo",
"spilo-cluster": name,
},
},
Type: v1.SecretTypeOpaque,
Data: map[string][]byte{
"superuser-password": []byte("emFsYW5kbw=="),
"replication-password": []byte("cmVwLXBhc3M="),
"admin-password": []byte("YWRtaW4="),
},
}
_, err := z.Clientset.Secrets(ns).Create(&secret)
if err != nil {
log.Printf("Secret error: %+v", err)
} else {
log.Printf("Secret created: %+v", secret)
}
}
func (z *SpiloSupervisor) CreateService(ns, name string) {
service := v1.Service{
ObjectMeta: v1.ObjectMeta{
Name: name,
Labels: map[string]string{
"application": "spilo",
"spilo-cluster": name,
},
},
Spec: v1.ServiceSpec{
Type: v1.ServiceTypeClusterIP,
Ports: []v1.ServicePort{{Port: 5432, TargetPort: intstr.IntOrString{IntVal: 5432}}},
},
}
_, err := z.Clientset.Services(ns).Create(&service)
if err != nil {
log.Printf("Service error: %+v", err)
} else {
log.Printf("Service created: %+v", service)
}
}
func (z *SpiloSupervisor) CreateEndPoint(ns, name string) {
endPoint := v1.Endpoints{
ObjectMeta: v1.ObjectMeta{
Name: name,
Labels: map[string]string{
"application": "spilo",
"spilo-cluster": name,
},
},
}
_, err := z.Clientset.Endpoints(ns).Create(&endPoint)
if err != nil {
log.Printf("Endpoint error: %+v", err)
} else {
log.Printf("Endpoint created: %+v", endPoint)
}
}
+137
View File
@@ -0,0 +1,137 @@
package controller
import (
"fmt"
"log"
"time"
"net/http"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/pkg/api"
apierrors "k8s.io/client-go/pkg/api/errors"
"k8s.io/client-go/pkg/api/unversioned"
"k8s.io/client-go/pkg/api/v1"
"k8s.io/client-go/pkg/apis/extensions/v1beta1"
"k8s.io/client-go/pkg/runtime"
"k8s.io/client-go/pkg/runtime/serializer"
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/clientcmd"
)
var (
etcdHostOutside string
VENDOR = "acid.zalan.do"
VERSION = "0.0.1.dev"
resyncPeriod = 5 * time.Minute
etcdKeyTemplate = "/service/%s"
)
type Options struct {
KubeConfig string
}
type Pgconf struct {
Parameter string `json:"param"`
Value string `json:"value"`
}
type SpiloSpec struct {
EtcdHost string `json:"etcd_host"`
VolumeSize int `json:"volume_size"`
NumberOfInstances int32 `json:"number_of_instances"`
DockerImage string `json:"docker_image"`
PostgresConfiguration []Pgconf `json:"postgres_configuration"`
ResourceCPU string `json:"resource_cpu"`
ResourceMemory string `json:"resource_memory"`
}
type Spilo struct {
unversioned.TypeMeta `json:",inline"`
Metadata api.ObjectMeta `json:"metadata"`
Spec SpiloSpec `json:"spec"`
}
type SpiloList struct {
unversioned.TypeMeta `json:",inline"`
Metadata unversioned.ListMeta `json:"metadata"`
Items []Spilo `json:"items"`
}
func KubernetesConfig(options Options) *rest.Config {
rules := clientcmd.NewDefaultClientConfigLoadingRules()
overrides := &clientcmd.ConfigOverrides{}
if options.KubeConfig != "" {
rules.ExplicitPath = options.KubeConfig
}
config, err := clientcmd.NewNonInteractiveDeferredLoadingClientConfig(rules, overrides).ClientConfig()
etcdHostOutside = config.Host
if err != nil {
log.Fatalf("Couldn't get Kubernetes default config: %s", err)
}
return config
}
func newKubernetesSpiloClient(c *rest.Config) (*rest.RESTClient, error) {
c.APIPath = "/apis"
c.GroupVersion = &unversioned.GroupVersion{
Group: VENDOR,
Version: "v1",
}
c.NegotiatedSerializer = serializer.DirectCodecFactory{CodecFactory: api.Codecs}
schemeBuilder := runtime.NewSchemeBuilder(
func(scheme *runtime.Scheme) error {
scheme.AddKnownTypes(
*c.GroupVersion,
&Spilo{},
&SpiloList{},
&api.ListOptions{},
&api.DeleteOptions{},
)
return nil
})
schemeBuilder.AddToScheme(api.Scheme)
return rest.RESTClientFor(c)
}
//TODO: Move to separate package
func IsKubernetesResourceNotFoundError(err error) bool {
se, ok := err.(*apierrors.StatusError)
if !ok {
return false
}
if se.Status().Code == http.StatusNotFound && se.Status().Reason == unversioned.StatusReasonNotFound {
return true
}
return false
}
func EnsureSpiloThirdPartyResource(client *kubernetes.Clientset) error {
// The resource doesn't exist, so we create it.
tpr := v1beta1.ThirdPartyResource{
ObjectMeta: v1.ObjectMeta{
Name: fmt.Sprintf("spilo.%s", VENDOR),
},
Description: "A specification of Spilo StatefulSets",
Versions: []v1beta1.APIVersion{
{Name: "v1"},
},
}
_, err := client.ExtensionsV1beta1().ThirdPartyResources().Create(&tpr)
if IsKubernetesResourceNotFoundError(err) {
return err
}
return nil
}
+109
View File
@@ -0,0 +1,109 @@
package controller
import (
"encoding/json"
"log"
"sync"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/pkg/api/meta"
"k8s.io/client-go/pkg/api/unversioned"
"k8s.io/client-go/rest"
"net/url"
"fmt"
"strings"
)
type SpiloOperator struct {
Options
ClientSet *kubernetes.Clientset
Client *rest.RESTClient
Supervisor *SpiloSupervisor
}
func New(options Options) *SpiloOperator {
config := KubernetesConfig(options)
clientSet, err := kubernetes.NewForConfig(config)
if err != nil {
log.Fatalf("Couldn't create Kubernetes client: %s", err)
}
etcdService, _ := clientSet.Services("default").Get("etcd-client")
if len(etcdService.Spec.Ports) != 1 {
log.Fatalln("Can't find Etcd cluster")
}
ports := etcdService.Spec.Ports[0]
nodeurl, _ := url.Parse(config.Host)
etcdHostOutside = fmt.Sprintf("http://%s:%d", strings.Split(nodeurl.Host, ":")[0], ports.NodePort)
spiloClient, err := newKubernetesSpiloClient(config)
if err != nil {
log.Fatalf("Couldn't create Spilo client: %s", err)
}
operator := &SpiloOperator{
Options: options,
ClientSet: clientSet,
Client: spiloClient,
Supervisor: newSupervisor(spiloClient, clientSet),
}
return operator
}
func (o *SpiloOperator) Run(stopCh <-chan struct{}, wg *sync.WaitGroup) {
log.Printf("Spilo operator %v\n", VERSION)
go o.Supervisor.Run(stopCh, wg)
log.Println("Started working in background")
}
// The code below is used only to work around a known problem with third-party
// resources and ugorji. If/when these issues are resolved, the code below
// should no longer be required.
//
func (s *Spilo) GetObjectKind() unversioned.ObjectKind {
return &s.TypeMeta
}
func (s *Spilo) GetObjectMeta() meta.Object {
return &s.Metadata
}
func (sl *SpiloList) GetObjectKind() unversioned.ObjectKind {
return &sl.TypeMeta
}
func (sl *SpiloList) GetListMeta() unversioned.List {
return &sl.Metadata
}
type SpiloListCopy SpiloList
type SpiloCopy Spilo
func (e *Spilo) UnmarshalJSON(data []byte) error {
tmp := SpiloCopy{}
err := json.Unmarshal(data, &tmp)
if err != nil {
return err
}
tmp2 := Spilo(tmp)
*e = tmp2
return nil
}
func (el *SpiloList) UnmarshalJSON(data []byte) error {
tmp := SpiloListCopy{}
err := json.Unmarshal(data, &tmp)
if err != nil {
return err
}
tmp2 := SpiloList(tmp)
*el = tmp2
return nil
}
+340
View File
@@ -0,0 +1,340 @@
package controller
import (
"fmt"
"log"
"sync"
"time"
etcdclient "github.com/coreos/etcd/client"
"k8s.io/client-go/kubernetes"
"k8s.io/client-go/pkg/api"
"k8s.io/client-go/pkg/api/v1"
"k8s.io/client-go/pkg/fields"
"k8s.io/client-go/rest"
"k8s.io/client-go/tools/cache"
)
const (
ACTION_DELETE = "delete"
ACTION_UPDATE = "update"
ACTION_ADD = "add"
)
type podEvent struct {
namespace string
name string
actionType string
}
type podWatcher struct {
podNamespace string
podName string
eventsChannel chan podEvent
subscribe bool
}
type SpiloSupervisor struct {
podEvents chan podEvent
podWatchers chan podWatcher
SpiloClient *rest.RESTClient
Clientset *kubernetes.Clientset
spiloInformer cache.SharedIndexInformer
podInformer cache.SharedIndexInformer
etcdApiClient etcdclient.KeysAPI
}
func podsListWatch(client *kubernetes.Clientset) *cache.ListWatch {
return cache.NewListWatchFromClient(client.Core().RESTClient(), "pods", api.NamespaceAll, fields.Everything())
}
func newSupervisor(spiloClient *rest.RESTClient, clientset *kubernetes.Clientset) *SpiloSupervisor {
spiloSupervisor := &SpiloSupervisor{
SpiloClient: spiloClient,
Clientset: clientset,
}
spiloInformer := cache.NewSharedIndexInformer(
cache.NewListWatchFromClient(spiloClient, "spilos", api.NamespaceAll, fields.Everything()),
&Spilo{},
resyncPeriod,
cache.Indexers{cache.NamespaceIndex: cache.MetaNamespaceIndexFunc},
)
spiloInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: spiloSupervisor.spiloAdd,
UpdateFunc: spiloSupervisor.spiloUpdate,
DeleteFunc: spiloSupervisor.spiloDelete,
})
podInformer := cache.NewSharedIndexInformer(
podsListWatch(clientset),
&v1.Pod{},
resyncPeriod,
cache.Indexers{cache.NamespaceIndex: cache.MetaNamespaceIndexFunc},
)
podInformer.AddEventHandler(cache.ResourceEventHandlerFuncs{
AddFunc: spiloSupervisor.podAdd,
UpdateFunc: spiloSupervisor.podUpdate,
DeleteFunc: spiloSupervisor.podDelete,
})
spiloSupervisor.spiloInformer = spiloInformer
spiloSupervisor.podInformer = podInformer
cfg := etcdclient.Config{
Endpoints: []string{etcdHostOutside},
Transport: etcdclient.DefaultTransport,
HeaderTimeoutPerRequest: time.Second,
}
c, err := etcdclient.New(cfg)
if err != nil {
log.Fatal(err)
}
spiloSupervisor.etcdApiClient = etcdclient.NewKeysAPI(c)
spiloSupervisor.podEvents = make(chan podEvent)
return spiloSupervisor
}
func (d *SpiloSupervisor) podAdd(obj interface{}) {
pod := obj.(*v1.Pod)
d.podEvents <- podEvent{
namespace: pod.Namespace,
name: pod.Name,
actionType: ACTION_ADD,
}
}
func (d *SpiloSupervisor) podDelete(obj interface{}) {
pod := obj.(*v1.Pod)
d.podEvents <- podEvent{
namespace: pod.Namespace,
name: pod.Name,
actionType: ACTION_DELETE,
}
}
func (d *SpiloSupervisor) podUpdate(old, cur interface{}) {
oldPod := old.(*v1.Pod)
d.podEvents <- podEvent{
namespace: oldPod.Namespace,
name: oldPod.Name,
actionType: ACTION_UPDATE,
}
}
func (z *SpiloSupervisor) Run(stopCh <-chan struct{}, wg *sync.WaitGroup) {
defer wg.Done()
wg.Add(1)
if err := EnsureSpiloThirdPartyResource(z.Clientset); err != nil {
log.Fatalf("Couldn't create ThirdPartyResource: %s", err)
}
go z.spiloInformer.Run(stopCh)
go z.podInformer.Run(stopCh)
go z.podWatcher(stopCh)
<-stopCh
}
func (z *SpiloSupervisor) spiloAdd(obj interface{}) {
spilo := obj.(*Spilo)
clusterName := (*spilo).Metadata.Name
ns := (*spilo).Metadata.Namespace
//TODO: check if object already exists before creating
z.CreateEndPoint(ns, clusterName)
z.CreateService(ns, clusterName)
z.CreateSecrets(ns, clusterName)
z.CreateStatefulSet(spilo)
}
func (z *SpiloSupervisor) spiloUpdate(old, cur interface{}) {
oldSpilo := old.(*Spilo)
curSpilo := cur.(*Spilo)
if oldSpilo.Spec.NumberOfInstances != curSpilo.Spec.NumberOfInstances {
z.UpdateStatefulSet(curSpilo)
}
if oldSpilo.Spec.DockerImage != curSpilo.Spec.DockerImage {
log.Printf("Updating DockerImage: %s.%s",
curSpilo.Metadata.Namespace,
curSpilo.Metadata.Name)
z.UpdateStatefulSetImage(curSpilo)
}
log.Printf("Update spilo old: %+v\ncurrent: %+v", *oldSpilo, *curSpilo)
}
func (z *SpiloSupervisor) spiloDelete(obj interface{}) {
spilo := obj.(*Spilo)
err := z.DeleteStatefulSet(spilo.Metadata.Namespace, spilo.Metadata.Name)
if err != nil {
log.Printf("Error while deleting stateful set: %+v", err)
}
}
func (z *SpiloSupervisor) DeleteStatefulSet(ns, clusterName string) error {
orphanDependents := false
deleteOptions := v1.DeleteOptions{
OrphanDependents: &orphanDependents,
}
listOptions := v1.ListOptions{
LabelSelector: fmt.Sprintf("%s=%s", "spilo-cluster", clusterName),
}
podList, err := z.Clientset.Pods(ns).List(listOptions)
if err != nil {
log.Printf("Error: %+v", err)
}
err = z.Clientset.StatefulSets(ns).Delete(clusterName, &deleteOptions)
if err != nil {
return err
}
log.Printf("StatefulSet %s.%s has been deleted\n", ns, clusterName)
for _, pod := range podList.Items {
err = z.Clientset.Pods(pod.Namespace).Delete(pod.Name, &deleteOptions)
if err != nil {
log.Printf("Error while deleting Pod %s: %+v", pod.Name, err)
return err
}
log.Printf("Pod %s.%s has been deleted\n", pod.Namespace, pod.Name)
}
serviceList, err := z.Clientset.Services(ns).List(listOptions)
if err != nil {
return err
}
for _, service := range serviceList.Items {
err = z.Clientset.Services(service.Namespace).Delete(service.Name, &deleteOptions)
if err != nil {
log.Printf("Error while deleting Service %s: %+v", service.Name, err)
return err
}
log.Printf("Service %s.%s has been deleted\n", service.Namespace, service.Name)
}
z.DeleteEtcdKey(clusterName)
return nil
}
func (z *SpiloSupervisor) UpdateStatefulSet(spilo *Spilo) {
ns := (*spilo).Metadata.Namespace
statefulSet := z.createSetFromSpilo(spilo)
_, err := z.Clientset.StatefulSets(ns).Update(&statefulSet)
if err != nil {
log.Printf("Error while updating StatefulSet: %s", err)
}
}
func (z *SpiloSupervisor) UpdateStatefulSetImage(spilo *Spilo) {
ns := (*spilo).Metadata.Namespace
z.UpdateStatefulSet(spilo)
listOptions := v1.ListOptions{
LabelSelector: fmt.Sprintf("%s=%s", "spilo-cluster", (*spilo).Metadata.Name),
}
pods, err := z.Clientset.Pods(ns).List(listOptions)
if err != nil {
log.Printf("Error while getting pods: %s", err)
}
orphanDependents := true
deleteOptions := v1.DeleteOptions{
OrphanDependents: &orphanDependents,
}
var masterPodName string
for _, pod := range pods.Items {
log.Printf("Pod processing: %s", pod.Name)
role, ok := pod.Labels["spilo-role"]
if ok == false {
log.Println("No spilo-role label")
continue
}
if role == "master" {
masterPodName = pod.Name
log.Printf("Skipping master: %s", masterPodName)
continue
}
err := z.Clientset.Pods(ns).Delete(pod.Name, &deleteOptions)
if err != nil {
log.Printf("Error while deleting Pod %s.%s: %s", pod.Namespace, pod.Name, err)
} else {
log.Printf("Pod deleted: %s.%s", pod.Namespace, pod.Name)
}
w1 := podWatcher{
podNamespace: pod.Namespace,
podName: pod.Name,
eventsChannel: make(chan podEvent, 1),
subscribe: true,
}
log.Printf("Watching pod %s.%s being recreated", pod.Namespace, pod.Name)
z.podWatchers <- w1
for e := range w1.eventsChannel {
if e.actionType == ACTION_ADD { break }
}
log.Printf("Pod %s.%s has been recreated", pod.Namespace, pod.Name)
}
//TODO: do manual failover
err = z.Clientset.Pods(ns).Delete(masterPodName, &deleteOptions)
if err != nil {
log.Printf("Error while deleting Pod %s.%s: %s", ns, masterPodName, err)
} else {
log.Printf("Pod deleted: %s.%s", ns, masterPodName)
}
}
func (z *SpiloSupervisor) podWatcher(stopCh <-chan struct{}) {
//TODO: mind the namespace of the pod
watchers := make(map[string] podWatcher)
for {
select {
case watcher := <-z.podWatchers:
if watcher.subscribe {
watchers[watcher.podName] = watcher
} else {
close(watcher.eventsChannel)
delete(watchers, watcher.podName)
}
case event := <-z.podEvents:
log.Printf("Pod watcher event: %s.%s - %s", event.namespace, event.name, event.actionType)
log.Printf("Current watchers: %+v", watchers)
podWatcher, ok := watchers[event.name]
if ok == false {
continue
}
podWatcher.eventsChannel <- event
}
}
}