feat(deploy): add a cfs-deploy tool

Signed-off-by: ytyuanxi <1206812491@qq.com>
This commit is contained in:
ytyuanxi 2023-10-27 19:02:32 +08:00 committed by jiaruo bai
parent 1fba8dbee3
commit 06f46d9a52
18 changed files with 2134 additions and 2 deletions

View File

@ -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)

View File

@ -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
;;

291
deploy/cmd/cluster.go Normal file
View File

@ -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")
}

133
deploy/cmd/config.go Normal file
View File

@ -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
}

234
deploy/cmd/datanode.go Normal file
View File

@ -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
}

206
deploy/cmd/docker.go Normal file
View File

@ -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
}

74
deploy/cmd/firewall.go Normal file
View File

@ -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
}

173
deploy/cmd/master.go Normal file
View File

@ -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
}

243
deploy/cmd/metanode.go Normal file
View File

@ -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
}

136
deploy/cmd/restart.go Normal file
View File

@ -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")
}

76
deploy/cmd/root.go Normal file
View File

@ -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())
}
},
}

53
deploy/cmd/ssh.go Normal file
View File

@ -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
}

116
deploy/cmd/start.go Normal file
View File

@ -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")
}

View File

@ -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
}

118
deploy/cmd/stop.go Normal file
View File

@ -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")
}

121
deploy/cmd/transform.go Normal file
View File

@ -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
}

50
deploy/config.yaml Normal file
View File

@ -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

36
deploy/deploy_cli.go Normal file
View File

@ -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)
}
}