diff --git a/.github/workflows/main.yml b/.github/workflows/main.yml index 3a9aa50..c4f21e1 100644 --- a/.github/workflows/main.yml +++ b/.github/workflows/main.yml @@ -18,7 +18,7 @@ jobs: - name: run golangci-lint uses: golangci/golangci-lint-action@v2 with: - version: v1.35 + version: v1.35.2 static-test: runs-on: ubuntu-20.04 steps: @@ -45,7 +45,7 @@ jobs: e2e: needs: [image-build] runs-on: self-hosted - timeout-minutes: 40 + timeout-minutes: 70 strategy: matrix: config: diff --git a/Makefile b/Makefile index bc3329c..8dba075 100644 --- a/Makefile +++ b/Makefile @@ -143,6 +143,8 @@ e2e-deploy: registry docker-build docker-push deploy hack/e2e.sh bootstrap hack/e2e.sh update_cm_after_delete hack/e2e.sh update_secret_after_delete + hack/e2e.sh add_host + hack/e2e.sh add_disk hack/e2e.sh delete_cluster e2e: e2e-deploy clean diff --git a/api/v1alpha1/cephcluster_types.go b/api/v1alpha1/cephcluster_types.go index 172b6ff..ada3d1a 100644 --- a/api/v1alpha1/cephcluster_types.go +++ b/api/v1alpha1/cephcluster_types.go @@ -53,19 +53,30 @@ type CephClusterSpec struct { Config map[string]string `json:"config,omitempty"` } -// CephClusterState is the current state of CephCluster -type CephClusterState string +// CephClusterCondition is the current condition of CephCluster +type CephClusterCondition string const ( // ConditionReadyToUse indicates CephCluster is ready to use - ConditionReadyToUse = "ReadyToUse" + ConditionReadyToUse CephClusterCondition = "ReadyToUse" + // ConditionBootstrapped indicates CephCluster is bootstrapped + ConditionBootstrapped CephClusterCondition = "Bootstrapped" + // ConditionOsdDeployed indicates CephCluster is deployed with osds + ConditionOsdDeployed CephClusterCondition = "OsdDeployed" ) +// CephClusterState is the current state of CephCluster +type CephClusterState string + const ( + // CephClusterStatePending indicates CephClusterState is creating or updating + CephClusterStatePending CephClusterState = "Pending" // CephClusterStateCreating indicates CephClusterState is creating CephClusterStateCreating CephClusterState = "Creating" - // CephClusterStateCompleted indicates CephClusterState is completed - CephClusterStateCompleted CephClusterState = "Completed" + // CephClusterStateUpdating indicates CephClusterState is updating + CephClusterStateUpdating CephClusterState = "Updating" + // CephClusterStateRunning indicates CephClusterState is available + CephClusterStateRunning CephClusterState = "Running" // CephClusterStateError indicates CephClusterState is error CephClusterStateError CephClusterState = "Error" ) @@ -74,7 +85,7 @@ const ( type CephClusterStatus struct { // INSERT ADDITIONAL STATUS FIELD - define observed state of cluster // Important: Run "make" to regenerate code after modifying this file - + DeployNode Node `json:"deployNode,omitempty"` State CephClusterState `json:"state"` Conditions []metav1.Condition `json:"conditions,omitempty"` } diff --git a/api/v1alpha1/zz_generated.deepcopy.go b/api/v1alpha1/zz_generated.deepcopy.go index b753157..9999d01 100644 --- a/api/v1alpha1/zz_generated.deepcopy.go +++ b/api/v1alpha1/zz_generated.deepcopy.go @@ -157,6 +157,7 @@ func (in *CephClusterSpec) DeepCopy() *CephClusterSpec { // DeepCopyInto is an autogenerated deepcopy function, copying the receiver, writing into out. in must be non-nil. func (in *CephClusterStatus) DeepCopyInto(out *CephClusterStatus) { *out = *in + out.DeployNode = in.DeployNode if in.Conditions != nil { in, out := &in.Conditions, &out.Conditions *out = make([]v1.Condition, len(*in)) diff --git a/config/crd/bases/hypersds.tmax.io_cephclusters.yaml b/config/crd/bases/hypersds.tmax.io_cephclusters.yaml index db666f0..1add827 100644 --- a/config/crd/bases/hypersds.tmax.io_cephclusters.yaml +++ b/config/crd/bases/hypersds.tmax.io_cephclusters.yaml @@ -165,6 +165,25 @@ spec: - type type: object type: array + deployNode: + description: 'INSERT ADDITIONAL STATUS FIELD - define observed state + of cluster Important: Run "make" to regenerate code after modifying + this file' + properties: + hostName: + type: string + ip: + type: string + password: + type: string + userId: + type: string + required: + - hostName + - ip + - password + - userId + type: object state: description: CephClusterState is the current state of CephCluster type: string diff --git a/config/manager/kustomization.yaml b/config/manager/kustomization.yaml index 6d06efc..4676d57 100644 --- a/config/manager/kustomization.yaml +++ b/config/manager/kustomization.yaml @@ -1,7 +1,13 @@ resources: - manager.yaml patches: -- target: +- path: manager_resource_patch.yaml + target: kind: Deployment name: controller-manager - path: manager_resource_patch.yaml +apiVersion: kustomize.config.k8s.io/v1beta1 +kind: Kustomization +images: +- name: controller + newName: 192.168.9.21:5000/hypersds-operator + newTag: latest diff --git a/config/samples/hypersds_v1alpha1_cephcluster_diskadd_unordered.yaml b/config/samples/hypersds_v1alpha1_cephcluster_diskadd_unordered.yaml new file mode 100644 index 0000000..f9cf556 --- /dev/null +++ b/config/samples/hypersds_v1alpha1_cephcluster_diskadd_unordered.yaml @@ -0,0 +1,24 @@ +apiVersion: hypersds.tmax.io/v1alpha1 +kind: CephCluster +metadata: + name: cephcluster-sample +spec: + mon: + count: 1 + osd: + - hostName: centos-node5 + devices: + - /dev/sdb + - /dev/sdc + - hostName: centos-node4 + devices: + - /dev/sdb + nodes: + - ip: 192.168.33.21 + userId: root + password: "ck@3434" + hostName: centos-node4 + - ip: 192.168.33.22 + userId: root + password: "ck@3434" + hostName: centos-node5 diff --git a/config/samples/hypersds_v1alpha1_cephcluster_hostadd.yaml b/config/samples/hypersds_v1alpha1_cephcluster_hostadd.yaml new file mode 100644 index 0000000..bde084b --- /dev/null +++ b/config/samples/hypersds_v1alpha1_cephcluster_hostadd.yaml @@ -0,0 +1,23 @@ +apiVersion: hypersds.tmax.io/v1alpha1 +kind: CephCluster +metadata: + name: cephcluster-sample +spec: + mon: + count: 1 + osd: + - hostName: centos-node4 + devices: + - /dev/sdb + - hostName: centos-node5 + devices: + - /dev/sdb + nodes: + - ip: 192.168.33.21 + userId: root + password: "ck@3434" + hostName: centos-node4 + - ip: 192.168.33.22 + userId: root + password: "ck@3434" + hostName: centos-node5 diff --git a/controllers/cephcluster_controller.go b/controllers/cephcluster_controller.go index 3d6822f..7a29892 100644 --- a/controllers/cephcluster_controller.go +++ b/controllers/cephcluster_controller.go @@ -72,10 +72,17 @@ func (r *CephClusterReconciler) Reconcile(req ctrl.Request) (ctrl.Result, error) return nil } if err := syncAll(); err != nil { - if err2 := r.updateStateWithReadyToUse(hypersdsv1alpha1.CephClusterStateError, metav1.ConditionFalse, "SeeMessages", err.Error()); err2 != nil { + if err2 := r.updateState(hypersdsv1alpha1.CephClusterStateError); err2 != nil { + return ctrl.Result{}, err2 + } + if err2 := r.updateCondition(&metav1.Condition{ + Type: string(hypersdsv1alpha1.ConditionReadyToUse), + Status: metav1.ConditionFalse, + Reason: "SeeMessages", + Message: err.Error(), + }); err2 != nil { return ctrl.Result{}, err2 } - return ctrl.Result{}, err } return ctrl.Result{}, nil @@ -92,14 +99,12 @@ func (r *CephClusterReconciler) SetupWithManager(mgr ctrl.Manager) error { Complete(r) } -func (r *CephClusterReconciler) updateStateWithReadyToUse(state hypersdsv1alpha1.CephClusterState, readyToUseStatus metav1.ConditionStatus, - reason, message string) error { - meta.SetStatusCondition(&r.Cluster.Status.Conditions, metav1.Condition{ - Type: hypersdsv1alpha1.ConditionReadyToUse, - Status: readyToUseStatus, - Reason: reason, - Message: message, - }) +func (r *CephClusterReconciler) updateState(state hypersdsv1alpha1.CephClusterState) error { r.Cluster.Status.State = state return r.Client.Status().Update(context.TODO(), r.Cluster) } + +func (r *CephClusterReconciler) updateCondition(condition *metav1.Condition) error { + meta.SetStatusCondition(&r.Cluster.Status.Conditions, *condition) + return r.Client.Status().Update(context.TODO(), r.Cluster) +} diff --git a/controllers/configmap.go b/controllers/configmap.go index d3d4ba4..fe0ab8c 100644 --- a/controllers/configmap.go +++ b/controllers/configmap.go @@ -23,7 +23,21 @@ func (r *CephClusterReconciler) syncConfigMap() error { } klog.Infof("syncConfigMap: creating config map %s", r.Cluster.Name) - if err := r.updateStateWithReadyToUse(v1alpha1.CephClusterStateCreating, metav1.ConditionFalse, "CephClusterIsCreating", "CephCluster is creating"); err != nil { + if err := r.updateState(v1alpha1.CephClusterStateCreating); err != nil { + return err + } + if err := r.updateCondition(&metav1.Condition{ + Type: string(v1alpha1.ConditionBootstrapped), + Status: metav1.ConditionFalse, + Reason: "BootstrappingIsNotFinished", + }); err != nil { + return err + } + if err := r.updateCondition(&metav1.Condition{ + Type: string(v1alpha1.ConditionReadyToUse), + Status: metav1.ConditionFalse, + Reason: "CephClusterIsCreating", + }); err != nil { return err } diff --git a/controllers/configmap_test.go b/controllers/configmap_test.go index 5a5748e..865bed4 100644 --- a/controllers/configmap_test.go +++ b/controllers/configmap_test.go @@ -39,7 +39,7 @@ var _ = Describe("syncConfigMap", func() { cc := &hypersdsv1alpha1.CephCluster{} err := r.Client.Get(context.TODO(), types.NamespacedName{Namespace: r.Cluster.Namespace, Name: r.Cluster.Name}, cc) Expect(err).Should(BeNil()) - cond := meta.FindStatusCondition(cc.Status.Conditions, hypersdsv1alpha1.ConditionReadyToUse) + cond := meta.FindStatusCondition(cc.Status.Conditions, string(hypersdsv1alpha1.ConditionReadyToUse)) Expect(cond).ShouldNot(BeNil()) Expect(cond.Status).Should(Equal(metav1.ConditionFalse)) }) diff --git a/controllers/provisioner.go b/controllers/provisioner.go index 3b8a151..2dab901 100644 --- a/controllers/provisioner.go +++ b/controllers/provisioner.go @@ -1,9 +1,11 @@ package controllers import ( + "context" "github.com/tmax-cloud/hypersds-operator/api/v1alpha1" "github.com/tmax-cloud/hypersds-operator/pkg/provisioner/provisioner" "k8s.io/apimachinery/pkg/api/errors" + "k8s.io/apimachinery/pkg/api/meta" "k8s.io/apimachinery/pkg/apis/meta/v1" "k8s.io/klog" ) @@ -23,6 +25,10 @@ func (r *CephClusterReconciler) syncProvisioner() error { return err } + if err := r.updateDeployNode(); err != nil { + return err + } + klog.Infof("syncProvisioner: bootstrapping ceph cluster %s", r.Cluster.Name) provisionerInstance, err := provisioner.NewProvisioner(r.Cluster.Spec, r.Client, r.Cluster.Namespace, r.Cluster.Name) if err != nil { @@ -31,9 +37,40 @@ func (r *CephClusterReconciler) syncProvisioner() error { if err := provisionerInstance.Run(); err != nil { return err } - if err := r.updateStateWithReadyToUse(v1alpha1.CephClusterStateCompleted, v1.ConditionTrue, "CephClusterIsReady", "Ceph cluster is ready to use"); err != nil { + if err := r.updateState(v1alpha1.CephClusterStateRunning); err != nil { return err } + if err := r.updateCondition(&v1.Condition{ + Type: string(v1alpha1.ConditionReadyToUse), + Status: v1.ConditionTrue, + Reason: "CephClusterIsReady", + }); err != nil { + return err + } + if meta.IsStatusConditionFalse(r.Cluster.Status.Conditions, string(v1alpha1.ConditionBootstrapped)) { + if err := r.updateCondition(&v1.Condition{ + Type: string(v1alpha1.ConditionBootstrapped), + Status: v1.ConditionTrue, + Reason: "CephClusterIsBootstrapped", + }); err != nil { + return err + } + } return nil } + +func (r *CephClusterReconciler) updateDeployNode() error { + if r.Cluster.Status.DeployNode == (v1alpha1.Node{}) { + r.Cluster.Status.DeployNode = r.getDefaultDeployNode() + if err := r.Client.Status().Update(context.TODO(), r.Cluster); err != nil { + return err + } + } + return nil +} + +// Decide deploying node (currently, first node is deploying node) +func (r *CephClusterReconciler) getDefaultDeployNode() v1alpha1.Node { + return r.Cluster.Spec.Nodes[0] +} diff --git a/hack/e2e.sh b/hack/e2e.sh index bf5fafa..3d96b56 100755 --- a/hack/e2e.sh +++ b/hack/e2e.sh @@ -29,7 +29,7 @@ case "${1:-}" in bootstrap) echo "deploying ceph cluster cr ..." kubectl apply -f config/samples/hypersds_v1alpha1_cephcluster.yaml - wait_condition "kubectl get cephclusters.hypersds.tmax.io | grep Completed" 2400 60 + wait_condition "kubectl get cephclusters.hypersds.tmax.io | grep Running" 2400 60 ;; update_cm_after_delete) echo "deleting configmap ..." @@ -43,6 +43,16 @@ case "${1:-}" in sleep 3 wait_condition "kubectl describe secret cephcluster-sample-keyring | grep keyring:" 300 60 ;; + add_host) + echo "adding host ..." + kubectl apply -f config/samples/hypersds_v1alpha1_cephcluster_hostadd.yaml + sleep 600 + ;; + add_disk) + echo "adding disk ..." + kubectl apply -f config/samples/hypersds_v1alpha1_cephcluster_diskadd_unordered.yaml + sleep 600 + ;; delete_cluster) echo "deleting ceph cluster cr ..." kubectl delete -f config/samples/hypersds_v1alpha1_cephcluster.yaml diff --git a/pkg/provisioner/provisioner/cmd.go b/pkg/provisioner/provisioner/cmd.go new file mode 100644 index 0000000..b5da2ab --- /dev/null +++ b/pkg/provisioner/provisioner/cmd.go @@ -0,0 +1,314 @@ +package provisioner + +import ( + "fmt" + "github.com/tmax-cloud/hypersds-operator/pkg/common/util" + "github.com/tmax-cloud/hypersds-operator/pkg/common/wrapper" + "github.com/tmax-cloud/hypersds-operator/pkg/provisioner/node" + "github.com/tmax-cloud/hypersds-operator/pkg/provisioner/osd" + "strings" +) + +const ( + installCmd = "install -y " + updateCmd = "update -y" +) + +func isDockerInstalled(targetNode *node.Node) (bool, error) { + const checkDockerCmd = "docker -v" + output, err := targetNode.RunSSHCmd(wrapper.SSHWrapper, checkDockerCmd) + if err != nil { + if strings.Contains(output.String(), cmdNotFound) { + return false, nil + } + fmt.Println("[Provisioner] docker installation check is failed") + return false, err + } + return true, nil +} + +func installCommonPackage(n *node.Node) error { + fmt.Println("\n----------------Start to install base package---------------") + var err error + packager := n.GetOs().Packager + + fmt.Println("[installBasePackage] executing packager update") + if packager == node.Apt { + err = processCmdOnNode(n, string(packager)+updateCmd) + if err != nil { + return err + } + } + + fmt.Println("[installBasePackage] install common packages") + const commonPkgs = "chrony python3" + err = processCmdOnNode(n, string(packager)+installCmd+commonPkgs) + if err != nil { + return err + } + + fmt.Println("[installBasePackage] enable chronyd") + const enableChronydCmd = "systemctl enable chronyd" + err = processCmdOnNode(n, enableChronydCmd) + if err != nil { + return err + } + + fmt.Println("[installBasePackage] restart chronyd") + const restartChronydCmd = "systemctl restart chronyd" + err = processCmdOnNode(n, restartChronydCmd) + if err != nil { + return err + } + + return nil +} + +func installDocker(n *node.Node) error { + fmt.Println("[installBasePackage] install docker dependencies") + var err error + packager := n.GetOs().Packager + + const dockerDepAptPkgs = "apt-transport-https ca-certificates curl gnupg lsb-release" + const dockerDepYumPkgs = "yum-utils" + pkgs := dockerDepAptPkgs + if packager == node.Yum { + pkgs = dockerDepYumPkgs + } + err = processCmdOnNode(n, string(packager)+installCmd+pkgs) + if err != nil { + return err + } + + fmt.Println("[installBasePackage] executing curl docker ...") + const dockerUbuntuGpgKeyCmd = "curl -s https://download.docker.com/linux/ubuntu/gpg | apt-key add - &>/dev/null" + osName := n.GetOs().Distro + if osName == node.Ubuntu { + err = processCmdOnNode(n, dockerUbuntuGpgKeyCmd) + if err != nil { + return err + } + } + + fmt.Println("[installBasePackage] executing add-apt-repo docker ...") + const dockerAptRepoCmd = `add-apt-repository "deb [arch=amd64] https://download.docker.com/linux/ubuntu $(lsb_release -cs) stable"` + const dockerYumRepoCmd = "yum-config-manager --add-repo https://download.docker.com/linux/centos/docker-ce.repo" + addDockerRepoCmd := "" + switch packager { + case node.Apt: + addDockerRepoCmd += dockerAptRepoCmd + case node.Yum: + addDockerRepoCmd += dockerYumRepoCmd + } + err = processCmdOnNode(n, addDockerRepoCmd) + if err != nil { + return err + } + + fmt.Println("[installBasePackage] executing install docker-ce") + const dockerPkgs = "docker-ce" + if packager == node.Apt { + err = processCmdOnNode(n, string(packager)+updateCmd) + if err != nil { + return err + } + } + err = processCmdOnNode(n, string(packager)+installCmd+dockerPkgs) + if err != nil { + return err + } + + fmt.Println("[installBasePackage] executing sysctl docker") + const restartDockerCmd = "systemctl restart docker" + err = processCmdOnNode(n, restartDockerCmd) + if err != nil { + return err + } + + return nil +} + +func installCephadmPackage(targetNode *node.Node) error { + fmt.Println("\n----------------Start to install cephadm---------------") + + fmt.Println("[installCephadm] executing curl cephadm") + curlCephadmCmd := fmt.Sprintf("curl --silent --remote-name --location https://github.com/ceph/ceph/raw/v%s/src/cephadm/cephadm", cephVersion) + err := processCmdOnNode(targetNode, curlCephadmCmd) + if err != nil { + return err + } + + fmt.Println("[installCephadm] executing chmod") + const chmodCmd = "chmod +x cephadm" + err = processCmdOnNode(targetNode, chmodCmd) + if err != nil { + return err + } + + // TODO: Specify release version + fmt.Println("[installCephadm] executing cephadm add-repo") + admAddRepoCmd := fmt.Sprintf("./cephadm add-repo --version %s", cephVersion) + err = processCmdOnNode(targetNode, admAddRepoCmd) + if err != nil { + return err + } + + if targetNode.GetOs().Distro == node.Ubuntu { + // need for ubuntu 18.04 & ceph version 15.2.8 + fmt.Println("[installCephadm] executing curl cephadm gpg key") + addCephadmRepoCmd := "curl https://download.ceph.com/keys/release.asc | " + + "gpg --no-default-keyring --keyring /tmp/fix.gpg --import - && " + + "gpg --no-default-keyring --keyring /tmp/fix.gpg --export > /etc/apt/trusted.gpg.d/ceph.release.gpg && rm /tmp/fix.gpg" + err = processCmdOnNode(targetNode, addCephadmRepoCmd) + if err != nil { + return err + } + + fmt.Println("[installCephadm] executing cephadm apt-get update") + const aptUpdateCmd = "apt-get update" + err = processCmdOnNode(targetNode, aptUpdateCmd) + if err != nil { + return err + } + } + + fmt.Println("[installCephadm] executing cephadm install") + const admInstallCmd = "./cephadm install" + err = processCmdOnNode(targetNode, admInstallCmd) + if err != nil { + return err + } + + fmt.Println("[installCephadm] executing mkdir") + const mkdirCmd = "mkdir -p /etc/ceph" + err = processCmdOnNode(targetNode, mkdirCmd) + + return err +} + +func bootstrapCephadm(targetNode *node.Node, pathConfFromCr string) error { + fmt.Println("\n----------------Start to bootstrap ceph---------------") + + fmt.Println("[bootstrapCephadm] copying conf file") + err := createDstDir(targetNode, pathConfFromCr) + if err != nil { + return err + } + err = copyFile(targetNode, node.DESTINATION, pathConfFromCr, pathConfFromCr) + if err != nil { + return err + } + + deployNodeHostSpec := targetNode.GetHostSpec() + + monIP := deployNodeHostSpec.GetAddr() + + fmt.Println("[bootstrapCephadm] executing bootstrap") + admBootstrapCmd := fmt.Sprintf("cephadm --image %s bootstrap --mon-ip %s --config %s", + cephImageName, monIP, pathConfFromCr) + err = processCmdOnNode(targetNode, admBootstrapCmd) + if err != nil { + return err + } + + fmt.Println("[bootstrapCephadm] checking status") + const admHealthCheckCmd = "cephadm shell -- ceph -s" + err = processCmdOnNode(targetNode, admHealthCheckCmd) + + return err +} + +func (p *Provisioner) applyOsd(cephConf, cephKeyring []byte) error { + var err error + + fmt.Println("[applyOsd] get osds from CephOrch") + + cephName := p.getCephName() + pathConfigDir := p.getPathConfigDir() + + cmd := []string{"orch", "ls", "--service_type", "osd", "--export", "--refresh"} + output, err := util.RunCephCmd(wrapper.OsWrapper, wrapper.ExecWrapper, wrapper.IoUtilWrapper, cephConf, cephKeyring, cephName, cmd...) + if err != nil { + return processExecError(err, output) + } + + var osdsFromOrch []*osd.Osd + if !strings.Contains(output.String(), "No services reported") { + rawOsdsFromOrch := output.Bytes() + + osdsFromOrch, err = osd.NewOsdsFromCephOrch(wrapper.YamlWrapper, rawOsdsFromOrch) + if err != nil { + return err + } + } + osdsFromCephCr, err := osd.NewOsdsFromCephCr(p.cephCluster) + if err != nil { + return err + } + + var osdMap map[string]*osd.Osd + var removeOsdMap map[string]bool + + osdMap = make(map[string]*osd.Osd) + removeOsdMap = make(map[string]bool) + + for _, osdOrch := range osdsFromOrch { + osdService := osdOrch.GetService() + + osdServiceID := osdService.GetServiceID() + + osdMap[osdServiceID] = osdOrch + removeOsdMap[osdServiceID] = true + } + + fmt.Println("[applyOsd] compare osds between CephCR and CephOrch") + + for _, osdCephCr := range osdsFromCephCr { + osdService := osdCephCr.GetService() + osdServiceID := osdService.GetServiceID() + osdOrch, exist := osdMap[osdServiceID] + + var addDeviceList, removeDeviceList []string + + if exist { + addDeviceList, removeDeviceList, err = osdOrch.CompareDataDevices(osdCephCr) + if err != nil { + return err + } + removeOsdMap[osdServiceID] = false + // todo remove disk .... + fmt.Printf("[applyOsd] osd service: %s, add: %+q, remove: %+q\n", osdServiceID, addDeviceList, removeDeviceList) + } + + fmt.Println("[applyOsd] make osd yaml") + + osdFileName := pathConfigDir + osdServiceID + ".yaml" + err = osdCephCr.MakeYmlFile(wrapper.YamlWrapper, wrapper.IoUtilWrapper, osdFileName) + if err != nil { + return err + } + + fmt.Printf("[applyOsd] apply osd service: %s\n", osdServiceID) + + applyCmd := []string{"orch", "apply", "-i", osdFileName} + output, err = util.RunCephCmd(wrapper.OsWrapper, wrapper.ExecWrapper, wrapper.IoUtilWrapper, cephConf, cephKeyring, cephName, applyCmd...) + if err != nil { + return processExecError(err, output) + } + } + for osdServiceID, value := range removeOsdMap { + if !value { + continue + } + osdServiceName := "osd." + osdServiceID + + fmt.Printf("[applyOsd] remove osd service: %s\n", osdServiceName) + + removeCmd := []string{"orch", "rm", osdServiceName} + output, err = util.RunCephCmd(wrapper.OsWrapper, wrapper.ExecWrapper, wrapper.IoUtilWrapper, cephConf, cephKeyring, cephName, removeCmd...) + if err != nil { + return processExecError(err, output) + } + } + return nil +} diff --git a/pkg/provisioner/provisioner/host.go b/pkg/provisioner/provisioner/host.go index b6da28a..768d646 100644 --- a/pkg/provisioner/provisioner/host.go +++ b/pkg/provisioner/provisioner/host.go @@ -116,7 +116,7 @@ func (p *Provisioner) applyHost(yamlWrapper wrapper.YamlInterface, execWrapper w nodeIP := hostToApply.GetAddr() nodePw := cephHostNodesToApply[hostNameToApply].GetUserPw() - const sshKeyCheckOpt = "-oStrictHostKeyChecking=no -oUserKnownHostsFile=/dev/null" + const sshKeyCheckOpt = "-o StrictHostKeyChecking=no -o UserKnownHostsFile=/dev/null" sshPassCmd := fmt.Sprintf("sshpass -f <(printf '%%s\\n' %s)", nodePw) hostAuthApplyCmd := fmt.Sprintf("%s ssh-copy-id %s -f -i %s %s@%s", sshPassCmd, sshKeyCheckOpt, pathHostPub, nodeID, nodeIP) @@ -131,7 +131,7 @@ func (p *Provisioner) applyHost(yamlWrapper wrapper.YamlInterface, execWrapper w hostApplyCmd := []string{"orch", "apply", "-i", hostFileName} fmt.Println("Executing: " + strings.Join(hostApplyCmd, ",")) - hostApplyBuf, err = util.RunCephCmd(wrapper.OsWrapper, execWrapper, ioUtilWrapper, cephConf, cephKeyring, cephName, hostAuthGetCmd...) + hostApplyBuf, err = util.RunCephCmd(wrapper.OsWrapper, execWrapper, ioUtilWrapper, cephConf, cephKeyring, cephName, hostApplyCmd...) if err != nil { fmt.Println("Error: " + hostApplyBuf.String()) diff --git a/pkg/provisioner/provisioner/install.go b/pkg/provisioner/provisioner/install.go index 8d6b37a..0696b10 100644 --- a/pkg/provisioner/provisioner/install.go +++ b/pkg/provisioner/provisioner/install.go @@ -9,9 +9,10 @@ import ( "github.com/tmax-cloud/hypersds-operator/pkg/common/util" "github.com/tmax-cloud/hypersds-operator/pkg/common/wrapper" "github.com/tmax-cloud/hypersds-operator/pkg/provisioner/node" - "github.com/tmax-cloud/hypersds-operator/pkg/provisioner/osd" ) +const cmdNotFound = "command not found" + func (p *Provisioner) updateCephClusterToOp() error { fmt.Println("\n----------------Start to update conf and keyring to operator---------------") @@ -42,305 +43,65 @@ func (p *Provisioner) updateCephClusterToOp() error { return err } -func bootstrapCephadm(targetNode *node.Node, pathConfFromCr string) error { - fmt.Println("\n----------------Start to bootstrap ceph---------------") - - fmt.Println("[bootstrapCephadm] copying conf file") - err := createDstDir(targetNode, pathConfFromCr) - if err != nil { - return err - } - err = copyFile(targetNode, node.DESTINATION, pathConfFromCr, pathConfFromCr) +func isCephadmInstalled(deployNode *node.Node) (bool, error) { + const checkCephadmInstalledCmd = "cephadm version" + output, err := deployNode.RunSSHCmd(wrapper.SSHWrapper, checkCephadmInstalledCmd) if err != nil { - return err - } - - deployNodeHostSpec := targetNode.GetHostSpec() - - monIP := deployNodeHostSpec.GetAddr() - - fmt.Println("[bootstrapCephadm] executing bootstrap") - admBootstrapCmd := fmt.Sprintf("cephadm --image %s bootstrap --mon-ip %s --config %s", - cephImageName, monIP, pathConfFromCr) - err = processCmdOnNode(targetNode, admBootstrapCmd) - if err != nil { - return err + if strings.Contains(output.String(), cmdNotFound) { + return false, nil + } + fmt.Println("[Provisioner] cephadm installation check is failed") + return false, err } - - fmt.Println("[bootstrapCephadm] checking status") - const admHealthCheckCmd = "cephadm shell -- ceph -s" - err = processCmdOnNode(targetNode, admHealthCheckCmd) - - return err + return true, nil } func installCephadm(targetNode *node.Node) error { - fmt.Println("\n----------------Start to install cephadm---------------") - - fmt.Println("[installCephadm] executing curl cephadm") - curlCephadmCmd := fmt.Sprintf("curl --silent --remote-name --location https://github.com/ceph/ceph/raw/v%s/src/cephadm/cephadm", cephVersion) - err := processCmdOnNode(targetNode, curlCephadmCmd) + installed, err := isCephadmInstalled(targetNode) if err != nil { return err } - - fmt.Println("[installCephadm] executing chmod") - const chmodCmd = "chmod +x cephadm" - err = processCmdOnNode(targetNode, chmodCmd) - if err != nil { - return err + if installed { + return nil } - - // TODO: Specify release version - fmt.Println("[installCephadm] executing cephadm add-repo") - admAddRepoCmd := fmt.Sprintf("./cephadm add-repo --version %s", cephVersion) - err = processCmdOnNode(targetNode, admAddRepoCmd) - if err != nil { + if err := installCephadmPackage(targetNode); err != nil { return err } - - if targetNode.GetOs().Distro == node.Ubuntu { - // need for ubuntu 18.04 & ceph version 15.2.8 - fmt.Println("[installCephadm] executing curl cephadm gpg key") - addCephadmRepoCmd := "curl https://download.ceph.com/keys/release.asc | " + - "gpg --no-default-keyring --keyring /tmp/fix.gpg --import - && " + - "gpg --no-default-keyring --keyring /tmp/fix.gpg --export > /etc/apt/trusted.gpg.d/ceph.release.gpg && rm /tmp/fix.gpg" - err = processCmdOnNode(targetNode, addCephadmRepoCmd) - if err != nil { - return err - } - - fmt.Println("[installCephadm] executing cephadm apt-get update") - const aptUpdateCmd = "apt-get update" - err = processCmdOnNode(targetNode, aptUpdateCmd) - if err != nil { - return err - } - } - - fmt.Println("[installCephadm] executing cephadm install") - const admInstallCmd = "./cephadm install" - err = processCmdOnNode(targetNode, admInstallCmd) - if err != nil { - return err - } - - fmt.Println("[installCephadm] executing mkdir") - const mkdirCmd = "mkdir -p /etc/ceph" - err = processCmdOnNode(targetNode, mkdirCmd) - - return err -} - -func fetchOSInfo(targetNodeList []*node.Node) error { - for _, n := range targetNodeList { - if err := node.GetDistro(n); err != nil { - return err - } - } - return nil } -func installBasePackage(targetNodeList []*node.Node) error { - var err error - fmt.Println("\n----------------Start to install base package---------------") - - fmt.Println("[installBasePackage] executing packager update") - const updateCmd = "update -y" - for _, n := range targetNodeList { - packager := n.GetOs().Packager - if packager != node.Apt { - continue - } - fmt.Println(string(packager) + updateCmd) - err = processCmdOnNode(n, string(packager)+updateCmd) - if err != nil { - return err - } - } - - fmt.Println("[installBasePackage] install common packages") - const installCmd = "install -y " - const ntpPkgs = "chrony python3" - for _, n := range targetNodeList { - packager := n.GetOs().Packager - fmt.Println(string(packager) + installCmd + ntpPkgs) - err = processCmdOnNode(n, string(packager)+installCmd+ntpPkgs) - if err != nil { - return err - } - } - - fmt.Println("[installBasePackage] install docker dependencies") - const dockerDepAptPkgs = "apt-transport-https ca-certificates curl gnupg lsb-release" - const dockerDepYumPkgs = "yum-utils" - for _, n := range targetNodeList { - pkgs := "" - packager := n.GetOs().Packager - switch packager { - case node.Apt: - pkgs += dockerDepAptPkgs - case node.Yum: - pkgs += dockerDepYumPkgs - } - fmt.Println(string(packager) + installCmd + pkgs) - err = processCmdOnNode(n, string(packager)+installCmd+pkgs) - if err != nil { - return err - } - } - - fmt.Println("[installBasePackage] executing curl docker ...") - const dockerUbuntuGpgKeyCmd = "curl -s https://download.docker.com/linux/ubuntu/gpg | apt-key add - &>/dev/null" - for _, n := range targetNodeList { - osName := n.GetOs().Distro - if osName == node.Ubuntu { - err = processCmdOnNode(n, dockerUbuntuGpgKeyCmd) - if err != nil { - return err - } - } +func installPackages(nodeList []*node.Node) error { + if err := fetchOSInfo(nodeList); err != nil { + return err } - fmt.Println("[installBasePackage] executing add-apt-repo docker ...") - const dockerAptRepoCmd = `add-apt-repository "deb [arch=amd64] https://download.docker.com/linux/ubuntu $(lsb_release -cs) stable"` - const dockerYumRepoCmd = "yum-config-manager --add-repo https://download.docker.com/linux/centos/docker-ce.repo" - for _, n := range targetNodeList { - addDockerRepoCmd := "" - packager := n.GetOs().Packager - switch packager { - case node.Apt: - addDockerRepoCmd += dockerAptRepoCmd - case node.Yum: - addDockerRepoCmd += dockerYumRepoCmd + for _, n := range nodeList { + installed, err := isDockerInstalled(n) + if installed { + continue } - err = processCmdOnNode(n, addDockerRepoCmd) if err != nil { return err } - } - - fmt.Println("[installBasePackage] executing install docker-ce") - const dockerPkgs = "docker-ce" - for _, n := range targetNodeList { - packager := n.GetOs().Packager - if packager == node.Apt { - err = processCmdOnNode(n, string(packager)+updateCmd) - if err != nil { - return err - } - } - err = processCmdOnNode(n, string(packager)+installCmd+dockerPkgs) + err = installCommonPackage(n) if err != nil { return err } - } - - fmt.Println("[installBasePackage] executing sysctl docker") - const restartDockerCmd = "systemctl restart docker" - for _, n := range targetNodeList { - err = processCmdOnNode(n, restartDockerCmd) + err = installDocker(n) if err != nil { return err } } - return nil } -func (p *Provisioner) applyOsd(cephConf, cephKeyring []byte) error { - var err error - - fmt.Println("[applyOsd] get osds from CephOrch") - - cephName := p.getCephName() - pathConfigDir := p.getPathConfigDir() - - cmd := []string{"orch", "ls", "--service_type", "osd", "--export", "--refresh"} - output, err := util.RunCephCmd(wrapper.OsWrapper, wrapper.ExecWrapper, wrapper.IoUtilWrapper, cephConf, cephKeyring, cephName, cmd...) - if err != nil { - return processExecError(err, output) - } - - var osdsFromOrch []*osd.Osd - if !strings.Contains(output.String(), "No services reported") { - rawOsdsFromOrch := output.Bytes() - - osdsFromOrch, err = osd.NewOsdsFromCephOrch(wrapper.YamlWrapper, rawOsdsFromOrch) - if err != nil { - return err - } - } - osdsFromCephCr, err := osd.NewOsdsFromCephCr(p.cephCluster) - if err != nil { - return err - } - - var osdMap map[string]*osd.Osd - var removeOsdMap map[string]bool - - osdMap = make(map[string]*osd.Osd) - removeOsdMap = make(map[string]bool) - - for _, osdOrch := range osdsFromOrch { - osdService := osdOrch.GetService() - - osdServiceID := osdService.GetServiceID() - - osdMap[osdServiceID] = osdOrch - removeOsdMap[osdServiceID] = true - } - - fmt.Println("[applyOsd] compare osds between CephCR and CephOrch") - - for _, osdCephCr := range osdsFromCephCr { - osdService := osdCephCr.GetService() - osdServiceID := osdService.GetServiceID() - osdOrch, exist := osdMap[osdServiceID] - - var addDeviceList, removeDeviceList []string - - if exist { - addDeviceList, removeDeviceList, err = osdOrch.CompareDataDevices(osdCephCr) - if err != nil { - return err - } - removeOsdMap[osdServiceID] = false - // todo remove disk .... - fmt.Printf("[applyOsd] osd service: %s, add: %+q, remove: %+q\n", osdServiceID, addDeviceList, removeDeviceList) - } - - fmt.Println("[applyOsd] make osd yaml") - - osdFileName := pathConfigDir + osdServiceID + ".yaml" - err = osdCephCr.MakeYmlFile(wrapper.YamlWrapper, wrapper.IoUtilWrapper, osdFileName) - if err != nil { +func fetchOSInfo(targetNodeList []*node.Node) error { + for _, n := range targetNodeList { + if err := node.GetDistro(n); err != nil { return err } - - fmt.Printf("[applyOsd] apply osd service: %s\n", osdServiceID) - - applyCmd := []string{"orch", "apply", "-i", osdFileName} - output, err = util.RunCephCmd(wrapper.OsWrapper, wrapper.ExecWrapper, wrapper.IoUtilWrapper, cephConf, cephKeyring, cephName, applyCmd...) - if err != nil { - return processExecError(err, output) - } } - for osdServiceID, value := range removeOsdMap { - if !value { - continue - } - osdServiceName := "osd." + osdServiceID - fmt.Printf("[applyOsd] remove osd service: %s\n", osdServiceName) - - removeCmd := []string{"orch", "rm", osdServiceName} - output, err = util.RunCephCmd(wrapper.OsWrapper, wrapper.ExecWrapper, wrapper.IoUtilWrapper, cephConf, cephKeyring, cephName, removeCmd...) - if err != nil { - return processExecError(err, output) - } - } return nil } diff --git a/pkg/provisioner/provisioner/provisioner.go b/pkg/provisioner/provisioner/provisioner.go index b82ef7a..8cbc282 100644 --- a/pkg/provisioner/provisioner/provisioner.go +++ b/pkg/provisioner/provisioner/provisioner.go @@ -6,12 +6,12 @@ import ( "github.com/tmax-cloud/hypersds-operator/pkg/common/wrapper" "github.com/tmax-cloud/hypersds-operator/pkg/provisioner/config" "github.com/tmax-cloud/hypersds-operator/pkg/provisioner/node" + v1 "k8s.io/api/apps/v1" corev1 "k8s.io/api/core/v1" kubeerrors "k8s.io/apimachinery/pkg/api/errors" "k8s.io/apimachinery/pkg/types" "context" - "errors" "fmt" "strings" @@ -21,7 +21,6 @@ import ( // Provisioner contains information for ceph deployment and current ceph deployment status type Provisioner struct { cephCluster hypersdsv1alpha1.CephClusterSpec - state provisionerState cephNamespace string cephName string @@ -29,10 +28,6 @@ type Provisioner struct { cephConfig *config.CephConfig } -func (p *Provisioner) getState() provisionerState { - return p.state -} - func (p *Provisioner) getNodes() ([]*node.Node, error) { return node.NewNodesFromCephCr(p.cephCluster) } @@ -52,12 +47,6 @@ func (p *Provisioner) getPathConfigDir() string { return util.PathGlobalConfigDir + p.cephName + "/" } -//nolint:unparam // now, setMethod always return nil, if setMethod needs return other value, deletes nolint and implement -func (p *Provisioner) setState(state provisionerState) error { - p.state = state - return nil -} - //nolint:unparam // now, setMethod always return nil, if setMethod needs return other value, deletes nolint and implement func (p *Provisioner) setCephCluster(cephCluster hypersdsv1alpha1.CephClusterSpec) error { p.cephCluster = cephCluster @@ -90,219 +79,202 @@ func (p *Provisioner) setCephConfig() error { // Run executes the following ceph deployment steps according to the current ceph deployment state func (p *Provisioner) Run() error { - // Decide deploying node (currently, first node is deploying node) - nodeList, err := p.getNodes() + // Create config object from Ceph CR + err := p.setCephConfig() if err != nil { return err } - deployNode := nodeList[0] - // Create config object from Ceph CR - err = p.setCephConfig() + nodeList, err := p.getNodes() if err != nil { return err } - cephConfig := p.getCephConfig() - pathConfigDir := p.getPathConfigDir() - pathConfFromCr := pathConfigDir + cephConfNameFromCr - pathConf := pathConfigDir + util.CephConfName - pathKeyring := pathConfigDir + util.CephKeyringName - - switch currentState := p.getState(); currentState { - case InstallStarted: - // Fetch OS information for each nodes - err = fetchOSInfo(nodeList) - if err != nil { - return err - } - // Install base package to all nodes - err = installBasePackage(nodeList) - if err != nil { - return err - } - - // Set provisioner state to BasePkgInstalled - err = p.setState(BasePkgInstalled) - if err != nil { - return err - } + err = installPackages(nodeList) + if err != nil { + return err + } - fallthrough + deployNode, err := p.getDeployNode() + if err != nil { + return err + } + err = p.bootstrap(deployNode) + if err != nil { + return err + } - case BasePkgInstalled: - // Install cephadm package to deploying node - err = installCephadm(deployNode) - if err != nil { - return err - } + err = p.updateKubeObject(deployNode) + if err != nil { + return err + } - // Set provisioner state to CephadmPkgInstalled - err = p.setState(CephadmPkgInstalled) - if err != nil { - return err + cephConf, cephKeyring, err := p.getCephData() + if err != nil { + if kubeerrors.IsNotFound(err) { + if err2 := p.updateKubeObject(deployNode); err2 != nil { + return err + } } + return err + } - fallthrough - - case CephadmPkgInstalled: - // Extract initial conf file of Ceph - err = cephConfig.MakeIniFile(wrapper.IoUtilWrapper, pathConfFromCr) - if err != nil { - return err - } + err = p.applyHost(wrapper.YamlWrapper, wrapper.ExecWrapper, wrapper.IoUtilWrapper, cephConf, cephKeyring) + if err != nil { + return err + } + err = p.applyOsd(cephConf, cephKeyring) + if err != nil { + return err + } - // Bootstrap ceph on deploy node with cephadm - err = bootstrapCephadm(deployNode, pathConfFromCr) - if err != nil { - return err - } + pathConfigDir := p.getPathConfigDir() + err = wrapper.OsWrapper.RemoveAll(pathConfigDir) + if err != nil { + fmt.Println("[Provisioner] CephName Directory Remove Error") + return err + } - // Set provisioner state to CephBootstrapped - err = p.setState(CephBootstrapped) - if err != nil { - return err - } + return nil +} - fallthrough +func (p *Provisioner) getDeployNode() (*node.Node, error) { + deployNode, err := p.fetchDeployNodeFromK8sStatus() + if err != nil { + return nil, err + } - case CephBootstrapped: - // Copy conf and keyring from deploy node - err = copyFile(deployNode, node.SOURCE, pathConfFromAdm, pathConf) - if err != nil { - return err - } - err = copyFile(deployNode, node.SOURCE, pathKeyringFromAdm, pathKeyring) - if err != nil { - return err - } + return deployNode, nil +} - // Update conf and keyring to ConfigMap and Secret - err = p.updateCephClusterToOp() - if err != nil { - return err - } +func (p *Provisioner) fetchDeployNodeFromK8sStatus() (*node.Node, error) { + clientSet := p.getClientSet() + cephNamespace := p.getCephNamespace() + cephName := p.getCephName() - // Set provisioner state to CephBootstrapCommitted - err = p.setState(CephBootstrapCommitted) - if err != nil { - return err - } + cephCluster := &hypersdsv1alpha1.CephCluster{} + if err := clientSet.Get(context.TODO(), types.NamespacedName{Namespace: cephNamespace, Name: cephName}, cephCluster); err != nil { + return nil, err + } - fallthrough + deployNode, err := convertNodeType(cephCluster.Status.DeployNode) + if err != nil { + return nil, err + } - case CephBootstrapCommitted: - //nolint:govet // err need to be redeclared - // Copy conf and keyring from deploy node for ceph-common - cephConf, cephKeyring, err := p.checkKubeObjectUpdated() - if err != nil { - return err - } - // Can not proceed. Need to reconcile again - if cephConf == nil || cephKeyring == nil { - return errors.New("k8s ConfigMap and Secret is not updated yet") - } + return deployNode, nil +} - err = p.applyHost(wrapper.YamlWrapper, wrapper.ExecWrapper, wrapper.IoUtilWrapper, cephConf, cephKeyring) - if err != nil { - return err - } +func convertNodeType(k8sNode hypersdsv1alpha1.Node) (*node.Node, error) { + hostSpec := node.HostSpec{ + Addr: k8sNode.IP, + HostName: k8sNode.HostName, + } - err = p.applyOsd(cephConf, cephKeyring) - if err != nil { - return err - } - err = p.setState(CephOsdDeployed) - if err != nil { - return err - } + var n node.Node + var err error + err = n.SetUserID(k8sNode.UserID) + if err != nil { + return nil, err } - err = wrapper.OsWrapper.RemoveAll(pathConfigDir) + err = n.SetUserPw(k8sNode.Password) if err != nil { - fmt.Println("[Provisioner] CephName Driectory Remove Error") - return err + return nil, err + } + err = n.SetHostSpec(&hostSpec) + if err != nil { + return nil, err } - return nil + return &n, nil } -func (p *Provisioner) identifyProvisionerState() (provisionerState, error) { - // Decide deploying node (currently, first node is deploying node) - nodes, err := p.getNodes() +func (p *Provisioner) bootstrap(deployNode *node.Node) error { + bootstrapped, err := isBootstrapped(deployNode) if err != nil { - return "", err + return err + } + if bootstrapped { + return nil } - deployNode := nodes[0] - // Check base pkgs are installed - // TODO: May contain error if user removed docker but did not purge dpkg - const checkDockerWorkingCmd = "docker -v" - _, err = deployNode.RunSSHCmd(wrapper.SSHWrapper, checkDockerWorkingCmd) - // It considers any error that base pkgs are not installed + err = installCephadm(deployNode) if err != nil { - // TODO: Replace stdout to log out - fmt.Println("[identifyProvisionerState] docker is not installed") - return InstallStarted, nil + return err } - // Check Cephadm is installed - const checkCephadmInstalledCmd = "cephadm version" - outputCephadm, err := deployNode.RunSSHCmd(wrapper.SSHWrapper, checkCephadmInstalledCmd) + cephConfig := p.getCephConfig() + pathConfigDir := p.getPathConfigDir() + pathConfFromCr := pathConfigDir + cephConfNameFromCr + + // Extract initial conf file of Ceph + err = cephConfig.MakeIniFile(wrapper.IoUtilWrapper, pathConfFromCr) if err != nil { - const cephadmNotFound = "command not found" - if strings.Contains(outputCephadm.String(), cephadmNotFound) { - return BasePkgInstalled, nil - } + return err + } - // Other error occurred on RunSSHCmd - // TODO: Replace stdout to log out - fmt.Println("[identifyProvisionerState] cephadm installation check is failed") - return BasePkgInstalled, err + // Bootstrap ceph on deploy node with cephadm + err = bootstrapCephadm(deployNode, pathConfFromCr) + if err != nil { + return err } + return nil +} - // Check Ceph is bootstrapped +func isBootstrapped(deployNode *node.Node) (bool, error) { const checkBootstrappedCmd = "cephadm shell -- ceph -s" - outputBootstrap, err := deployNode.RunSSHCmd(wrapper.SSHWrapper, checkBootstrappedCmd) + output, err := deployNode.RunSSHCmd(wrapper.SSHWrapper, checkBootstrappedCmd) if err != nil { const objectNotFound = "ObjectNotFound" - if strings.Contains(outputBootstrap.String(), objectNotFound) { - return CephadmPkgInstalled, nil + if strings.Contains(output.String(), objectNotFound) || strings.Contains(output.String(), cmdNotFound) { + return false, nil } - - // Other error occurred on RunSSHCmd - // TODO: Replace stdout to log out - fmt.Println("[identifyProvisionerState] ceph bootstrap check is failed") - return CephadmPkgInstalled, err + // Confirm some Ceph health is returned + const cephHealthPrefix = "HEALTH_ERR" + if strings.Contains(output.String(), cephHealthPrefix) { + fmt.Println("[Provisioner] ceph cluster status is error") + return false, err + } + fmt.Println("[Provisioner] ceph bootstrap check is failed") + return false, err } + return true, nil +} - // Confirm some Ceph health is returned - const cephHealthPrefix = "HEALTH_" - if !strings.Contains(outputBootstrap.String(), cephHealthPrefix) { - // TODO: Replace stdout to log out - fmt.Println("[identifyProvisionerState] ceph status does not return HEALTH_*") +func (p *Provisioner) updateKubeObject(deployNode *node.Node) error { + pathConfigDir := p.getPathConfigDir() + pathConf := pathConfigDir + util.CephConfName + pathKeyring := pathConfigDir + util.CephKeyringName - // TODO: Make own error pkg of hypersds-provisioner - return CephadmPkgInstalled, - errors.New("Error on Ceph bootstrap, cech status result: \n" + - outputBootstrap.String()) + // Copy conf and keyring from deploy node + err := copyFile(deployNode, node.SOURCE, pathConfFromAdm, pathConf) + if err != nil { + return err + } + err = copyFile(deployNode, node.SOURCE, pathKeyringFromAdm, pathKeyring) + if err != nil { + return err } - // Check Ceph bootstrap is committed - conf, keyring, err := p.checkKubeObjectUpdated() + // Update conf and keyring to ConfigMap and Secret + err = p.updateCephClusterToOp() if err != nil { - // Other error occurred on checkKubeObjectUpdated - // TODO: Replace stdout to log out - fmt.Println("[identifyProvisionerState] k8s configmap and secret check is failed") - return CephBootstrapped, err + return err } + return nil +} +func (p *Provisioner) getCephData() (confDataBuf, keyringDataBuf []byte, err error) { + conf, keyring, err := p.checkKubeObjectUpdated() + if err != nil { + fmt.Println("[Provisioner] k8s configmap and secret check is failed") + return conf, keyring, err + } if conf == nil || keyring == nil { - // TODO: Replace stdout to log out - fmt.Println("[identifyProvisionerState] k8s configmap and secret are not updated") - return CephBootstrapped, nil + fmt.Println("[Provisioner] k8s configmap and secret are not updated") + return conf, keyring, kubeerrors.NewNotFound(v1.Resource("CephCluster"), p.cephName) } - - return CephBootstrapCommitted, nil + return conf, keyring, nil } // TODO: Replace config const to inputs (e.g. K8sConfigMap, etc) @@ -314,7 +286,6 @@ func (p *Provisioner) checkKubeObjectUpdated() (confDataBuf, keyringDataBuf []by if kubeerrors.IsNotFound(err) { // TODO: Replace stdout to log out fmt.Println("ConfigMap must exist") - return nil, nil, err } return nil, nil, err } @@ -331,7 +302,6 @@ func (p *Provisioner) checkKubeObjectUpdated() (confDataBuf, keyringDataBuf []by if kubeerrors.IsNotFound(err) { // TODO: Replace stdout to log out fmt.Println("Secret must exist") - return nil, nil, err } return nil, nil, err } @@ -347,7 +317,6 @@ func (p *Provisioner) checkKubeObjectUpdated() (confDataBuf, keyringDataBuf []by // NewProvisioner creates Provisioner using ceph deployment information and checks the current ceph deployment status in node func NewProvisioner(cephClusterSpec hypersdsv1alpha1.CephClusterSpec, clientSet client.Client, cephNamespace, cephName string) (*Provisioner, error) { var err error - var currentState provisionerState provisionerInstance := &Provisioner{} @@ -382,23 +351,7 @@ func NewProvisioner(cephClusterSpec hypersdsv1alpha1.CephClusterSpec, clientSet err = wrapper.OsWrapper.MkdirAll(pathConfigDir, 0644) if err != nil { - fmt.Println("[Provisioner] CephName Driectory Create Error") - return nil, err - } - - // identifyProvisionerState is only called once, on init - // No one is allowed to modify provisionerState - currentState, err = provisionerInstance.identifyProvisionerState() - if err != nil { - // TODO: Replace stdout to log out - fmt.Println("[Provisioner] identifyProvisionerState Error") - return nil, err - } - - err = provisionerInstance.setState(currentState) - if err != nil { - // TODO: Replace stdout to log out - fmt.Println("[Provisioner] setState Error") + fmt.Println("[Provisioner] CephName Directory Create Error") return nil, err }