diff --git a/Makefile b/Makefile index 347496c34..8a9a10f90 100644 --- a/Makefile +++ b/Makefile @@ -8,12 +8,17 @@ default: all phony := all all: build -phony += build server authtool client cli libsdk fsck fdstore preload bcache blobstore -build: server authtool client cli libsdk fsck fdstore preload bcache blobstore +phony += build server authtool client cli libsdk fsck fdstore preload bcache blobstore deploy +build: server authtool client cli libsdk fsck fdstore preload bcache blobstore deploy server: @build/build.sh server $(GOMOD) --threads=$(threads) + +deploy: + @build/build.sh deploy $(GOMOD) --threads=$(threads) + + blobstore: @build/build.sh blobstore $(GOMOD) --threads=$(threads) diff --git a/build/build.sh b/build/build.sh index db94e0676..b50188e7b 100755 --- a/build/build.sh +++ b/build/build.sh @@ -333,6 +333,16 @@ build_cli() { popd >/dev/null } + + +build_cfs_deploy() { + #cfs_deploy need gorocksdb too + pushd $SrcPath >/dev/null + echo -n "build cfs-deploy " + CGO_ENABLED=1 go build ${MODFLAGS} -gcflags=all=-trimpath=${SrcPath} -asmflags=all=-trimpath=${SrcPath} -ldflags="${LDFlags}" -o ${BuildBinPath}/cfs-deploy ${SrcPath}/deploy/*.go && echo "success" || echo "failed" + popd >/dev/null +} + build_fsck() { pushd $SrcPath >/dev/null echo -n "build cfs-fsck " @@ -470,6 +480,9 @@ case "$cmd" in "cli") build_cli ;; + "deploy") + build_cfs_deploy + ;; "fsck") build_fsck ;; diff --git a/deploy/cmd/cluster.go b/deploy/cmd/cluster.go new file mode 100644 index 000000000..2a5253b07 --- /dev/null +++ b/deploy/cmd/cluster.go @@ -0,0 +1,291 @@ +package cmd + +import ( + "fmt" + "log" + "net" + "os" + "strconv" + "strings" + + "github.com/spf13/cobra" +) + +var ClusterCmd = &cobra.Command{ + Use: "cluster", + Short: "cluster manager", + Long: `This command will manager the cluster.`, + + Run: func(cmd *cobra.Command, args []string) { + fmt.Println(cmd.UsageString()) + }, +} + +var initCommand = &cobra.Command{ + Use: "init", + Short: "init the cluster from config.yaml", + Long: "init the cluster from config.yaml", + + Run: func(cmd *cobra.Command, args []string) { + initCluster() + }, +} + +var infoCommand = &cobra.Command{ + Use: "info", + Short: "Display cluster information", + Long: "Display cluster information", + Run: func(cmd *cobra.Command, args []string) { + err := infoOfCluster() + if err != nil { + log.Println(err) + } + }, +} + +var clearCommand = &cobra.Command{ + Use: "clear", + Short: "Clear cluster files and information", + Long: "Clear cluster files and information", + Run: func(cmd *cobra.Command, args []string) { + cleanContainerAndLogs() + }, +} + +var configCommand = &cobra.Command{ + Use: "config", + Short: "Loading configuration files into the cluster", + Long: "Loading configuration files into the cluster", + Run: func(cmd *cobra.Command, args []string) { + err := convertToJosn() + if err != nil { + log.Println(err) + } + + }, +} + +// Obtain the IP address of the current host +func getCurrentIP() (string, error) { + // Get Host Name + hostname, err := os.Hostname() + if err != nil { + return "", err + } + + addrs, err := net.LookupIP(hostname) + if err != nil { + return "", err + } + + for _, addr := range addrs { + if ipv4 := addr.To4(); ipv4 != nil { + + return ipv4.String(), nil + + } + } + + return "", fmt.Errorf("IPv4 address not found") +} + +func printTable(services []Service) { + fmt.Println("Server Type | Container Name | Node IP | Status") + fmt.Println("-----------------------------------------------------") + for _, service := range services { + fmt.Printf("%-12s | %-14s | %-13s | %s\n", service.ServerType, service.ContainerName, service.NodeIP, service.Status) + } +} + +func infoOfCluster() error { + config, err := readConfig() + if err != nil { + log.Fatal(err) + } + servers := []Service{} + + for id, node := range config.DeployHostsList.Master.Hosts { + server := Service{} + server.NodeIP = node + server.ServerType = MasterServer + server.Status = Stopped + server.ContainerName = MasterName + strconv.Itoa(id+1) + + ps, _ := containerStatus(RemoteUser, node, MasterName+strconv.Itoa(id+1)) + + if strings.Contains(ps, "running") { + server.Status = Running + } + servers = append(servers, server) + + } + + for id, node := range config.DeployHostsList.MetaNode.Hosts { + server := Service{} + server.NodeIP = node + server.ServerType = MetaNodeServer + server.Status = Stopped + server.ContainerName = MetaNodeName + strconv.Itoa(id+1) + ps, _ := containerStatus(RemoteUser, node, MetaNodeName+strconv.Itoa(id+1)) + + if strings.Contains(ps, "running") { + server.Status = Running + } + servers = append(servers, server) + } + + for id, node := range config.DeployHostsList.DataNode { + server := Service{} + server.NodeIP = node.Hosts + server.ServerType = DataNodeServer + server.Status = Stopped + server.ContainerName = DataNodeName + strconv.Itoa(id+1) + + ps, _ := containerStatus(RemoteUser, node.Hosts, DataNodeName+strconv.Itoa(id+1)) + + if strings.Contains(ps, "running") { + server.Status = Running + } + servers = append(servers, server) + } + + printTable(servers) + return nil +} + +func removeDuplicates(slice []string) []string { + encountered := map[string]bool{} + result := []string{} + + for _, item := range slice { + if encountered[item] == true { + continue + } + encountered[item] = true + result = append(result, item) + } + + return result +} +func initCluster() { + + config, err := readConfig() + if err != nil { + log.Fatal(err) + } + + hosts := []string{} + hosts = append(hosts, config.DeployHostsList.Master.Hosts...) + hosts = append(hosts, config.DeployHostsList.MetaNode.Hosts...) + for i := 0; i < len(config.DeployHostsList.DataNode); i++ { + hosts = append(hosts, config.DeployHostsList.DataNode[i].Hosts) + } + + newHosts := removeDuplicates(hosts) + + // Obtain the IP address of the current host + currentNode, err := getCurrentIP() + if err != nil { + log.Fatal(err) + } + + log.Println("The IP address of the current host:", currentNode) + + // Establish a secure connection from the current node to other nodes + for _, node := range newHosts { + if node == currentNode || node == "" { + continue + } + err := establishSSHConnectionWithoutPassword(currentNode, RemoteUser, node) + if err != nil { + log.Fatal(err) + } + } + + log.Println("Password free connection establishment completed") + + for _, node := range newHosts { + // Check if Docker is installed and installed + if node == "" { + continue + } + checkAndInstallDocker(RemoteUser, node) + + // Check if the Docker service is started and started + err = checkAndStartDockerService(RemoteUser, node) + if err != nil { + log.Printf("Failed to start Docker service on node% s:% v", node, err) + } else { + log.Printf("The docker for node %s is ready", node) + } + + // Pull Mirror + err = pullImageOnNode(RemoteUser, node, config.Global.ContainerImage) + if err != nil { + log.Printf("Failed to pull mirror% s on node% s:% v", node, config.Global.ContainerImage, err) + } else { + log.Printf("Successfully pulled mirror % s on node % s", config.Global.ContainerImage, node) + } + + //create dir + err = createRemoteFolder(RemoteUser, node, config.Global.DataDir) + //check firewall + if err != nil { + log.Println(err) + } + err = transferDirectoryToRemote(BinDir, config.Global.DataDir, RemoteUser, node) + if err != nil { + log.Println(err) + } + err = transferDirectoryToRemote(ScriptDir, config.Global.DataDir, RemoteUser, node) + if err != nil { + log.Println(err) + } + //create conf dir + err = createRemoteFolder(RemoteUser, node, config.Global.DataDir+"/conf") + //check firewall + if err != nil { + log.Println(err) + } + stopFirewall(RemoteUser, node) + + } + + log.Println("*******Cluster environment initialization completed******") +} + +func cleanContainerAndLogs() { + config, err := readConfig() + if err != nil { + log.Fatal(err) + } + hosts := []string{} + hosts = append(hosts, config.DeployHostsList.Master.Hosts...) + hosts = append(hosts, config.DeployHostsList.MetaNode.Hosts...) + for i := 0; i < len(config.DeployHostsList.DataNode); i++ { + hosts = append(hosts, config.DeployHostsList.DataNode[i].Hosts) + } + + newHosts := removeDuplicates(hosts) + + // only remove container image + for _, node := range newHosts { + if node == "" { + continue + } + err := removeImageOnNode(RemoteUser, node, config.Global.ContainerImage) + if err != nil { + log.Fatalln(err) + } + } + +} + +func init() { + ClusterCmd.AddCommand(initCommand) + ClusterCmd.AddCommand(infoCommand) + ClusterCmd.AddCommand(clearCommand) + ClusterCmd.AddCommand(configCommand) + configCommand.PersistentFlags().StringVarP(&ConfigFileName, "file", "f", "", "Specify the location of the configuration file relative to cfs deploy") + +} diff --git a/deploy/cmd/config.go b/deploy/cmd/config.go new file mode 100644 index 000000000..11cf8480b --- /dev/null +++ b/deploy/cmd/config.go @@ -0,0 +1,133 @@ +package cmd + +import ( + "io/ioutil" + "log" + "strconv" + + "gopkg.in/yaml.v2" +) + +type Config struct { + Global GlobalConfig `yaml:"global"` + Master MasterConfig `yaml:"master"` + MetaNode MetaNodeConfig `yaml:"metanode"` + DataNode DataNodeConfig `yaml:"datanode"` + DeployHostsList DeployHostsListConfig `yaml:"deplopy_hosts_list"` +} + +type GlobalConfig struct { + ContainerImage string `yaml:"container_image"` + DataDir string `yaml:"data_dir"` + Variable struct { + Target string `yaml:"target"` + } `yaml:"variable"` +} + +type MasterConfig struct { + Config struct { + Listen string `yaml:"listen"` + Prof string `yaml:"prof"` + DataDir string `yaml:"data_dir"` + } `yaml:"config"` +} + +type MetaNodeConfig struct { + Config struct { + Listen string `yaml:"listen"` + Prof string `yaml:"prof"` + DataDir string `yaml:"data_dir"` + } `yaml:"config"` +} + +type DataNodeConfig struct { + Config struct { + Listen string `yaml:"listen"` + Prof string `yaml:"prof"` + DataDir string `yaml:"data_dir"` + } `yaml:"config"` +} + +type DeployHostsListConfig struct { + Master struct { + Hosts []string `yaml:"hosts"` + } `yaml:"master"` + MetaNode struct { + Hosts []string `yaml:"hosts"` + } `yaml:"metanode"` + DataNode []struct { + Hosts string `yaml:"hosts"` + Disk []DiskInfo `yaml:"disk"` + } `yaml:"datanode"` +} + +type DiskInfo struct { + Path string `yaml:"path"` + Size string `yaml:"size"` +} + +func readConfig() (*Config, error) { + + data, err := ioutil.ReadFile(ConfigFileName) + if err != nil { + log.Println("Unable to read configuration file:", err) + return nil, err + } + + config := &Config{} + err = yaml.Unmarshal(data, &config) + if err != nil { + log.Println("Unable to parse configuration file:", err) + return nil, err + } + return config, nil + +} + +func convertToJosn() error { + err := createFolder(ConfDir) + if err != nil { + return err + } + config, err := readConfig() + if err != nil { + return err + } + + for id, node := range config.DeployHostsList.Master.Hosts { + peers := getMasterPeers(config) + err := writeMaster(ClusterName, strconv.Itoa(id+1), node, config.Master.Config.Listen, config.Master.Config.Prof, peers) + if err != nil { + return err + } + } + + masterAddr, err := getMasterAddrAndPort() + if err != nil { + return err + } + for id, node := range config.DeployHostsList.MetaNode.Hosts { + err = writeMetaNode(config.MetaNode.Config.Listen, config.MetaNode.Config.Prof, strconv.Itoa(id+1), node, masterAddr) + if err != nil { + return err + } + } + + disksInfo := []string{} + for _, node := range config.DeployHostsList.DataNode { + diskMap := "" + for _, info := range node.Disk { + diskMap += " -v " + info.Path + ":/cfs" + info.Path + disksInfo = append(disksInfo, "/cfs"+info.Path+":"+info.Size) + } + } + + for id, node := range config.DeployHostsList.MetaNode.Hosts { + err = writeDataNode(config.DataNode.Config.Listen, config.DataNode.Config.Prof, strconv.Itoa(id+1), node, masterAddr, disksInfo) + if err != nil { + return err + } + } + + return nil +} diff --git a/deploy/cmd/datanode.go b/deploy/cmd/datanode.go new file mode 100644 index 000000000..e380c1f5c --- /dev/null +++ b/deploy/cmd/datanode.go @@ -0,0 +1,234 @@ +package cmd + +import ( + "encoding/json" + "fmt" + "io/ioutil" + "log" + "strconv" + "strings" +) + +type DataNode struct { + Role string `json:"role"` + Listen string `json:"listen"` + Prof string `json:"prof"` + RaftHeartbeat string `json:"raftHeartbeat"` + RaftReplica string `json:"raftReplica"` + RaftDir string `json:"raftDir"` + LocalIP string `json:"localIP"` + ConsulAddr string `json:"consulAddr"` + ExporterPort int `json:"exporterPort"` + Cell string `json:"cell"` + LogDir string `json:"logDir"` + LogLevel string `json:"logLevel"` + Disks []string `json:"disks"` + DiskIopsReadLimit string `json:"diskIopsReadLimit"` + DiskIopsWriteLimit string `json:"diskIopsWriteLimit"` + DiskFlowReadLimit string `json:"diskFlowReadLimit"` + DiskFlowWriteLimit string `json:"diskFlowWriteLimit"` + MasterAddr []string `json:"masterAddr"` + EnableSmuxConnPool bool `json:"enableSmuxConnPool"` +} + +func readDataNode(filename string) (*DataNode, error) { + data, err := ioutil.ReadFile(filename) + if err != nil { + return nil, err + } + + dataNode := &DataNode{} + err = json.Unmarshal(data, dataNode) + if err != nil { + return nil, err + } + return dataNode, nil +} + +func writeDataNode(listen, prof, id, localIP string, masterAddrs, disks []string) error { + // 将DataNode配置写入DataNode.json文件 + datanode := DataNode{ + Role: "datanode", + Listen: listen, + Prof: prof, + LocalIP: localIP, + RaftHeartbeat: "17330", + RaftReplica: "17340", + RaftDir: "/cfs/log", + ConsulAddr: "http://192.168.0.101:8500", + ExporterPort: 9500, + Cell: "cell-01", + LogDir: "/cfs/log", + LogLevel: "debug", + Disks: disks, + DiskIopsReadLimit: "20000", + DiskIopsWriteLimit: "5000", + DiskFlowReadLimit: "1024000000", + DiskFlowWriteLimit: "524288000", + MasterAddr: masterAddrs, + EnableSmuxConnPool: true, + } + + dataNodeData, err := json.MarshalIndent(datanode, "", " ") + if err != nil { + log.Println("Unable to encode DataNode configuration:", err) + return err + } + + err = ioutil.WriteFile(ConfDir+"/datanode"+id+".json", dataNodeData, 0644) + if err != nil { + log.Println("Unable to write to DataNode.json file:", err) + return err + } + return nil +} + +func startAllDataNode() error { + config, err := readConfig() + if err != nil { + log.Println(err) + } + disksInfo := []string{} + for _, node := range config.DeployHostsList.DataNode { + diskMap := "" + for _, info := range node.Disk { + diskMap += " -v " + info.Path + ":/cfs" + info.Path + //disksInfo = append(disksInfo, "/cfs"+info.Path+":"+info.Size) + } + disksInfo = append(disksInfo, diskMap) + } + + var dataDir string + if config.DataNode.Config.DataDir == "" { + dataDir = config.Global.DataDir + } else { + dataDir = config.DataNode.Config.DataDir + } + index := 0 + files, err := ioutil.ReadDir(ConfDir) + if err != nil { + return err + } + for _, file := range files { + if strings.HasPrefix(file.Name(), "datanode") && !file.IsDir() { + data, err := readDataNode(ConfDir + "/" + file.Name()) + if err != nil { + fmt.Printf("Error reading file %s: %s\n", file.Name(), err) + return nil + } + confFilePath := ConfDir + "/" + file.Name() + err = transferConfigFileToRemote(confFilePath, dataDir+"/conf", RemoteUser, data.LocalIP) + if err != nil { + return err + } + + err = checkAndDeleteContainerOnNode(RemoteUser, data.LocalIP, strings.Split(file.Name(), ".")[0]) + if err != nil { + return err + } + status, err := startDatanodeContainerOnNode(RemoteUser, data.LocalIP, strings.Split(file.Name(), ".")[0], dataDir, disksInfo[index]) + if err != nil { + return err + } + index++ + log.Println(status) + } + + } + + log.Println("start all datanode services") + return nil +} + +func startDatanodeInSpecificNode(node string) error { + + config, err := readConfig() + if err != nil { + return err + } + + var dataDir string + if config.DataNode.Config.DataDir == "" { + dataDir = config.Global.DataDir + } else { + dataDir = config.DataNode.Config.DataDir + } + for id, n := range config.DeployHostsList.DataNode { + if n.Hosts == node { + confFilePath := ConfDir + "/" + "datanode" + strconv.Itoa(id+1) + ".json" + err = transferConfigFileToRemote(confFilePath, dataDir+"/conf", RemoteUser, node) + if err != nil { + return err + } + + err = checkAndDeleteContainerOnNode(RemoteUser, node, DataNodeName+strconv.Itoa(id+1)) + if err != nil { + return err + } + diskMap := "" + for _, info := range n.Disk { + diskMap += " -v " + info.Path + ":/cfs" + info.Path + + } + status, err := startDatanodeContainerOnNode(RemoteUser, node, DataNodeName+strconv.Itoa(id+1), dataDir, diskMap) + if err != nil { + return err + } + + log.Println(status) + break + } + } + return nil +} + +func stopDatanodeInSpecificNode(node string) error { + config, err := readConfig() + if err != nil { + return err + } + for id, n := range config.DeployHostsList.DataNode { + if n.Hosts == node { + status, err := stopContainerOnNode(RemoteUser, node, DataNodeName+strconv.Itoa(id+1)) + if err != nil { + return err + } + log.Println(status) + status, err = rmContainerOnNode(RemoteUser, node, DataNodeName+strconv.Itoa(id+1)) + if err != nil { + return err + } + log.Println(status) + } + } + return nil +} + +func stopAllDataNode() error { + + files, err := ioutil.ReadDir(ConfDir) + if err != nil { + return err + } + for _, file := range files { + if strings.HasPrefix(file.Name(), "datanode") && !file.IsDir() { + data, err := readDataNode(ConfDir + "/" + file.Name()) + if err != nil { + fmt.Printf("Error reading file %s: %s\n", file.Name(), err) + return nil + } + status, err := stopContainerOnNode(RemoteUser, data.LocalIP, strings.Split(file.Name(), ".")[0]) + if err != nil { + return err + } + log.Println(status) + status, err = rmContainerOnNode(RemoteUser, data.LocalIP, strings.Split(file.Name(), ".")[0]) + if err != nil { + return err + } + log.Println(status) + + } + } + return nil +} diff --git a/deploy/cmd/docker.go b/deploy/cmd/docker.go new file mode 100644 index 000000000..6cf0f20f9 --- /dev/null +++ b/deploy/cmd/docker.go @@ -0,0 +1,206 @@ +package cmd + +import ( + "fmt" + "log" + "os" + "os/exec" + "strings" +) + +func stopContainerOnNode(nodeUser, node, containerName string) (string, error) { + if ok, _ := checkContainerExistence(nodeUser, node, containerName); ok { + cmd := exec.Command("ssh", nodeUser+"@"+node, "docker stop ", containerName) + _, err := cmd.Output() + if err != nil { + return fmt.Sprintf("failed stop %s on node %s", containerName, node), err + } + return fmt.Sprintf("successful stop %s on node %s", containerName, node), nil + } + return fmt.Sprintf("%s on node %s already stopped", containerName, node), nil + +} + +func rmContainerOnNode(nodeUser, node, containerName string) (string, error) { + if ok, _ := checkContainerExistence(nodeUser, node, containerName); ok { + cmd := exec.Command("ssh", nodeUser+"@"+node, "docker rm ", containerName) + _, err := cmd.Output() + if err != nil { + return fmt.Sprintf("failed rm %s on node %s", containerName, node), err + } + return fmt.Sprintf("successful rm %s on node %s", containerName, node), nil + } + return fmt.Sprintf("%s on node %s already removed", containerName, node), nil + +} + +func startMasterContainerOnNode(nodeUser, node, containerName, dataDir string) (string, error) { + cmd := exec.Command("ssh", nodeUser+"@"+node, + "docker run -d --name "+containerName+ + " -v "+dataDir+"/disk/"+containerName+"/data:/cfs/data"+ + " -v "+dataDir+"/bin"+":/cfs/bin:ro"+ + " -v "+dataDir+"/disk/"+containerName+"/log:/cfs/log"+ + " -v "+dataDir+"/conf/"+containerName+".json:/cfs/conf/master.json"+ + " -v "+dataDir+"/script/start_master.sh:/cfs/script/start.sh"+ + " --restart on-failure --privileged --network host "+ImageName+ + " /bin/sh /cfs/script/start.sh ") + _, err := cmd.Output() + if err != nil { + return fmt.Sprintf("failed start %s on node %s", containerName, node), err + } + return fmt.Sprintf("successfully started the container %s on node %s", containerName, node), nil +} + +func startMetanodeContainerOnNode(nodeUser, node, containerName, dataDir string) (string, error) { + cmd := exec.Command("ssh", nodeUser+"@"+node, + "docker run -d --name "+containerName+ + " -v "+dataDir+"/disk/"+containerName+"/data:/cfs/data"+ + " -v "+dataDir+"/bin"+":/cfs/bin:ro"+ + " -v "+dataDir+"/disk/"+containerName+"/log:/cfs/log"+ + " -v "+dataDir+"/conf/"+containerName+".json:/cfs/conf/metanode.json"+ + " -v "+dataDir+"/script/start_meta.sh:/cfs/script/start.sh"+ + " --restart on-failure --privileged --network host "+ImageName+ + " /bin/sh /cfs/script/start.sh ") + _, err := cmd.Output() + if err != nil { + return fmt.Sprintf("failed start %s on node %s", containerName, node), err + } + return fmt.Sprintf("successfully started the container %s on node %s", containerName, node), nil +} + +func startDatanodeContainerOnNode(nodeUser, node, containerName, dataDir, diskMap string) (string, error) { + cmd := exec.Command("ssh", nodeUser+"@"+node, + "docker run -d --name "+containerName+ + " -v "+dataDir+"/disk/"+containerName+"/data:/cfs/data"+ + " -v "+dataDir+"/bin"+":/cfs/bin:ro"+ + " -v "+dataDir+"/disk/"+containerName+"/log:/cfs/log"+ + " -v "+dataDir+"/conf/"+containerName+".json:/cfs/conf/datanode.json"+diskMap+ + " -v "+dataDir+"/script/start_datanode.sh:/cfs/script/start.sh"+ + " --restart on-failure --privileged --network host "+ImageName+ + " /bin/sh /cfs/script/start.sh ") + _, err := cmd.Output() + if err != nil { + return fmt.Sprintf("failed start %s on node %s", containerName, node), err + } + return fmt.Sprintf("successfully started the container %s on node %s", containerName, node), nil +} + +func checkContainerExistence(nodeUser, node, containerName string) (bool, error) { + cmd := exec.Command("ssh", nodeUser+"@"+node, "docker ps -a --format "+`{{.Names}}`+" | grep "+`"`+containerName+`"`) + output, err := cmd.Output() + if err != nil { + return false, err + } + for _, name := range strings.Fields(string(output)) { + if name == containerName { + return true, nil + } + } + return false, nil +} + +func checkAndDeleteContainerOnNode(nodeUser, node, containerName string) error { + if ok, _ := checkContainerExistence(nodeUser, node, containerName); ok { + log.Printf("container %s already exists on node %s", containerName, node) + _, err := stopContainerOnNode(nodeUser, node, containerName) + log.Printf("stop container %s on node %s successfully", containerName, node) + if err != nil { + return err + } + _, err = rmContainerOnNode(nodeUser, node, containerName) + log.Printf("rm container %s on node %s successfully", containerName, node) + if err != nil { + return err + } + } + return nil +} + +// Check if Docker is installed and installed +func checkAndInstallDocker(nodeUser, node string) error { + // Check if Docker is installed + cmd := exec.Command("ssh", nodeUser+"@"+node, "docker --version") + output, err := cmd.Output() + if err == nil && strings.Contains(string(output), "Docker version") { + //log.Println("Docker installed") + return nil + } + + // Docker not installed, installing Docker + cmd = exec.Command("ssh", nodeUser+"@"+node, "yum", "install", "docker", "-y") + + // Set output to standard output and standard error output + cmd.Stdout = os.Stdout + cmd.Stderr = os.Stderr + + // Execute Command + err = cmd.Run() + if err != nil { + return fmt.Errorf("failed to install Docker on node %s", node) + } + + return nil + +} + +// Check if the Docker service is started and started +func checkAndStartDockerService(nodeUser, node string) error { + // Check Docker Service Status + cmd := exec.Command("ssh", nodeUser+"@"+node, "systemctl is-active docker.service") + output, err := cmd.Output() + + if err == nil && strings.TrimSpace(string(output)) == "active" { + return nil + } + + // Docker service not started, starting Docker service + cmd = exec.Command("ssh", nodeUser+"@"+node, "systemctl start docker") + + // Set output to standard output and standard error output + cmd.Stdout = os.Stdout + cmd.Stderr = os.Stderr + + err = cmd.Run() + if err != nil { + return fmt.Errorf("failed to start Docker service on node %s", node) + } + log.Println("docker start") + + return nil +} + +// Pull image from configuration file +func pullImageOnNode(nodeUser, node, imageName string) error { + // Remote execution of commands to pull images + cmd := exec.Command("ssh", nodeUser+"@"+node, "docker pull "+imageName) + + err := cmd.Run() + if err != nil { + return fmt.Errorf("failed to pull mirror %s on node %s", imageName, node) + } + + return nil +} + +func removeImageOnNode(nodeUser, node, imageName string) error { + cmd := exec.Command("ssh", nodeUser+"@"+node, "docker rmi "+imageName) + + err := cmd.Run() + if err != nil { + return fmt.Errorf("failed to remove mirror %s on node %s", imageName, node) + } + log.Printf("success to remove mirror %s on node %s \n", imageName, node) + return nil +} + +func containerStatus(nodeUser, node, containerName string) (string, error) { + + cmd := exec.Command("ssh", nodeUser+"@"+node, "docker inspect --format='{{.State.Status}}' "+containerName) + output, err := cmd.Output() + if err != nil { + return "", fmt.Errorf("failed to pull mirror %s on node %s", containerName, node) + } + + return string(output), nil + +} diff --git a/deploy/cmd/firewall.go b/deploy/cmd/firewall.go new file mode 100644 index 000000000..780ea39be --- /dev/null +++ b/deploy/cmd/firewall.go @@ -0,0 +1,74 @@ +package cmd + +import ( + "fmt" + "log" + "os/exec" +) + +// func openRemotePortFirewall(hostname, username string, privateKeyPath string, port int) error { +// key, err := ioutil.ReadFile(privateKeyPath) +// if err != nil { +// return err +// } + +// signer, err := ssh.ParsePrivateKey(key) +// if err != nil { +// return err +// } + +// config := &ssh.ClientConfig{ +// User: username, +// Auth: []ssh.AuthMethod{ +// ssh.PublicKeys(signer), +// }, +// HostKeyCallback: ssh.InsecureIgnoreHostKey(), +// } + +// conn, err := ssh.Dial("tcp", fmt.Sprintf("%s:%d", hostname, 22), config) +// if err != nil { +// return err +// } +// defer conn.Close() + +// session, err := conn.NewSession() +// if err != nil { +// return err +// } +// defer session.Close() + +// command := fmt.Sprintf("sudo firewall-cmd --zone=public --add-port=%d/tcp --permanent", port) +// err = session.Run(command) +// if err != nil { +// return err +// } + +// return nil +// // } + +// func reloadFirewall() { + +// cmd := exec.Command("firewall-cmd", "--reload") +// cmd.Run() +// } + +func stopFirewall(nodeUser, node string) { + + cmd := exec.Command("ssh", nodeUser+"@"+node, "systemctl", "stop firewalld") + err := cmd.Run() + if err != nil { + log.Println(err) + } + log.Println(cmd) + +} + +func checkPortStatus(nodeUser, node string, port string) (string, error) { + cmd := exec.Command("ssh", nodeUser+"@"+node, "firewall-cmd --list-all | grep "+port) + fmt.Println(cmd) + _, err := cmd.Output() + if err != nil { + return fmt.Sprintf("Port %s %s is closed", node, port), err + } + return fmt.Sprintf("Port %s is open", port), nil +} diff --git a/deploy/cmd/master.go b/deploy/cmd/master.go new file mode 100644 index 000000000..8931b6eae --- /dev/null +++ b/deploy/cmd/master.go @@ -0,0 +1,173 @@ +package cmd + +import ( + "encoding/json" + "fmt" + "io/ioutil" + "log" + "strconv" + "strings" +) + +type Master struct { + ClusterName string `json:"clusterName"` + ID string `json:"id"` + Role string `json:"role"` + IP string `json:"ip"` + Listen string `json:"listen"` + Prof string `json:"prof"` + HeartbeatPort string `json:"heartbeatPort"` + ReplicaPort string `json:"replicaPort"` + Peers string `json:"peers"` + RetainLogs string `json:"retainLogs"` + ConsulAddr string `json:"consulAddr"` + ExporterPort int `json:"exporterPort"` + LogLevel string `json:"logLevel"` + LogDir string `json:"logDir"` + WALDir string `json:"walDir"` + StoreDir string `json:"storeDir"` + MetaNodeReservedMem string `json:"metaNodeReservedMem"` + EBSAddr string `json:"ebsAddr"` + EBSServicePath string `json:"ebsServicePath"` +} + +func getMasterPeers(config *Config) string { + peers := "" + for id, node := range config.DeployHostsList.Master.Hosts { + if id != len(config.DeployHostsList.Master.Hosts)-1 { + peers = peers + strconv.Itoa(id+1) + ":" + node + ":" + config.Master.Config.Listen + "," + } else { + peers = peers + strconv.Itoa(id+1) + ":" + node + ":" + config.Master.Config.Listen + } + + } + return peers +} + +func readMaster(filename string) (*Master, error) { + data, err := ioutil.ReadFile(filename) + if err != nil { + return nil, err + } + + master := &Master{} + err = json.Unmarshal(data, master) + if err != nil { + return nil, err + } + + return master, nil +} + +func writeMaster(clusterName, id, ip, listen, prof, peers string) error { + master := Master{ + ClusterName: clusterName, + ID: id, + Role: "master", + IP: ip, + Listen: listen, + Prof: prof, + HeartbeatPort: "5901", + ReplicaPort: "5902", + Peers: peers, + RetainLogs: "20000", + ConsulAddr: "http://192.168.0.101:8500", + ExporterPort: 9500, + LogLevel: "debug", + LogDir: "/cfs/log", + WALDir: "/cfs/data/wal", + StoreDir: "/cfs/data/store", + MetaNodeReservedMem: "67108864", + EBSAddr: "10.177.40.215:8500", + EBSServicePath: "access", + } + + masterData, err := json.MarshalIndent(master, "", " ") + if err != nil { + return fmt.Errorf("cannot be resolved to master.json %v", err) + } + fileName := ConfDir + "/master" + id + ".json" + err = ioutil.WriteFile(fileName, masterData, 0644) + if err != nil { + + return fmt.Errorf("unable to write %s %v", fileName, err) + } + return nil +} + +func startAllMaster() error { + + config, err := readConfig() + if err != nil { + return err + } + + files, err := ioutil.ReadDir(ConfDir) + if err != nil { + return err + } + + for _, file := range files { + if strings.HasPrefix(file.Name(), "master") && !file.IsDir() { + data, err := readMaster(ConfDir + "/" + file.Name()) + if err != nil { + fmt.Printf("Error reading file %s: %s\n", file.Name(), err) + return nil + } + confFilePath := ConfDir + "/" + file.Name() + var dataDir string + if config.Master.Config.DataDir == "" { + dataDir = config.Global.DataDir + } else { + dataDir = config.Master.Config.DataDir + + } + err = transferConfigFileToRemote(confFilePath, dataDir+"/conf", RemoteUser, data.IP) + if err != nil { + return err + } + err = checkAndDeleteContainerOnNode(RemoteUser, data.IP, data.Role+data.ID) + if err != nil { + return err + } + status, err := startMasterContainerOnNode(RemoteUser, data.IP, data.Role+data.ID, dataDir) + if err != nil { + return err + } + log.Println(status) + } + } + //Detect successful deployment + log.Println("start all master services") + return nil +} + +func stopAllMaster() error { + + files, err := ioutil.ReadDir(ConfDir) + if err != nil { + return err + } + + for _, file := range files { + if strings.HasPrefix(file.Name(), "master") && !file.IsDir() { + data, err := readMaster(ConfDir + "/" + file.Name()) + if err != nil { + fmt.Printf("Error reading file %s: %s\n", file.Name(), err) + return nil + } + status, err := stopContainerOnNode(RemoteUser, data.IP, data.Role+data.ID) + if err != nil { + return err + } + log.Println(status) + status, err = rmContainerOnNode(RemoteUser, data.IP, data.Role+data.ID) + if err != nil { + return err + } + log.Println(status) + } + + } + return nil +} diff --git a/deploy/cmd/metanode.go b/deploy/cmd/metanode.go new file mode 100644 index 000000000..b5ce401a5 --- /dev/null +++ b/deploy/cmd/metanode.go @@ -0,0 +1,243 @@ +package cmd + +import ( + "encoding/json" + "fmt" + "io/ioutil" + "log" + "strings" +) + +type MetaNode struct { + Role string `json:"role"` + Listen string `json:"listen"` + Prof string `json:"prof"` + RaftHeartbeatPort string `json:"raftHeartbeatPort"` + RaftReplicaPort string `json:"raftReplicaPort"` + LocalIP string `json:"localIP"` + ConsulAddr string `json:"consulAddr"` + ExporterPort int `json:"exporterPort"` + LogLevel string `json:"logLevel"` + LogDir string `json:"logDir"` + WarnLogDir string `json:"warnLogDir"` + TotalMem string `json:"totalMem"` + MetadataDir string `json:"metadataDir"` + RaftDir string `json:"raftDir"` + MasterAddr []string `json:"masterAddr"` +} + +func readMetaNode(filename string) (*MetaNode, error) { + data, err := ioutil.ReadFile(filename) + if err != nil { + return nil, err + } + + metaNode := &MetaNode{} + err = json.Unmarshal(data, metaNode) + if err != nil { + return nil, err + } + return metaNode, nil +} + +func writeMetaNode(listen, prof, id, localIP string, masterAddrs []string) error { + metanode := MetaNode{ + Role: "metanode", + Listen: listen, + Prof: prof, + RaftHeartbeatPort: "17230", + RaftReplicaPort: "17240", + LocalIP: localIP, + ConsulAddr: "http://192.168.0.101:8500", + ExporterPort: 9500, + LogLevel: "debug", + LogDir: "/cfs/log", + WarnLogDir: "/cfs/log", + TotalMem: "536870912", + MetadataDir: "/cfs/data/meta", + RaftDir: "/cfs/data/raft", + MasterAddr: masterAddrs, + } + + metaNodeData, err := json.MarshalIndent(metanode, "", " ") + if err != nil { + return err + } + err = ioutil.WriteFile(ConfDir+"/metanode"+id+".json", metaNodeData, 0644) + if err != nil { + return err + } + return nil + +} + +func stopMetanodeInSpecificNode(node string) error { + + files, err := ioutil.ReadDir(ConfDir) + if err != nil { + return err + } + + for _, file := range files { + if strings.HasPrefix(file.Name(), "metanode") && !file.IsDir() { + data, err := readMetaNode(ConfDir + "/" + file.Name()) + if err != nil { + fmt.Printf("Error reading file %s: %s\n", file.Name(), err) + return nil + } + if data.LocalIP == node { + status, err := stopContainerOnNode(RemoteUser, node, strings.Split(file.Name(), ".")[0]) + if err != nil { + return err + } + log.Println(status) + status, err = rmContainerOnNode(RemoteUser, node, strings.Split(file.Name(), ".")[0]) + if err != nil { + return err + } + log.Println(status) + } + } + } + + return nil + +} + +func startMetanodeInSpecificNode(node string) error { + config, err := readConfig() + if err != nil { + return err + } + files, err := ioutil.ReadDir(ConfDir) + if err != nil { + return err + } + + for _, file := range files { + if strings.HasPrefix(file.Name(), "metanode") && !file.IsDir() { + data, err := readMetaNode(ConfDir + "/" + file.Name()) + if err != nil { + fmt.Printf("Error reading file %s: %s\n", file.Name(), err) + return nil + } + if data.LocalIP == node { + var dataDir string + if config.MetaNode.Config.DataDir == "" { + dataDir = config.Global.DataDir + } else { + dataDir = config.MetaNode.Config.DataDir + + } + confFilePath := ConfDir + "/" + file.Name() + err = transferConfigFileToRemote(confFilePath, dataDir+"/conf", RemoteUser, node) + if err != nil { + return err + } + + err = checkAndDeleteContainerOnNode(RemoteUser, node, strings.Split(file.Name(), ".")[0]) + if err != nil { + return err + } + status, err := startMetanodeContainerOnNode(RemoteUser, node, strings.Split(file.Name(), ".")[0], dataDir) + if err != nil { + return err + } + log.Println(status) + break + } + + } + } + + return nil +} + +func getMasterAddrAndPort() ([]string, error) { + config, err := readConfig() + if err != nil { + return []string{}, err + } + masterAddr := make([]string, len(config.DeployHostsList.Master.Hosts)) + for id, node := range config.DeployHostsList.Master.Hosts { + masterAddr[id] = node + ":" + config.Master.Config.Listen + } + return masterAddr, nil +} + +func startAllMetaNode() error { + + config, err := readConfig() + if err != nil { + return err + } + + files, err := ioutil.ReadDir(ConfDir) + if err != nil { + return err + } + + for _, file := range files { + if strings.HasPrefix(file.Name(), "metanode") && !file.IsDir() { + data, err := readMetaNode(ConfDir + "/" + file.Name()) + if err != nil { + fmt.Printf("Error reading file %s: %s\n", file.Name(), err) + return nil + } + confFilePath := ConfDir + "/" + file.Name() + var dataDir string + if config.MetaNode.Config.DataDir == "" { + dataDir = config.Global.DataDir + } else { + dataDir = config.MetaNode.Config.DataDir + + } + err = transferConfigFileToRemote(confFilePath, dataDir+"/conf", RemoteUser, data.LocalIP) + if err != nil { + return err + } + err = checkAndDeleteContainerOnNode(RemoteUser, data.LocalIP, strings.Split(file.Name(), ".")[0]) + if err != nil { + return err + } + status, err := startMetanodeContainerOnNode(RemoteUser, data.LocalIP, strings.Split(file.Name(), ".")[0], dataDir) + if err != nil { + return err + } + log.Println(status) + } + } + + //Detect successful deployment + log.Println("start all metanode services") + return nil +} + +func stopAllMetaNode() error { + + files, err := ioutil.ReadDir(ConfDir) + if err != nil { + return err + } + + for _, file := range files { + if strings.HasPrefix(file.Name(), "metanode") && !file.IsDir() { + data, err := readMetaNode(ConfDir + "/" + file.Name()) + if err != nil { + fmt.Printf("Error reading file %s: %s\n", file.Name(), err) + return nil + } + status, err := stopContainerOnNode(RemoteUser, data.LocalIP, strings.Split(file.Name(), ".")[0]) + if err != nil { + return err + } + log.Println(status) + status, err = rmContainerOnNode(RemoteUser, data.LocalIP, strings.Split(file.Name(), ".")[0]) + if err != nil { + return err + } + log.Println(status) + } + } + return nil +} diff --git a/deploy/cmd/restart.go b/deploy/cmd/restart.go new file mode 100644 index 000000000..ec6df2f4a --- /dev/null +++ b/deploy/cmd/restart.go @@ -0,0 +1,136 @@ +package cmd + +import ( + "fmt" + "log" + + "github.com/spf13/cobra" +) + +var allRestart bool + +var RestartCmd = &cobra.Command{ + Use: "restart", + Short: "start service", + Long: `This command will start service.`, + Run: func(cmd *cobra.Command, args []string) { + if allRestart { + fmt.Println("restart all services......") + err := stopAllMaster() + if err != nil { + log.Println(err) + } + err = startAllMaster() + if err != nil { + log.Println(err) + } + err = startAllMetaNode() + if err != nil { + log.Println(err) + } + err = stopAllMetaNode() + if err != nil { + log.Println(err) + } + + err = startAllDataNode() + if err != nil { + log.Println(err) + } + err = stopAllDataNode() + if err != nil { + log.Println(err) + } + + } else { + fmt.Println(cmd.UsageString()) + } + }, +} + +var restartMasterCommand = &cobra.Command{ + Use: "master", + Short: "", + Long: "", + Run: func(cmd *cobra.Command, args []string) { + + err := stopAllMaster() + if err != nil { + log.Println(err) + } + err = startAllMaster() + if err != nil { + log.Println(err) + } + log.Println("restart all master") + }, +} + +var restartMetanodeCommand = &cobra.Command{ + Use: "metanode", + Short: "", + Long: "", + Run: func(cmd *cobra.Command, args []string) { + if cmd.Flags().Changed("ip") { + err := stopMetanodeInSpecificNode(ip) + if err != nil { + log.Println(err) + } + err = startMetanodeInSpecificNode(ip) + if err != nil { + log.Println(err) + } + log.Println("restart metanode in ", ip) + } else { + err := stopAllMetaNode() + if err != nil { + log.Println(err) + } + err = startAllMetaNode() + if err != nil { + log.Println(err) + } + log.Println("restart all metanode") + } + }, +} + +var restartDatanodeCommand = &cobra.Command{ + Use: "datanode", + Short: "", + Long: "", + Run: func(cmd *cobra.Command, args []string) { + if cmd.Flags().Changed("ip") { + err := stopDatanodeInSpecificNode(ip) + if err != nil { + log.Println(err) + } + err = startDatanodeInSpecificNode(ip) + if err != nil { + log.Println(err) + } + log.Println("restart datanode in ", ip) + + } else { + err := stopAllDataNode() + if err != nil { + log.Println(err) + } + err = startAllDataNode() + if err != nil { + log.Println(err) + } + log.Println("restart all datanode") + } + }, +} + +func init() { + RestartCmd.AddCommand(restartMasterCommand) + RestartCmd.AddCommand(restartMetanodeCommand) + RestartCmd.AddCommand(restartDatanodeCommand) + RestartCmd.Flags().BoolVarP(&allRestart, "all", "a", false, "restart all services") + RestartCmd.PersistentFlags().StringVarP(&ip, "ip", "", "", "specify an IP address to start services") + restartDatanodeCommand.Flags().StringVarP(&datanodeDisk, "disk", "d", "", "specify the disk where datanode mount") + +} diff --git a/deploy/cmd/root.go b/deploy/cmd/root.go new file mode 100644 index 000000000..65849839f --- /dev/null +++ b/deploy/cmd/root.go @@ -0,0 +1,76 @@ +package cmd + +import ( + "fmt" + "os" + + "github.com/spf13/cobra" +) + +var Version bool +var CubeFSPath string +var ConfDir string +var ScriptDir string +var BinDir string +var ConfigFileName string +var ImageName string + +const envVar = "CUBEFS" +const ClusterName = "cubeFS" +const RemoteUser = "root" + +const BinVersion = "release-3.2.1" + +func init() { + CubeFSPath = os.Getenv(envVar) + ConfDir = CubeFSPath + "/deploy/conf" + ScriptDir = CubeFSPath + "/docker/script" + BinDir = CubeFSPath + "/build/bin" + ConfigFileName = CubeFSPath + "/deploy/config.yaml" + config, _ := readConfig() + ImageName = config.Global.ContainerImage +} + +type ServerType string + +const ( + MasterServer ServerType = "master" + MetaNodeServer ServerType = "metanode" + DataNodeServer ServerType = "datanode" +) + +const ( + MasterName = "master" + MetaNodeName = "metanode" + DataNodeName = "datanode" +) + +type Status string + +const ( + Running Status = "running" + Stopped Status = "stopped" + Created Status = "created" + Paused Status = "paused" +) + +type Service struct { + ServerType ServerType + ContainerName string + NodeIP string + Status Status +} + +var RootCmd = &cobra.Command{ + Use: "deploy-cli", + Short: "CLI for managing CubeFS server and client using Docker", + Long: `cubefs is a CLI application for managing CubeFS, an open-source distributed file system, using Docker containers.`, + Run: func(cmd *cobra.Command, args []string) { + + if Version { + fmt.Printf("deploy-cli version 0.0.1 cubefs version %s \n", BinVersion) + } else { + fmt.Println(cmd.UsageString()) + } + }, +} diff --git a/deploy/cmd/ssh.go b/deploy/cmd/ssh.go new file mode 100644 index 000000000..89374efa9 --- /dev/null +++ b/deploy/cmd/ssh.go @@ -0,0 +1,53 @@ +package cmd + +import ( + "fmt" + "log" + "os" + "os/exec" +) + +// Check if the private and public key files already exist and generate an SSH key pair +func generateSSHKey() error { + // Check if the private and public key files already exist + privateKeyPath := os.Getenv("HOME") + "/.ssh/id_rsa" + //publicKeyPath := privateKeyPath + ".pub" + if _, err := os.Stat(privateKeyPath); err == nil { + return fmt.Errorf("SSH key already exists") + } + + // Generate SSH key pairs + cmd := exec.Command("ssh-keygen", "-t", "rsa", "-N", "", "-f", privateKeyPath) + err := cmd.Run() + if err != nil { + return fmt.Errorf("failed to generate SSH key: %v", err) + } + + log.Printf("SSH key generated successfully.\n") + + return nil +} + +// Establishing a secure connection +func establishSSHConnectionWithoutPassword(sourceNode, targetNodeUser, targetNode string) error { + + // Check if the private and public key files already exist and generate an SSH key pairs + generateSSHKey() + + // Check if it is possible to connect to the target node without a password + cmd := exec.Command("ssh", "-o", "BatchMode=yes", "-o", "ConnectTimeout=5", targetNodeUser+"@"+targetNode, "echo", "connection successful") + err := cmd.Run() + if err != nil { + // Copy public key to remote host + privateKeyPath := os.Getenv("HOME") + "/.ssh/id_rsa" + publicKeyPath := privateKeyPath + ".pub" + cmd = exec.Command("ssh-copy-id", "-i", publicKeyPath, targetNodeUser+"@"+targetNode) + err = cmd.Run() + if err != nil { + return fmt.Errorf("failed to establish passwordless SSH connection with %s@%s: %v", targetNodeUser, targetNode, err) + } + } + log.Printf("Passwordless SSH connection is established with %s@%s.\n", targetNodeUser, targetNode) + + return nil +} diff --git a/deploy/cmd/start.go b/deploy/cmd/start.go new file mode 100644 index 000000000..560de460c --- /dev/null +++ b/deploy/cmd/start.go @@ -0,0 +1,116 @@ +package cmd + +import ( + "fmt" + "log" + + "github.com/spf13/cobra" +) + +var ip string +var allStart bool +var datanodeDisk string +var disk string + +var StartCmd = &cobra.Command{ + Use: "start", + Short: "start service", + Long: `This command will start services.`, + Run: func(cmd *cobra.Command, args []string) { + if allStart { + fmt.Println("start all services......") + err := startAllMaster() + if err != nil { + log.Println(err) + } + + err = startAllMetaNode() + if err != nil { + log.Println(err) + } + + err = startAllDataNode() + if err != nil { + log.Println(err) + } + + } else { + fmt.Println(cmd.UsageString()) + } + + }, +} + +var startMasterCommand = &cobra.Command{ + Use: "master", + Short: "start master", + Long: "", + Run: func(cmd *cobra.Command, args []string) { + err := startAllMaster() + if err != nil { + log.Println(err) + } + }, +} + +var startFromDockerCompose = &cobra.Command{ + Use: "test", + Short: "start test for on node", + Long: "", + Run: func(cmd *cobra.Command, args []string) { + err := startALLFromDockerCompose(disk) + if err != nil { + log.Println(err) + } + }, +} + +var startMetanodeCommand = &cobra.Command{ + Use: "metanode", + Short: "start metanode", + Long: "", + Run: func(cmd *cobra.Command, args []string) { + if cmd.Flags().Changed("ip") { + err := startMetanodeInSpecificNode(ip) + if err != nil { + log.Println(err) + } + } else { + err := startAllMetaNode() + if err != nil { + log.Println(err) + } + } + + }, +} + +var startDatanodeCommand = &cobra.Command{ + Use: "datanode", + Short: "start datanode", + Long: "", + Run: func(cmd *cobra.Command, args []string) { + if cmd.Flags().Changed("ip") { + startDatanodeInSpecificNode(ip) + fmt.Println("start datanode in ", ip) + + } else { + err := startAllDataNode() + if err != nil { + log.Println(err) + } + } + }, +} + +func init() { + StartCmd.AddCommand(startMasterCommand) + StartCmd.AddCommand(startMetanodeCommand) + StartCmd.AddCommand(startDatanodeCommand) + StartCmd.AddCommand(startFromDockerCompose) + startFromDockerCompose.PersistentFlags().StringVarP(&disk, "disk", "d", "", "disk option description") + startFromDockerCompose.MarkPersistentFlagRequired("disk") + StartCmd.Flags().BoolVarP(&allStart, "all", "a", false, "start all services") + StartCmd.PersistentFlags().StringVarP(&ip, "ip", "", "", "specify an IP address to start services") + startDatanodeCommand.Flags().StringVarP(&datanodeDisk, "disk", "d", "", "specify the disk where datanode mount") +} diff --git a/deploy/cmd/start_docker_compose.go b/deploy/cmd/start_docker_compose.go new file mode 100644 index 000000000..4af5ac9ef --- /dev/null +++ b/deploy/cmd/start_docker_compose.go @@ -0,0 +1,54 @@ +package cmd + +import ( + "fmt" + "log" + "os" + "os/exec" +) + +func startALLFromDockerCompose(disk string) error { + scriptPath := CubeFSPath + "/docker/run_docker.sh" + args := []string{"-r", "-d", disk} + err := os.Chmod(scriptPath, 0700) + if err != nil { + log.Fatal(err) + } + err = runScript(scriptPath, args...) + + if err != nil { + log.Fatal(err) + } + return nil +} + +func stopALLFromDockerCompose() error { + + err := os.Chdir(CubeFSPath + "/docker") + if err != nil { + return err + } + + cmd := exec.Command("docker-compose", "down") + cmd.Stdout = os.Stdout + cmd.Stderr = os.Stderr + err = cmd.Run() + if err != nil { + return err + } + + return nil +} + +func runScript(scriptPath string, args ...string) error { + cmd := exec.Command(scriptPath, args...) + cmd.Stdout = os.Stdout + cmd.Stderr = os.Stderr + + err := cmd.Run() + if err != nil { + return fmt.Errorf("failed to run script: %w", err) + } + + return nil +} diff --git a/deploy/cmd/stop.go b/deploy/cmd/stop.go new file mode 100644 index 000000000..ab704f19c --- /dev/null +++ b/deploy/cmd/stop.go @@ -0,0 +1,118 @@ +package cmd + +import ( + "fmt" + "log" + + "github.com/spf13/cobra" +) + +var allStop bool + +var StopCmd = &cobra.Command{ + Use: "stop", + Short: "start service", + Long: `This command will start service.`, + Run: func(cmd *cobra.Command, args []string) { + if allStop { + fmt.Println("stop all services......") + err := stopAllMaster() + if err != nil { + log.Println(err) + } + fmt.Println("stop all master services") + + err = stopAllMetaNode() + if err != nil { + log.Println(err) + } + fmt.Println("stop all metanode services") + err = stopAllDataNode() + if err != nil { + log.Println(err) + } + fmt.Println("stop all datanode services") + + } else { + fmt.Println(cmd.UsageString()) + } + }, +} +var stopFromDockerCompose = &cobra.Command{ + Use: "test", + Short: "start test for on node", + Long: "", + Run: func(cmd *cobra.Command, args []string) { + err := stopALLFromDockerCompose() + if err != nil { + log.Println(err) + } + }, +} + +var stopMasterCommand = &cobra.Command{ + Use: "master", + Short: "", + Long: "", + Run: func(cmd *cobra.Command, args []string) { + if cmd.Flags().Changed("ip") { + fmt.Println("ip:", ip) + } else { + err := stopAllMaster() + if err != nil { + log.Println(err) + } + log.Println("stop all master services") + } + }, +} + +var stopMetanodeCommand = &cobra.Command{ + Use: "metanode", + Short: "", + Long: "", + Run: func(cmd *cobra.Command, args []string) { + if cmd.Flags().Changed("ip") { + err := stopMetanodeInSpecificNode(ip) + if err != nil { + log.Println(err) + } + log.Println("stop metanode in ", ip) + } else { + err := stopAllMetaNode() + if err != nil { + log.Println(err) + } + log.Println("stop all metanode services") + } + }, +} + +var stopDatanodeCommand = &cobra.Command{ + Use: "datanode", + Short: "", + Long: "", + Run: func(cmd *cobra.Command, args []string) { + if cmd.Flags().Changed("ip") { + stopDatanodeInSpecificNode(ip) + fmt.Println("stop datanode in ", ip) + } else { + err := stopAllDataNode() + if err != nil { + log.Println(err) + } + fmt.Println("stop all datanode services") + } + }, +} + +func init() { + StopCmd.AddCommand(stopMasterCommand) + StopCmd.AddCommand(stopMetanodeCommand) + StopCmd.AddCommand(stopDatanodeCommand) + StopCmd.AddCommand(stopFromDockerCompose) + StopCmd.Flags().BoolVarP(&allStop, "all", "a", false, "stop all services") + StopCmd.PersistentFlags().StringVarP(&ip, "ip", "", "", "specify an IP address to start services") + stopDatanodeCommand.Flags().StringVarP(&datanodeDisk, "disk", "d", "", "specify the disk where datanode mount") + +} diff --git a/deploy/cmd/transform.go b/deploy/cmd/transform.go new file mode 100644 index 000000000..b4782f4e4 --- /dev/null +++ b/deploy/cmd/transform.go @@ -0,0 +1,121 @@ +package cmd + +import ( + "fmt" + "io" + "log" + "os" + "os/exec" + "path/filepath" +) + +func transferConfigFileToRemote(localFilePath string, remoteFilePath string, remoteUser string, remoteHost string) error { + localFile, err := os.Open(localFilePath) + if err != nil { + return fmt.Errorf("failed to open local file: %s", err) + } + defer localFile.Close() + + cmd := exec.Command("scp", localFilePath, remoteUser+"@"+remoteHost+":"+remoteFilePath) + err = cmd.Run() + if err != nil { + return fmt.Errorf("file %s transferred to %s@%s:%s failed", localFilePath, remoteUser, remoteHost, remoteFilePath) + } + + log.Printf("file '%s' transferred to '%s@%s:%s' successfully.\n", localFilePath, remoteUser, remoteHost, remoteFilePath) + + return nil +} + +func transferDirectoryToRemote(localFilePath string, remoteFilePath string, remoteUser string, remoteHost string) error { + + localFile, err := os.Open(localFilePath) + if err != nil { + return fmt.Errorf("failed to open local file: %s", err) + } + defer localFile.Close() + + cmd := exec.Command("scp", "-r", localFilePath, remoteUser+"@"+remoteHost+":"+remoteFilePath) + + err = cmd.Run() + if err != nil { + return fmt.Errorf("file %s transferred to %s@%s:%s failed", localFilePath, remoteUser, remoteHost, remoteFilePath) + } + + log.Printf("file '%s' transferred to '%s@%s:%s' successfully.\n", localFilePath, remoteUser, remoteHost, remoteFilePath) + + return nil +} + +func copyFolder(sourcePath, destinationPath string) error { + err := filepath.Walk(sourcePath, func(path string, info os.FileInfo, err error) error { + if err != nil { + return err + } + + relativePath, err := filepath.Rel(sourcePath, path) + if err != nil { + return err + } + + destinationFilePath := filepath.Join(destinationPath, relativePath) + + if info.IsDir() { + err := os.MkdirAll(destinationFilePath, info.Mode()) + if err != nil { + return err + } + } else { + sourceFile, err := os.Open(path) + if err != nil { + return err + } + defer sourceFile.Close() + + destinationFile, err := os.Create(destinationFilePath) + if err != nil { + return err + } + defer destinationFile.Close() + + _, err = io.Copy(destinationFile, sourceFile) + if err != nil { + return err + } + } + + return nil + }) + + if err != nil { + return err + } + + return nil +} + +func createRemoteFolder(nodeUser, node, folderPath string) error { + cmdStr := "mkdir -p " + folderPath + cmd := exec.Command("ssh", nodeUser+"@"+node, cmdStr) + _, err := cmd.Output() + if err != nil { + return err + } + return nil + +} +func createFolder(folderPath string) error { + _, err := os.Stat(folderPath) + if os.IsNotExist(err) { + err := os.MkdirAll(folderPath, 0755) + if err != nil { + return err + } + + } else if err != nil { + return err + } else { + return nil + } + return nil +} diff --git a/deploy/config.yaml b/deploy/config.yaml new file mode 100644 index 000000000..84e6345f1 --- /dev/null +++ b/deploy/config.yaml @@ -0,0 +1,50 @@ +global: + container_image: docker.io/cubefs/cbfs-base:1.0-golang-1.17.13 + data_dir: /data + variable: + target: 0.0.1 + +master: + config: + listen: 17010 + prof: 17020 + data_dir: /data + +metanode: + config: + listen: 17210 + prof: 17220 + data_dir: /data + +datanode: + config: + listen: 17310 + prof: 17320 + data_dir: /data + + +deplopy_hosts_list: + master: + hosts: + - 10.1.0.44 + - 10.1.0.45 + - 10.1.0.46 + metanode: + hosts: + - 10.1.0.44 + - 10.1.0.45 + - 10.1.0.46 + + datanode: + - hosts: 10.1.0.44 + disk: + - path: /data/disk0 + size: 10737418240 + - hosts: 10.1.0.45 + disk: + - path: /data/disk0 + size: 10737418240 + - hosts: 10.1.0.46 + disk: + - path: /data/disk0 + size: 10737418240 diff --git a/deploy/deploy_cli.go b/deploy/deploy_cli.go new file mode 100644 index 000000000..08420b381 --- /dev/null +++ b/deploy/deploy_cli.go @@ -0,0 +1,36 @@ +package main + +import ( + "fmt" + "github.com/cubefs/cubefs/deploy/cmd" + "github.com/cubefs/cubefs/util/log" + "os" +) + +func init() { + cmd.RootCmd.Flags().BoolVarP(&cmd.Version, "version", "v", false, "show version information") + cmd.RootCmd.AddCommand(cmd.ClusterCmd) + cmd.RootCmd.AddCommand(cmd.StartCmd) + cmd.RootCmd.AddCommand(cmd.StopCmd) + cmd.RootCmd.AddCommand(cmd.RestartCmd) + +} + +func main() { + + _, err := log.InitLog("/tmp/cfs", "deploy", log.DebugLevel, nil, log.DefaultLogLeftSpaceLimit) + if err != nil { + fmt.Fprintf(os.Stderr, "Error: %v\n", err) + log.LogFlush() + os.Exit(1) + } + + err = cmd.RootCmd.Execute() + defer log.LogFlush() + if err != nil { + fmt.Fprintf(os.Stderr, "Error: %v\n", err) + log.LogFlush() + os.Exit(1) + } + +}