diff --git a/pkg/operator/api/api.go b/pkg/operator/api/api.go index 2d39dbb81..f4715b5fe 100644 --- a/pkg/operator/api/api.go +++ b/pkg/operator/api/api.go @@ -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), }, } diff --git a/pkg/operator/cluster.go b/pkg/operator/cluster.go index a528ff32a..4f4136ad6 100644 --- a/pkg/operator/cluster.go +++ b/pkg/operator/cluster.go @@ -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 diff --git a/pkg/operator/cluster/cluster.go b/pkg/operator/cluster/cluster.go index 9ee58ca0a..68c876c0e 100644 --- a/pkg/operator/cluster/cluster.go +++ b/pkg/operator/cluster/cluster.go @@ -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 diff --git a/pkg/operator/cluster/cluster_list.go b/pkg/operator/cluster/cluster_list.go index 919b040a3..7965034bb 100644 --- a/pkg/operator/cluster/cluster_list.go +++ b/pkg/operator/cluster/cluster_list.go @@ -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) diff --git a/pkg/operator/cluster/pool.go b/pkg/operator/cluster/pool.go index caf270b7f..f157cd024 100644 --- a/pkg/operator/cluster/pool.go +++ b/pkg/operator/cluster/pool.go @@ -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 { diff --git a/pkg/operator/cluster/pool_list.go b/pkg/operator/cluster/pool_list.go index 69c3dc5f5..2ee6c48a5 100644 --- a/pkg/operator/cluster/pool_list.go +++ b/pkg/operator/cluster/pool_list.go @@ -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) diff --git a/pkg/operator/cluster/spec.go b/pkg/operator/cluster/spec.go index c59d50616..705731dff 100644 --- a/pkg/operator/cluster/spec.go +++ b/pkg/operator/cluster/spec.go @@ -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) } diff --git a/pkg/operator/k8sutil/k8sutil.go b/pkg/operator/k8sutil/k8sutil.go index 3b0bc497f..32eb69294 100644 --- a/pkg/operator/k8sutil/k8sutil.go +++ b/pkg/operator/k8sutil/k8sutil.go @@ -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" ) diff --git a/pkg/operator/k8sutil/placement.go b/pkg/operator/k8sutil/placement.go index 0cb09dc22..d43c12e76 100644 --- a/pkg/operator/k8sutil/placement.go +++ b/pkg/operator/k8sutil/placement.go @@ -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 { diff --git a/pkg/operator/k8sutil/placement_test.go b/pkg/operator/k8sutil/placement_test.go index 42625b1e7..775e9b7ed 100644 --- a/pkg/operator/k8sutil/placement_test.go +++ b/pkg/operator/k8sutil/placement_test.go @@ -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 diff --git a/pkg/operator/k8sutil/pod.go b/pkg/operator/k8sutil/pod.go index 11391ecb0..26ab103bc 100644 --- a/pkg/operator/k8sutil/pod.go +++ b/pkg/operator/k8sutil/pod.go @@ -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 -} diff --git a/pkg/operator/k8sutil/volume.go b/pkg/operator/k8sutil/volume.go index 7193610fb..af941ad7f 100644 --- a/pkg/operator/k8sutil/volume.go +++ b/pkg/operator/k8sutil/volume.go @@ -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])? diff --git a/pkg/operator/kit/client.go b/pkg/operator/kit/client.go index dd8e08295..880b38901 100644 --- a/pkg/operator/kit/client.go +++ b/pkg/operator/kit/client.go @@ -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 ( diff --git a/pkg/operator/kit/log.go b/pkg/operator/kit/log.go index 815816685..4690e87d6 100644 --- a/pkg/operator/kit/log.go +++ b/pkg/operator/kit/log.go @@ -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" diff --git a/pkg/operator/kit/resource.go b/pkg/operator/kit/resource.go index 215a84a4f..ebc96f07e 100644 --- a/pkg/operator/kit/resource.go +++ b/pkg/operator/kit/resource.go @@ -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 ( diff --git a/pkg/operator/kit/retry.go b/pkg/operator/kit/retry.go index ef576ca80..b966afc43 100644 --- a/pkg/operator/kit/retry.go +++ b/pkg/operator/kit/retry.go @@ -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 ( diff --git a/pkg/operator/kit/timer.go b/pkg/operator/kit/timer.go index 012ddcf5a..dbd227e7c 100644 --- a/pkg/operator/kit/timer.go +++ b/pkg/operator/kit/timer.go @@ -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" diff --git a/pkg/operator/kit/watcher.go b/pkg/operator/kit/watcher.go index 018f49c32..6e953e5c1 100644 --- a/pkg/operator/kit/watcher.go +++ b/pkg/operator/kit/watcher.go @@ -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 ( diff --git a/pkg/operator/log.go b/pkg/operator/log.go deleted file mode 100644 index 3352eee89..000000000 --- a/pkg/operator/log.go +++ /dev/null @@ -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") diff --git a/pkg/operator/mds/mds.go b/pkg/operator/mds/mds.go index 8b1abfef3..75d01e5aa 100644 --- a/pkg/operator/mds/mds.go +++ b/pkg/operator/mds/mds.go @@ -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(), }, diff --git a/pkg/operator/mgr/mgr.go b/pkg/operator/mgr/mgr.go index 8665c1610..1a43bb98a 100644 --- a/pkg/operator/mgr/mgr.go +++ b/pkg/operator/mgr/mgr.go @@ -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(), }, diff --git a/pkg/operator/mon/health.go b/pkg/operator/mon/health.go index 1236965b6..bfaf0d68e 100644 --- a/pkg/operator/mon/health.go +++ b/pkg/operator/mon/health.go @@ -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 { diff --git a/pkg/operator/mon/mon.go b/pkg/operator/mon/mon.go index dca58413a..c7d230b22 100644 --- a/pkg/operator/mon/mon.go +++ b/pkg/operator/mon/mon.go @@ -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) diff --git a/pkg/operator/mon/mon_test.go b/pkg/operator/mon/mon_test.go index 242e1a063..b78117f81 100644 --- a/pkg/operator/mon/mon_test.go +++ b/pkg/operator/mon/mon_test.go @@ -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() diff --git a/pkg/operator/mon/spec.go b/pkg/operator/mon/spec.go index b7b7b52c6..bf513b72f 100644 --- a/pkg/operator/mon/spec.go +++ b/pkg/operator/mon/spec.go @@ -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(), }, diff --git a/pkg/operator/mon/spec_test.go b/pkg/operator/mon/spec_test.go index f1718207f..7081d0815 100644 --- a/pkg/operator/mon/spec_test.go +++ b/pkg/operator/mon/spec_test.go @@ -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) diff --git a/pkg/operator/operator.go b/pkg/operator/operator.go index 5ea2cd91c..4e79e770e 100644 --- a/pkg/operator/operator.go +++ b/pkg/operator/operator.go @@ -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 { diff --git a/pkg/operator/osd/osd.go b/pkg/operator/osd/osd.go index ac9b4bfc2..d08e0a4ec 100644 --- a/pkg/operator/osd/osd.go +++ b/pkg/operator/osd/osd.go @@ -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(), diff --git a/pkg/operator/osd/spec.go b/pkg/operator/osd/spec.go index 9f9ef0355..05271d6bd 100644 --- a/pkg/operator/osd/spec.go +++ b/pkg/operator/osd/spec.go @@ -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"` diff --git a/pkg/operator/osd/spec_test.go b/pkg/operator/osd/spec_test.go index 372c4773b..c0ff3cc4d 100644 --- a/pkg/operator/osd/spec_test.go +++ b/pkg/operator/osd/spec_test.go @@ -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 diff --git a/pkg/operator/osd/storage.go b/pkg/operator/osd/storage.go index 50dd6aa51..5921e984c 100644 --- a/pkg/operator/osd/storage.go +++ b/pkg/operator/osd/storage.go @@ -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 diff --git a/pkg/operator/pool.go b/pkg/operator/pool.go index af0e290cb..66889f274 100644 --- a/pkg/operator/pool.go +++ b/pkg/operator/pool.go @@ -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 ( diff --git a/pkg/operator/provisioner.go b/pkg/operator/provisioner.go index 1a3b4d9be..5fe08497c 100644 --- a/pkg/operator/provisioner.go +++ b/pkg/operator/provisioner.go @@ -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 ( diff --git a/pkg/operator/rgw/rgw.go b/pkg/operator/rgw/rgw.go index 1284da38f..cae55642d 100644 --- a/pkg/operator/rgw/rgw.go +++ b/pkg/operator/rgw/rgw.go @@ -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(), }, diff --git a/pkg/operator/test/client.go b/pkg/operator/test/client.go index 6e6248209..7c31a785d 100644 --- a/pkg/operator/test/client.go +++ b/pkg/operator/test/client.go @@ -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++ { diff --git a/pkg/operator/test/info.go b/pkg/operator/test/info.go index 1202f6beb..7195a7a7e 100644 --- a/pkg/operator/test/info.go +++ b/pkg/operator/test/info.go @@ -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", diff --git a/pkg/operator/tracker.go b/pkg/operator/tracker.go index 19e6d8ce9..8aa954498 100644 --- a/pkg/operator/tracker.go +++ b/pkg/operator/tracker.go @@ -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"