golint cleanup for the operator

This commit is contained in:
Travis Nielsen
2017-07-24 17:24:24 -07:00
parent 88905f279b
commit f3640887f3
37 changed files with 214 additions and 94 deletions
+9 -2
View File
@@ -13,6 +13,8 @@ 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 api for the operator api manager.
package api
import (
@@ -34,6 +36,7 @@ import (
var logger = capnslog.NewPackageLogger("github.com/rook/rook", "op-api")
const (
// DeploymentName is the service name
DeploymentName = "rook-api"
)
@@ -55,6 +58,7 @@ var clusterAccessRules = []v1beta1.PolicyRule{
},
}
// Cluster has the api service properties
type Cluster struct {
context *clusterd.Context
Namespace string
@@ -63,6 +67,7 @@ type Cluster struct {
Replicas int32
}
// New creates an instance
func New(context *clusterd.Context, namespace, version string, placement k8sutil.Placement) *Cluster {
return &Cluster{
context: context,
@@ -73,6 +78,7 @@ func New(context *clusterd.Context, namespace, version string, placement k8sutil
}
}
// Start the api service
func (c *Cluster) Start() error {
logger.Infof("starting the Rook api")
@@ -103,6 +109,7 @@ func (c *Cluster) Start() error {
return nil
}
// make a cluster role
func (c *Cluster) makeClusterRole() error {
account := &v1.ServiceAccount{}
account.Name = DeploymentName
@@ -186,9 +193,9 @@ func (c *Cluster) apiContainer() v1.Container {
{Name: "ROOKD_VERSION_TAG", Value: c.Version},
k8sutil.NamespaceEnvVar(),
k8sutil.RepoPrefixEnvVar(),
opmon.MonSecretEnvVar(),
opmon.SecretEnvVar(),
opmon.AdminSecretEnvVar(),
opmon.MonEndpointEnvVar(),
opmon.EndpointEnvVar(),
opmon.ClusterNameEnvVar(c.Namespace),
},
}
+3 -1
View File
@@ -16,6 +16,8 @@ limitations under the License.
Some of the code below came from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package operator to manage Kubernetes storage.
package operator
import (
@@ -99,7 +101,7 @@ func (m *clusterManager) startTrack(c *cluster.Cluster) error {
existing, ok := m.clusters[c.Namespace]
if ok {
if c.Name != existing.Name {
return fmt.Errorf("cluster %s is already running in namespace %s. Multiple clusters per namespace not supported.", existing.Name, existing.Namespace)
return fmt.Errorf("cluster %s is already running in namespace %s. Multiple clusters per namespace not supported", existing.Name, existing.Namespace)
}
} else {
// only start the cluster if we're not already tracking it from a previous iteration
+9 -1
View File
@@ -16,6 +16,8 @@ limitations under the License.
Some of the code below came from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package cluster to manage a rook cluster.
package cluster
import (
@@ -50,7 +52,8 @@ var (
healthCheckInterval = 10 * time.Second
clientTimeout = 15 * time.Second
logger = capnslog.NewPackageLogger("github.com/rook/rook", "op-cluster")
ClusterResource = kit.CustomResource{
// ClusterResource is the definition of the cluster CRD
ClusterResource = kit.CustomResource{
Name: "cluster",
Group: k8sutil.CustomResourceGroup,
Version: kit.V1Alpha1,
@@ -58,6 +61,7 @@ var (
}
)
// Cluster controls an instance of a Rook cluster
type Cluster struct {
context *clusterd.Context
v1.ObjectMeta `json:"metadata,omitempty"`
@@ -70,10 +74,12 @@ type Cluster struct {
rclient rookclient.RookRestClient
}
// Init assigns the cluster context
func (c *Cluster) Init(context *clusterd.Context) {
c.context = context
}
// CreateInstance creates a new Rook cluster instance
func (c *Cluster) CreateInstance() error {
// Create the namespace if not already created
@@ -137,6 +143,7 @@ func (c *Cluster) CreateInstance() error {
return nil
}
// Monitor watches a cluster for failures and restarts failed components
func (c *Cluster) Monitor(stopCh <-chan struct{}) {
for {
select {
@@ -252,6 +259,7 @@ func (c *Cluster) createClientAccess(clusterInfo *cephmon.ClusterInfo) error {
return nil
}
// GetRookClient gets the REST api client
func (c *Cluster) GetRookClient() (rookclient.RookRestClient, error) {
if c.rclient != nil {
return c.rclient, nil
+7
View File
@@ -16,6 +16,8 @@ limitations under the License.
Some of the code below came from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package cluster to manage a rook cluster.
package cluster
import (
@@ -40,9 +42,13 @@ type ClusterList struct {
// - We include `Metadata` field in object explicitly.
// - we have the code below to work around a known problem with third-party resources and ugorji.
// ClusterListCopy is for deserialization
type ClusterListCopy ClusterList
// ClusterCopy is for deserializaion
type ClusterCopy Cluster
// UnmarshalJSON unmarshals the raw cluster
func (c *Cluster) UnmarshalJSON(data []byte) error {
tmp := ClusterCopy{}
err := json.Unmarshal(data, &tmp)
@@ -54,6 +60,7 @@ func (c *Cluster) UnmarshalJSON(data []byte) error {
return nil
}
// UnmarshalJSON unmarshals the raw cluster list
func (cl *ClusterList) UnmarshalJSON(data []byte) error {
tmp := ClusterListCopy{}
err := json.Unmarshal(data, &tmp)
+6 -2
View File
@@ -16,6 +16,8 @@ limitations under the License.
Some of the code below came from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package cluster to manage a rook cluster.
package cluster
import (
@@ -34,6 +36,7 @@ const (
erasureCodeType = "erasure-coded"
)
// PoolResource is the definition of the pool CRD
var PoolResource = kit.CustomResource{
Name: "pool",
Group: k8sutil.CustomResourceGroup,
@@ -41,17 +44,18 @@ var PoolResource = kit.CustomResource{
Description: "Managed Rook pools",
}
// Pool is the spec for the CRD
type Pool struct {
v1.ObjectMeta `json:"metadata,omitempty"`
PoolSpec `json:"spec"`
}
// Instantiate a new pool
// NewPool creates a new pool
func NewPool(spec PoolSpec) *Pool {
return &Pool{PoolSpec: spec}
}
// Create the pool
// Create a pool
func (p *Pool) Create(rclient rookclient.RookRestClient) error {
// validate the pool settings
if err := p.validate(); err != nil {
+7
View File
@@ -16,6 +16,8 @@ limitations under the License.
Some of the code below came from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package cluster to manage a rook cluster.
package cluster
import (
@@ -40,9 +42,13 @@ type PoolList struct {
// - We include `Metadata` field in object explicitly.
// - we have the code below to work around a known problem with third-party resources and ugorji.
// PoolListCopy is for deserialization
type PoolListCopy PoolList
// PoolCopy is for deserialization
type PoolCopy Pool
// UnmarshalJSON deserializes the pool
func (p *Pool) UnmarshalJSON(data []byte) error {
tmp := PoolCopy{}
err := json.Unmarshal(data, &tmp)
@@ -54,6 +60,7 @@ func (p *Pool) UnmarshalJSON(data []byte) error {
return nil
}
// UnmarshalJSON deserializes the pool list
func (pl *PoolList) UnmarshalJSON(data []byte) error {
tmp := PoolListCopy{}
err := json.Unmarshal(data, &tmp)
+15
View File
@@ -16,6 +16,8 @@ limitations under the License.
Some of the code below came from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package cluster to manage a rook cluster.
package cluster
import (
@@ -23,6 +25,7 @@ import (
"github.com/rook/rook/pkg/operator/osd"
)
// Spec for the cluster
type Spec struct {
// VersionTag is the expected version of the rook container to run in the cluster.
// The operator will eventually make the rook cluster version
@@ -42,6 +45,7 @@ type Spec struct {
Storage osd.StorageSpec `json:"storage"`
}
// PoolSpec is the specific spec for the redundancy
type PoolSpec struct {
// The replication settings
Replication ReplicationSpec `json:"replication"`
@@ -50,11 +54,13 @@ type PoolSpec struct {
ErasureCoding ErasureCodeSpec `json:"erasureCode"`
}
// ReplicationSpec specifies the number of replicas
type ReplicationSpec struct {
// Number of copies per object in a replicated storage pool, including the object itself (required for replicated pool type)
Size uint `json:"size"`
}
// ErasureCodeSpec specifies the erasure coding params
type ErasureCodeSpec struct {
// Number of coding chunks per object in an erasure coded storage pool (required for erasure-coded pool type)
CodingChunks uint `json:"codingChunks"`
@@ -73,8 +79,17 @@ type PlacementSpec struct {
RGW k8sutil.Placement `json:"rgw,omitempty"`
}
// GetAPI returns the placement for the API service
func (p PlacementSpec) GetAPI() k8sutil.Placement { return p.All.Merge(p.API) }
// GetMDS returns the placement for the MDS service
func (p PlacementSpec) GetMDS() k8sutil.Placement { return p.All.Merge(p.MDS) }
// GetMON returns the placement for the MON service
func (p PlacementSpec) GetMON() k8sutil.Placement { return p.All.Merge(p.MON) }
// GetOSD returns the placement for the OSD service
func (p PlacementSpec) GetOSD() k8sutil.Placement { return p.All.Merge(p.OSD) }
// GetRGW returns the placement for the RGW service
func (p PlacementSpec) GetRGW() k8sutil.Placement { return p.All.Merge(p.RGW) }
+15 -6
View File
@@ -16,6 +16,8 @@ limitations under the License.
Some of the code below came from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package k8sutil for Kubernetes helpers.
package k8sutil
import "github.com/coreos/pkg/capnslog"
@@ -23,11 +25,18 @@ import "github.com/coreos/pkg/capnslog"
var logger = capnslog.NewPackageLogger("github.com/rook/rook", "op-k8sutil")
const (
Namespace = "rook"
// Namespace for rook
Namespace = "rook"
// CustomResourceGroup for rook CRD
CustomResourceGroup = "rook.io"
DefaultNamespace = "default"
DataDirVolume = "rook-data"
DataDir = "/var/lib/rook"
RookType = "kubernetes.io/rook"
RbdType = "kubernetes.io/rbd"
// DefaultNamespace for the cluster
DefaultNamespace = "default"
// DataDirVolume data dir volume
DataDirVolume = "rook-data"
// DataDir folder
DataDir = "/var/lib/rook"
// RookType for the CRD
RookType = "kubernetes.io/rook"
// RbdType for the RBD mounts
RbdType = "kubernetes.io/rbd"
)
+3
View File
@@ -13,6 +13,8 @@ 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 k8sutil for Kubernetes helpers.
package k8sutil
import (
@@ -26,6 +28,7 @@ type Placement struct {
Tolerations []v1.Toleration `json:"tolerations,omitemtpy"`
}
// ApplyToPodSpec adds placement to a pod spec
func (p Placement) ApplyToPodSpec(t *v1.PodSpec) {
if p.NodeAffinity != nil {
if t.Affinity == nil {
+2 -2
View File
@@ -39,12 +39,12 @@ tolerations:
operator: Exists`)
// convert the raw spec yaml into JSON
rawJson, err := yaml.YAMLToJSON(specYaml)
rawJSON, err := yaml.YAMLToJSON(specYaml)
assert.Nil(t, err)
// unmarshal the JSON into a strongly typed placement spec object
var placement Placement
err = json.Unmarshal(rawJson, &placement)
err = json.Unmarshal(rawJSON, &placement)
assert.Nil(t, err)
// the unmarshalled placement spec should equal the expected spec below
+27 -18
View File
@@ -16,6 +16,8 @@ limitations under the License.
Some of the code below came from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package k8sutil for Kubernetes helpers.
package k8sutil
import (
@@ -27,41 +29,54 @@ import (
)
const (
AppAttr = "app"
ClusterAttr = "rook_cluster"
VersionAttr = "rook_version"
PodIPEnvVar = "ROOKD_PRIVATE_IPV4"
DefaultRepoPrefix = "rook"
repoPrefixEnvVar = "ROOKD_REPO_PREFIX"
defaultVersion = "latest"
// AppAttr app label
AppAttr = "app"
// ClusterAttr cluster label
ClusterAttr = "rook_cluster"
// VersionAttr version label
VersionAttr = "rook_version"
// PodIPEnvVar pod IP env var
PodIPEnvVar = "ROOKD_PRIVATE_IPV4"
// DefaultRepoPrefix repo prefix
DefaultRepoPrefix = "rook"
// ConfigOverrideName config override name
ConfigOverrideName = "rook-config-override"
ConfigOverrideVal = "config"
configMountDir = "/etc/rook"
overrideFilename = "override.conf"
// ConfigOverrideVal config override value
ConfigOverrideVal = "config"
repoPrefixEnvVar = "ROOKD_REPO_PREFIX"
defaultVersion = "latest"
configMountDir = "/etc/rook"
overrideFilename = "override.conf"
)
// ConfigOverrideMount is an override mount
func ConfigOverrideMount() v1.VolumeMount {
return v1.VolumeMount{Name: ConfigOverrideName, MountPath: configMountDir}
}
// ConfigOverrideVolume is an override volume
func ConfigOverrideVolume() v1.Volume {
cmSource := &v1.ConfigMapVolumeSource{Items: []v1.KeyToPath{{Key: ConfigOverrideVal, Path: overrideFilename}}}
cmSource.Name = ConfigOverrideName
return v1.Volume{Name: ConfigOverrideName, VolumeSource: v1.VolumeSource{ConfigMap: cmSource}}
}
// ConfigOverrideEnvVar config override env var
func ConfigOverrideEnvVar() v1.EnvVar {
return v1.EnvVar{Name: "ROOKD_CEPH_CONFIG_OVERRIDE", Value: path.Join(configMountDir, overrideFilename)}
}
// NamespaceEnvVar namespace env var
func NamespaceEnvVar() v1.EnvVar {
return v1.EnvVar{Name: "ROOKD_NAMESPACE", ValueFrom: &v1.EnvVarSource{FieldRef: &v1.ObjectFieldSelector{FieldPath: "metadata.namespace"}}}
}
// RepoPrefixEnvVar repo prefix env var
func RepoPrefixEnvVar() v1.EnvVar {
return v1.EnvVar{Name: repoPrefixEnvVar, Value: repoPrefix()}
}
// ConfigDirEnvVar config dir env var
func ConfigDirEnvVar() v1.EnvVar {
return v1.EnvVar{Name: "ROOKD_CONFIG_DIR", Value: DataDir}
}
@@ -82,18 +97,12 @@ func getVersion(version string) string {
return version
}
// MakeRookImage formats the container name
func MakeRookImage(version string) string {
return fmt.Sprintf("%s/rook:%v", repoPrefix(), getVersion(version))
}
// SetPodVersion sets the pod annotation
func SetPodVersion(pod *v1.Pod, key, version string) {
pod.Annotations[key] = version
}
func GetPodNames(pods []*v1.Pod) []string {
res := []string{}
for _, p := range pods {
res = append(res, p.Name)
}
return res
}
+3
View File
@@ -13,6 +13,8 @@ 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 k8sutil for Kubernetes helpers.
package k8sutil
import (
@@ -20,6 +22,7 @@ import (
"strings"
)
// PathToVolumeName converts a path to a valid volume name
func PathToVolumeName(path string) string {
// kubernetes volume names must match this regex: [a-z0-9]([-a-z0-9]*[a-z0-9])?
+2 -2
View File
@@ -1,6 +1,4 @@
/*
Package kit for Kubernetes operators
Copyright 2016 The Rook Authors. All rights reserved.
Licensed under the Apache License, Version 2.0 (the "License");
@@ -18,6 +16,8 @@ limitations under the License.
Some of the code below came from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package kit for Kubernetes operators
package kit
import (
+2 -2
View File
@@ -1,6 +1,4 @@
/*
Package kit for Kubernetes operators
Copyright 2016 The Rook Authors. All rights reserved.
Licensed under the Apache License, Version 2.0 (the "License");
@@ -18,6 +16,8 @@ limitations under the License.
Some of the code below came from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package kit for Kubernetes operators
package kit
import "github.com/coreos/pkg/capnslog"
+2 -2
View File
@@ -1,6 +1,4 @@
/*
Package kit for Kubernetes operators
Copyright 2016 The Rook Authors. All rights reserved.
Licensed under the Apache License, Version 2.0 (the "License");
@@ -18,6 +16,8 @@ limitations under the License.
Some of the code was modified from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package kit for Kubernetes operators
package kit
import (
+2 -2
View File
@@ -1,6 +1,4 @@
/*
Package kit for Kubernetes operators
Copyright 2016 The Rook Authors. All rights reserved.
Licensed under the Apache License, Version 2.0 (the "License");
@@ -18,6 +16,8 @@ limitations under the License.
Some of the code below came from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package kit for Kubernetes operators
package kit
import (
+2 -2
View File
@@ -1,6 +1,4 @@
/*
Package kit for Kubernetes operators
Copyright 2016 The Rook Authors. All rights reserved.
Licensed under the Apache License, Version 2.0 (the "License");
@@ -18,6 +16,8 @@ limitations under the License.
Some of the code below came from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package kit for Kubernetes operators
package kit
import "time"
+2 -2
View File
@@ -1,6 +1,4 @@
/*
Package kit for Kubernetes operators
Copyright 2016 The Rook Authors. All rights reserved.
Licensed under the Apache License, Version 2.0 (the "License");
@@ -18,6 +16,8 @@ limitations under the License.
Some of the code was modified from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package kit for Kubernetes operators
package kit
import (
-20
View File
@@ -1,20 +0,0 @@
/*
Copyright 2016 The Rook Authors. All rights reserved.
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 operator
import "github.com/coreos/pkg/capnslog"
var logger = capnslog.NewPackageLogger("github.com/rook/rook", "operator")
+7 -2
View File
@@ -13,6 +13,8 @@ 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 mds for file systems.
package mds
import (
@@ -39,6 +41,7 @@ const (
keyringName = "keyring"
)
// Cluster for mds management
type Cluster struct {
Namespace string
Version string
@@ -48,6 +51,7 @@ type Cluster struct {
placement k8sutil.Placement
}
// New creates an instance of the mds manager
func New(context *clusterd.Context, namespace, version string, placement k8sutil.Placement) *Cluster {
return &Cluster{
context: context,
@@ -59,6 +63,7 @@ func New(context *clusterd.Context, namespace, version string, placement k8sutil
}
}
// Start the mds manager
func (c *Cluster) Start() error {
logger.Infof("start running mds")
@@ -162,8 +167,8 @@ func (c *Cluster) mdsContainer(id string) v1.Container {
Env: []v1.EnvVar{
{Name: "ROOKD_MDS_KEYRING", ValueFrom: &v1.EnvVarSource{SecretKeyRef: &v1.SecretKeySelector{LocalObjectReference: v1.LocalObjectReference{Name: appName}, Key: keyringName}}},
opmon.ClusterNameEnvVar(c.Namespace),
opmon.MonEndpointEnvVar(),
opmon.MonSecretEnvVar(),
opmon.EndpointEnvVar(),
opmon.SecretEnvVar(),
opmon.AdminSecretEnvVar(),
k8sutil.ConfigOverrideEnvVar(),
},
+7 -2
View File
@@ -13,6 +13,8 @@ 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 mgr for the Ceph manager.
package mgr
import (
@@ -36,6 +38,7 @@ const (
keyringName = "keyring"
)
// Cluster is the ceph mgr manager
type Cluster struct {
Namespace string
Version string
@@ -44,6 +47,7 @@ type Cluster struct {
dataDir string
}
// New creates an instance of the mgr
func New(context *clusterd.Context, namespace, version string) *Cluster {
return &Cluster{
context: context,
@@ -54,6 +58,7 @@ func New(context *clusterd.Context, namespace, version string) *Cluster {
}
}
// Start the mgr instance
func (c *Cluster) Start() error {
logger.Infof("start running mgr")
@@ -123,8 +128,8 @@ func (c *Cluster) mgrContainer(name string) v1.Container {
{Name: "ROOKD_MGR_NAME", Value: name},
{Name: "ROOKD_MGR_KEYRING", ValueFrom: &v1.EnvVarSource{SecretKeyRef: &v1.SecretKeySelector{LocalObjectReference: v1.LocalObjectReference{Name: name}, Key: keyringName}}},
opmon.ClusterNameEnvVar(c.Namespace),
opmon.MonEndpointEnvVar(),
opmon.MonSecretEnvVar(),
opmon.EndpointEnvVar(),
opmon.SecretEnvVar(),
opmon.AdminSecretEnvVar(),
k8sutil.ConfigOverrideEnvVar(),
},
+4 -1
View File
@@ -13,6 +13,8 @@ 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 mon for the Ceph monitors.
package mon
import (
@@ -24,6 +26,7 @@ import (
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
)
// CheckHealth for the monitors
func (c *Cluster) CheckHealth() error {
logger.Debugf("Checking health for mons. %+v", c.clusterInfo)
@@ -68,7 +71,7 @@ func (c *Cluster) failoverMon(name string) error {
logger.Infof("Failing over monitor %s", name)
// Start a new monitor
mons := []*MonConfig{{Name: fmt.Sprintf("mon%d", c.maxMonID+1), Port: int32(mon.Port)}}
mons := []*monConfig{{Name: fmt.Sprintf("mon%d", c.maxMonID+1), Port: int32(mon.Port)}}
logger.Infof("starting new mon %s", mons[0].Name)
err := c.startPods(mons)
if err != nil {
+14 -8
View File
@@ -13,6 +13,8 @@ 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 mon for the Ceph monitors.
package mon
import (
@@ -51,6 +53,7 @@ const (
maxMonIDKey = "maxMonId"
)
// Cluster is for the cluster of monitors
type Cluster struct {
context *clusterd.Context
Namespace string
@@ -67,11 +70,13 @@ type Cluster struct {
dataDirHostPath string
}
type MonConfig struct {
// monConfig for a single monitor
type monConfig struct {
Name string
Port int32
}
// New creates an instance of a mon cluster
func New(context *clusterd.Context, namespace, dataDirHostPath, version string, placement k8sutil.Placement) *Cluster {
return &Cluster{
context: context,
@@ -85,6 +90,7 @@ func New(context *clusterd.Context, namespace, dataDirHostPath, version string,
}
}
// Start the mon cluster
func (c *Cluster) Start() (*mon.ClusterInfo, error) {
logger.Infof("start running mons")
@@ -158,18 +164,18 @@ func (c *Cluster) initClusterInfo() error {
return nil
}
func (c *Cluster) getExpectedMonConfig() []*MonConfig {
mons := []*MonConfig{}
func (c *Cluster) getExpectedMonConfig() []*monConfig {
mons := []*monConfig{}
// initialize the mon pod info for mons that have been previously created
for _, monitor := range c.clusterInfo.Monitors {
mons = append(mons, &MonConfig{Name: monitor.Name, Port: int32(mon.Port)})
mons = append(mons, &monConfig{Name: monitor.Name, Port: int32(mon.Port)})
}
// initialize mon info if we don't have enough mons (at first startup)
for i := len(c.clusterInfo.Monitors); i < c.Size; i++ {
c.maxMonID++
mons = append(mons, &MonConfig{Name: fmt.Sprintf("%s%d", appName, c.maxMonID), Port: int32(mon.Port)})
mons = append(mons, &monConfig{Name: fmt.Sprintf("%s%d", appName, c.maxMonID), Port: int32(mon.Port)})
}
return mons
@@ -234,7 +240,7 @@ func (c *Cluster) createMonSecretsAndSave() error {
return nil
}
func (c *Cluster) startPods(mons []*MonConfig) error {
func (c *Cluster) startPods(mons []*monConfig) error {
// schedule the mons on different nodes if we have enough nodes to be unique
availableNodes, err := c.getAvailableMonNodes()
if err != nil {
@@ -273,7 +279,7 @@ func (c *Cluster) startPods(mons []*MonConfig) error {
return c.waitForMonsToJoin(mons)
}
func (c *Cluster) waitForMonsToJoin(mons []*MonConfig) error {
func (c *Cluster) waitForMonsToJoin(mons []*monConfig) error {
if !c.waitForStart {
return nil
}
@@ -522,7 +528,7 @@ func (c *Cluster) getNodesWithMons() (*util.Set, error) {
return nodes, nil
}
func (c *Cluster) startMon(m *MonConfig, nodeName string) error {
func (c *Cluster) startMon(m *monConfig, nodeName string) error {
rs := c.makeReplicaSet(m, nodeName)
logger.Debugf("Starting mon: %+v", rs.Name)
_, err := c.context.Clientset.Extensions().ReplicaSets(c.Namespace).Create(rs)
+2 -2
View File
@@ -210,7 +210,7 @@ func TestAvailableNodesInUse(t *testing.T) {
// start pods on two of the nodes so that only one node will be available
for i := 0; i < 2; i++ {
pod := c.makeMonPod(&MonConfig{Name: fmt.Sprintf("rook-ceph-mon%d", i)}, nodes[i].Name)
pod := c.makeMonPod(&monConfig{Name: fmt.Sprintf("rook-ceph-mon%d", i)}, nodes[i].Name)
_, err := clientset.CoreV1().Pods(c.Namespace).Create(pod)
assert.Nil(t, err)
}
@@ -221,7 +221,7 @@ func TestAvailableNodesInUse(t *testing.T) {
// start pods on the remaining node. We expect all nodes to be available for placement
// since there is no way to place a mon on an unused node.
pod := c.makeMonPod(&MonConfig{Name: "mon2"}, nodes[2].Name)
pod := c.makeMonPod(&monConfig{Name: "mon2"}, nodes[2].Name)
_, err = clientset.CoreV1().Pods(c.Namespace).Create(pod)
assert.Nil(t, err)
nodes, err = c.getAvailableMonNodes()
+13 -7
View File
@@ -13,6 +13,8 @@ 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 mon for the Ceph monitors.
package mon
import (
@@ -25,20 +27,24 @@ import (
"k8s.io/kubernetes/pkg/kubelet/apis"
)
// ClusterNameEnvVar is the cluster name environment var
func ClusterNameEnvVar(name string) v1.EnvVar {
return v1.EnvVar{Name: "ROOKD_CLUSTER_NAME", Value: name}
}
func MonEndpointEnvVar() v1.EnvVar {
// EndpointEnvVar is the mon endpoint environment var
func EndpointEnvVar() v1.EnvVar {
ref := &v1.ConfigMapKeySelector{LocalObjectReference: v1.LocalObjectReference{Name: monConfigMapName}, Key: monEndpointKey}
return v1.EnvVar{Name: "ROOKD_MON_ENDPOINTS", ValueFrom: &v1.EnvVarSource{ConfigMapKeyRef: ref}}
}
func MonSecretEnvVar() v1.EnvVar {
// SecretEnvVar is the mon secret environment var
func SecretEnvVar() v1.EnvVar {
ref := &v1.SecretKeySelector{LocalObjectReference: v1.LocalObjectReference{Name: appName}, Key: monSecretName}
return v1.EnvVar{Name: "ROOKD_MON_SECRET", ValueFrom: &v1.EnvVarSource{SecretKeyRef: ref}}
}
// AdminSecretEnvVar is the admin secret environment var
func AdminSecretEnvVar() v1.EnvVar {
ref := &v1.SecretKeySelector{LocalObjectReference: v1.LocalObjectReference{Name: appName}, Key: adminSecretName}
return v1.EnvVar{Name: "ROOKD_ADMIN_SECRET", ValueFrom: &v1.EnvVarSource{SecretKeyRef: ref}}
@@ -52,7 +58,7 @@ func (c *Cluster) getLabels(name string) map[string]string {
}
}
func (c *Cluster) makeReplicaSet(config *MonConfig, nodeName string) *extensions.ReplicaSet {
func (c *Cluster) makeReplicaSet(config *monConfig, nodeName string) *extensions.ReplicaSet {
rs := &extensions.ReplicaSet{}
rs.Name = config.Name
@@ -71,7 +77,7 @@ func (c *Cluster) makeReplicaSet(config *MonConfig, nodeName string) *extensions
return rs
}
func (c *Cluster) makeMonPod(config *MonConfig, nodeName string) *v1.Pod {
func (c *Cluster) makeMonPod(config *monConfig, nodeName string) *v1.Pod {
dataDirSource := v1.VolumeSource{EmptyDir: &v1.EmptyDirVolumeSource{}}
if c.dataDirHostPath != "" {
dataDirSource = v1.VolumeSource{HostPath: &v1.HostPathVolumeSource{Path: c.dataDirHostPath}}
@@ -103,7 +109,7 @@ func (c *Cluster) makeMonPod(config *MonConfig, nodeName string) *v1.Pod {
return pod
}
func (c *Cluster) monContainer(config *MonConfig, fsid string) v1.Container {
func (c *Cluster) monContainer(config *monConfig, fsid string) v1.Container {
return v1.Container{
Args: []string{
@@ -129,8 +135,8 @@ func (c *Cluster) monContainer(config *MonConfig, fsid string) v1.Container {
Env: []v1.EnvVar{
{Name: k8sutil.PodIPEnvVar, ValueFrom: &v1.EnvVarSource{FieldRef: &v1.ObjectFieldSelector{FieldPath: "status.podIP"}}},
ClusterNameEnvVar(c.Namespace),
MonEndpointEnvVar(),
MonSecretEnvVar(),
EndpointEnvVar(),
SecretEnvVar(),
AdminSecretEnvVar(),
k8sutil.ConfigOverrideEnvVar(),
},
+1 -1
View File
@@ -37,7 +37,7 @@ func testPodSpec(t *testing.T, dataDir string) {
clientset := testop.New(1)
c := New(&clusterd.Context{KubeContext: kit.KubeContext{Clientset: clientset}}, "ns", dataDir, "myversion", k8sutil.Placement{})
c.clusterInfo = testop.CreateClusterInfo(0)
config := &MonConfig{Name: "mon0", Port: 6790}
config := &monConfig{Name: "mon0", Port: 6790}
pod := c.makeMonPod(config, "foo")
assert.NotNil(t, pod)
+8
View File
@@ -16,12 +16,15 @@ limitations under the License.
Some of the code below came from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package operator to manage Kubernetes storage.
package operator
import (
"fmt"
"time"
"github.com/coreos/pkg/capnslog"
"github.com/kubernetes-incubator/external-storage/lib/controller"
"github.com/rook/rook/pkg/clusterd"
"github.com/rook/rook/pkg/operator/cluster"
@@ -39,6 +42,9 @@ const (
provisionerName = "rook.io/block"
)
var logger = capnslog.NewPackageLogger("github.com/rook/rook", "operator")
// Operator type for managing storage
type Operator struct {
context *clusterd.Context
resources []kit.CustomResource
@@ -58,6 +64,7 @@ type resourceManager interface {
Manage()
}
// New creates an operator instance
func New(context *clusterd.Context) *Operator {
poolInitiator := newPoolInitiator(context)
@@ -73,6 +80,7 @@ func New(context *clusterd.Context) *Operator {
}
}
// Run the operator instance
func (o *Operator) Run() error {
for {
+7 -2
View File
@@ -13,6 +13,8 @@ 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 osd for the Ceph OSDs.
package osd
import (
@@ -40,6 +42,7 @@ const (
appNameFmt = "rook-ceph-osd-%s"
)
// Cluster keeps track of the OSDs
type Cluster struct {
context *clusterd.Context
Namespace string
@@ -50,6 +53,7 @@ type Cluster struct {
dataDirHostPath string
}
// New creates an instance of the OSD manager
func New(context *clusterd.Context, namespace, version string, storageSpec StorageSpec, dataDirHostPath string, placement k8sutil.Placement) *Cluster {
return &Cluster{
context: context,
@@ -61,6 +65,7 @@ func New(context *clusterd.Context, namespace, version string, storageSpec Stora
}
}
// Start the osd management
func (c *Cluster) Start() error {
logger.Infof("start running osds in namespace %s", c.Namespace)
@@ -178,8 +183,8 @@ func (c *Cluster) osdContainer(devices []Device, directories []Directory, select
envVars := []v1.EnvVar{
nodeNameEnvVar(),
opmon.ClusterNameEnvVar(c.Namespace),
opmon.MonEndpointEnvVar(),
opmon.MonSecretEnvVar(),
opmon.EndpointEnvVar(),
opmon.SecretEnvVar(),
opmon.AdminSecretEnvVar(),
k8sutil.ConfigDirEnvVar(),
k8sutil.ConfigOverrideEnvVar(),
+8
View File
@@ -13,12 +13,15 @@ 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 osd for the Ceph OSDs.
package osd
import (
cephosd "github.com/rook/rook/pkg/ceph/osd"
)
// StorageSpec CRD settings
type StorageSpec struct {
Nodes []Node `json:"nodes,omitempty"`
UseAllNodes bool `json:"useAllNodes,omitempty"`
@@ -26,6 +29,7 @@ type StorageSpec struct {
Config
}
// Node specific CRD settings
type Node struct {
Name string `json:"name,omitempty"`
Devices []Device `json:"devices,omitempty"`
@@ -34,14 +38,17 @@ type Node struct {
Config
}
// Device CRD settings
type Device struct {
Name string `json:"name,omitempty"`
}
// Directory CRD settings
type Directory struct {
Path string `json:"path,omitempty"`
}
// Selection CRD settings
type Selection struct {
// Whether to consume all the storage devices found on a machine
UseAllDevices *bool `json:"useAllDevices,omitempty"`
@@ -52,6 +59,7 @@ type Selection struct {
MetadataDevice string `json:"metadataDevice,omitempty"`
}
// Config CRD settings
type Config struct {
StoreConfig cephosd.StoreConfig `json:"storeConfig,omitempty"`
Location string `json:"location,omitempty"`
+2 -2
View File
@@ -45,12 +45,12 @@ nodes:
deviceFilter: "^sd."`)
// convert the raw spec yaml into JSON
rawJson, err := yaml.YAMLToJSON(specYaml)
rawJSON, err := yaml.YAMLToJSON(specYaml)
assert.Nil(t, err)
// unmarshal the JSON into a strongly typed storage spec object
var storageSpec StorageSpec
err = json.Unmarshal(rawJson, &storageSpec)
err = json.Unmarshal(rawJSON, &storageSpec)
assert.Nil(t, err)
// the unmarshalled storage spec should equal the expected spec below
+4
View File
@@ -13,12 +13,15 @@ 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 osd for the Ceph OSDs.
package osd
import (
cephosd "github.com/rook/rook/pkg/ceph/osd"
)
// AnyUseAllDevices gets whether to use all devices
func (s *StorageSpec) AnyUseAllDevices() bool {
if s.Selection.getUseAllDevices() {
return true
@@ -33,6 +36,7 @@ func (s *StorageSpec) AnyUseAllDevices() bool {
return false
}
// ClearUseAllDevices clears all devices
func (s *StorageSpec) ClearUseAllDevices() {
clear := false
s.Selection.UseAllDevices = &clear
+2
View File
@@ -16,6 +16,8 @@ limitations under the License.
Some of the code below came from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package operator to manage Kubernetes storage.
package operator
import (
+2 -1
View File
@@ -12,8 +12,9 @@ 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 operator to manage Kubernetes storage.
package operator
import (
+7 -2
View File
@@ -13,6 +13,8 @@ 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 rgw for the Ceph object store.
package rgw
import (
@@ -37,6 +39,7 @@ const (
keyringName = "keyring"
)
// Cluster for rgw management
type Cluster struct {
context *clusterd.Context
Namespace string
@@ -45,6 +48,7 @@ type Cluster struct {
Replicas int32
}
// New creates an instance of an rgw manager
func New(context *clusterd.Context, namespace, version string, placement k8sutil.Placement) *Cluster {
return &Cluster{
context: context,
@@ -55,6 +59,7 @@ func New(context *clusterd.Context, namespace, version string, placement k8sutil
}
}
// Start the rgw manager
func (c *Cluster) Start() error {
logger.Infof("start running rgw")
@@ -165,8 +170,8 @@ func (c *Cluster) rgwContainer() v1.Container {
Env: []v1.EnvVar{
{Name: "ROOKD_RGW_KEYRING", ValueFrom: &v1.EnvVarSource{SecretKeyRef: &v1.SecretKeySelector{LocalObjectReference: v1.LocalObjectReference{Name: appName}, Key: keyringName}}},
opmon.ClusterNameEnvVar(c.Namespace),
opmon.MonEndpointEnvVar(),
opmon.MonSecretEnvVar(),
opmon.EndpointEnvVar(),
opmon.SecretEnvVar(),
opmon.AdminSecretEnvVar(),
k8sutil.ConfigOverrideEnvVar(),
},
+3
View File
@@ -13,6 +13,8 @@ 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 test for the operator tests.
package test
import (
@@ -22,6 +24,7 @@ import (
"k8s.io/client-go/kubernetes/fake"
)
// New creates a fake K8s cluster
func New(nodes int) *fake.Clientset {
clientset := fake.NewSimpleClientset()
for i := 0; i < nodes; i++ {
+3
View File
@@ -13,6 +13,8 @@ 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 test for the operator tests
package test
import (
@@ -21,6 +23,7 @@ import (
"github.com/rook/rook/pkg/ceph/mon"
)
// CreateClusterInfo creates a test cluster
func CreateClusterInfo(mons int) *mon.ClusterInfo {
c := &mon.ClusterInfo{
FSID: "12345",
+2
View File
@@ -16,6 +16,8 @@ limitations under the License.
Some of the code below came from https://github.com/coreos/etcd-operator
which also has the apache 2.0 license.
*/
// Package operator to manage Kubernetes storage.
package operator
import "sync"