From 3674fe385c2f3328fc96ac045bea7ba93200a77c Mon Sep 17 00:00:00 2001 From: f1v3-dev Date: Tue, 9 Jun 2026 11:07:24 +0900 Subject: [PATCH] FEATURE: Add zk deploy command --- cmd/zk/deploy.go | 12 ++++- internal/topology/zk.go | 40 +++++++++++++- internal/zk/config.go | 94 +++++++++++++++++++++++++++++++++ internal/zk/deploy.go | 112 ++++++++++++++++++++++++++++++++++++++++ internal/zk/download.go | 45 ++++++++++++++++ internal/zk/install.go | 99 +++++++++++++++++++++++++++++++++++ zk-sample-topology.yml | 39 ++++++++++++++ 7 files changed, 438 insertions(+), 3 deletions(-) create mode 100644 internal/zk/config.go create mode 100644 internal/zk/deploy.go create mode 100644 internal/zk/download.go create mode 100644 internal/zk/install.go create mode 100644 zk-sample-topology.yml diff --git a/cmd/zk/deploy.go b/cmd/zk/deploy.go index 9d9b520..30dba27 100644 --- a/cmd/zk/deploy.go +++ b/cmd/zk/deploy.go @@ -1,12 +1,20 @@ package zk -import "github.com/spf13/cobra" +import ( + "github.com/jam2in/arcusctl/internal/zk" + "github.com/spf13/cobra" +) var deployCmd = &cobra.Command{ Use: "deploy ", Short: "Deploy a new ZooKeeper ensemble", Args: cobra.ExactArgs(2), Run: func(cmd *cobra.Command, args []string) { - // TODO: deploy 구현 + version := args[0] + topologyPath := args[1] + + if err := zk.Deploy(version, topologyPath); err != nil { + panic(err) + } }, } diff --git a/internal/topology/zk.go b/internal/topology/zk.go index 999d06e..2bbc334 100644 --- a/internal/topology/zk.go +++ b/internal/topology/zk.go @@ -1,6 +1,9 @@ package topology -import "strings" +import ( + "fmt" + "strings" +) type ZKTopology struct { Name string `yaml:"name"` @@ -33,3 +36,38 @@ func (s *ZKServer) Host() string { host, _, _, _ := s.ParseAddress() return host } + +func (topo *ZKTopology) Validate() error { + if strings.TrimSpace(topo.Name) == "" { + return fmt.Errorf("ZooKeeper ensemble name is required") + } + + if strings.TrimSpace(topo.Path) == "" { + return fmt.Errorf("ZooKeeper installation path is required") + } + + if len(topo.Servers) == 0 { + return fmt.Errorf("no servers defined in topology") + } + + seenMyID := map[int]bool{} + seenAddress := map[string]bool{} + + for _, s := range topo.Servers { + if seenMyID[s.MyID] { + return fmt.Errorf("duplicate myid: %d", s.MyID) + } + seenMyID[s.MyID] = true + + if seenAddress[s.Address] { + return fmt.Errorf("duplicate address: %s", s.Address) + } + seenAddress[s.Address] = true + + if s.Config.DataDir == "" { + return fmt.Errorf("server myid=%d: data_dir is required", s.MyID) + } + } + + return nil +} diff --git a/internal/zk/config.go b/internal/zk/config.go new file mode 100644 index 0000000..e1dccf5 --- /dev/null +++ b/internal/zk/config.go @@ -0,0 +1,94 @@ +package zk + +import ( + "fmt" + "strings" + + "github.com/jam2in/arcusctl/internal/topology" +) + +const ( + defaultTickTime = 2000 + defaultInitLimit = 10 + defaultSyncLimit = 5 +) + +func buildConfig(server topology.ZKServer, topo *topology.ZKTopology) string { + var sb strings.Builder + + cfg := server.Config + dynamicConfigPath := fmt.Sprintf("%s/conf_myid_%d/zoo.cfg.dynamic", topo.Path, server.MyID) + + fmt.Fprintf(&sb, "tickTime=%d\n", cfg.TickTime) + fmt.Fprintf(&sb, "initLimit=%d\n", cfg.InitLimit) + fmt.Fprintf(&sb, "syncLimit=%d\n", cfg.SyncLimit) + fmt.Fprintf(&sb, "dataDir=%s/zk%d\n", cfg.DataDir, server.MyID) + fmt.Fprintf(&sb, "dataLogDir=%s/zk%d\n", cfg.DataLogDir, server.MyID) + fmt.Fprintf(&sb, "dynamicConfigFile=%s\n", dynamicConfigPath) + + sb.WriteString("standaloneEnabled=false\n") + sb.WriteString("reconfigEnabled=true\n") + sb.WriteString("4lw.commands.whitelist=*\n") + + for k, v := range cfg.Properties { + fmt.Fprintf(&sb, "%s=%s\n", k, v) + } + + return sb.String() +} + +func buildDynamicConfig(topo *topology.ZKTopology) string { + var sb strings.Builder + + for _, s := range topo.Servers { + host, clientPort, quorumPort, electionPort := s.ParseAddress() + fmt.Fprintf(&sb, "server.%d=%s:%s:%s;%s\n", + s.MyID, host, quorumPort, electionPort, clientPort) + } + + return sb.String() +} + +func mergeConfig(globalConfig topology.ZKConfig, nodeConfig *topology.ZKConfig) topology.ZKConfig { + merged := globalConfig + if nodeConfig != nil { + if nodeConfig.TickTime > 0 { + merged.TickTime = nodeConfig.TickTime + } + if nodeConfig.InitLimit > 0 { + merged.InitLimit = nodeConfig.InitLimit + } + if nodeConfig.SyncLimit > 0 { + merged.SyncLimit = nodeConfig.SyncLimit + } + if nodeConfig.DataDir != "" { + merged.DataDir = nodeConfig.DataDir + } + if nodeConfig.DataLogDir != "" { + merged.DataLogDir = nodeConfig.DataLogDir + } + if nodeConfig.Properties != nil { + if merged.Properties == nil { + merged.Properties = map[string]string{} + } + for k, v := range nodeConfig.Properties { + merged.Properties[k] = v + } + } + } + + if merged.TickTime == 0 { + merged.TickTime = defaultTickTime + } + if merged.InitLimit == 0 { + merged.InitLimit = defaultInitLimit + } + if merged.SyncLimit == 0 { + merged.SyncLimit = defaultSyncLimit + } + if merged.DataLogDir == "" { + merged.DataLogDir = merged.DataDir + } + + return merged +} diff --git a/internal/zk/deploy.go b/internal/zk/deploy.go new file mode 100644 index 0000000..1ed42f2 --- /dev/null +++ b/internal/zk/deploy.go @@ -0,0 +1,112 @@ +package zk + +import ( + "bufio" + "fmt" + "os" + "strings" + "text/tabwriter" + + "github.com/jam2in/arcusctl/internal/store" + "github.com/jam2in/arcusctl/internal/topology" +) + +func Deploy(version string, topologyPath string) error { + topo, topologyBytes, err := prepareTopology(topologyPath) + if err != nil { + return err + } + + printPlan(topo, version) + if !confirm() { + fmt.Println("Aborted.") + return nil + } + + localTarPath, err := ensureDownloaded(version) + if err != nil { + return err + } + + installed, err := installServers(topo, version, localTarPath) + if err != nil { + printRecoveryGuide(installed, topo) + return err + } + + if err := store.SaveZK(topo.Name, version, topologyBytes); err != nil { + printRecoveryGuide(installed, topo) + return fmt.Errorf("save metadata: %w", err) + } + + fmt.Printf("ZooKeeper ensemble %q deployed successfully.\n", topo.Name) + return nil +} + +func prepareTopology(topologyPath string) (*topology.ZKTopology, []byte, error) { + topo, rawData, err := topology.LoadZK(topologyPath) + if err != nil { + return nil, nil, err + } + + for i := range topo.Servers { + merged := mergeConfig(topo.GlobalConfig, topo.Servers[i].Config) + topo.Servers[i].Config = &merged + } + + if err := topo.Validate(); err != nil { + return nil, nil, err + } + + if store.ZKExists(topo.Name) { + return nil, nil, fmt.Errorf("ZooKeeper ensemble %q already exists", topo.Name) + } + + return topo, rawData, nil +} + +func printPlan(topo *topology.ZKTopology, version string) { + fmt.Printf("ZooKeeper ensemble %q will be deployed (version: %s)\n\n", topo.Name, version) + + w := tabwriter.NewWriter(os.Stdout, 0, 0, 2, ' ', 0) + fmt.Fprintln(w, "ROLE\tHOST\tPORTS\tDIRECTORIES") + fmt.Fprintln(w, "\t\t\t\t") + for _, s := range topo.Servers { + host, clientPort, quorumPort, electionPort := s.ParseAddress() + ports := fmt.Sprintf("%s/%s/%s", clientPort, quorumPort, electionPort) + fmt.Fprintf(w, "zookeeper\t%s\t%s\t%s\n", host, ports, topo.Path) + } + w.Flush() + + fmt.Println() + fmt.Println("Attention:") + fmt.Println(" 1. If the topology is not what you expected, check your yaml file.") + fmt.Println(" 2. Please confirm there is no port/directory conflicts in same host.") +} + +func printRecoveryGuide(deployed []topology.ZKServer, topo *topology.ZKTopology) { + if len(deployed) == 0 { + return + } + + fmt.Println("\nDeployment failed. Manual recovery required.") + fmt.Println("The following servers have been partially installed:") + for _, s := range deployed { + fmt.Printf(" - %s (myid=%d)\n", s.Host(), s.MyID) + } + fmt.Println("\nTo clean up, manually run on each server:") + fmt.Printf(" rm -rf %s/conf_myid_\n", topo.Path) +} + +func confirm() bool { + // FIXME: internal.ReadStdin()으로 변경 필요 + fmt.Print("\nProceed? [y/N]: ") + reader := bufio.NewReader(os.Stdin) + input, err := reader.ReadString('\n') + if err != nil { + return false + } + + input = strings.TrimSpace(strings.ToLower(input)) + return input == "y" || input == "yes" +} diff --git a/internal/zk/download.go b/internal/zk/download.go new file mode 100644 index 0000000..8b6496b --- /dev/null +++ b/internal/zk/download.go @@ -0,0 +1,45 @@ +package zk + +import ( + "fmt" + "os" + "os/exec" + "path/filepath" + + "github.com/jam2in/arcusctl/internal" +) + +const zkDownloadURLTemplate = "https://archive.apache.org/dist/zookeeper/zookeeper-%s/apache-zookeeper-%s-bin.tar.gz" + +func ensureDownloaded(version string) (string, error) { + dir := imageDir() + + if err := os.MkdirAll(dir, 0755); err != nil { + return "", fmt.Errorf("create image dir: %w", err) + } + + filename := fmt.Sprintf("apache-zookeeper-%s-bin.tar.gz", version) + localPath := filepath.Join(dir, filename) + + if _, err := os.Stat(localPath); err == nil { + fmt.Printf("Using existing file: %s\n", localPath) + return localPath, nil + } + + url := fmt.Sprintf(zkDownloadURLTemplate, version, version) + fmt.Printf("Downloading %s...\n", url) + + cmd := exec.Command("wget", "-q", "-O", localPath, url) + cmd.Stdout = os.Stdout + cmd.Stderr = os.Stderr + if err := cmd.Run(); err != nil { + _ = os.Remove(localPath) + return "", fmt.Errorf("download %q failed: %w", version, err) + } + + return localPath, nil +} + +func imageDir() string { + return filepath.Join(internal.Config.Home, "images", "zookeeper") +} diff --git a/internal/zk/install.go b/internal/zk/install.go new file mode 100644 index 0000000..9c918df --- /dev/null +++ b/internal/zk/install.go @@ -0,0 +1,99 @@ +package zk + +import ( + "fmt" + "os" + "path/filepath" + + "github.com/jam2in/arcusctl/internal/ssh" + "github.com/jam2in/arcusctl/internal/topology" +) + +func installServers(topo *topology.ZKTopology, version string, localTarPath string) ([]topology.ZKServer, error) { + var installed []topology.ZKServer + archiveInstalled := map[string]bool{} + + for _, server := range topo.Servers { + host := server.Host() + fmt.Printf("[%d/%d] %s: installing...\n", len(installed)+1, len(topo.Servers), host) + + if !archiveInstalled[host] { + if err := installArchive(host, topo.Path, version, localTarPath); err != nil { + return installed, fmt.Errorf("install %s: %w", host, err) + } + archiveInstalled[host] = true + } + + if err := configureServer(server, topo); err != nil { + return installed, fmt.Errorf("install %s: %w", host, err) + } + installed = append(installed, server) + } + + return installed, nil +} + +func installArchive(host string, installPath string, version string, localTarPath string) error { + if err := ssh.Run(host, fmt.Sprintf("mkdir -p %s", installPath)); err != nil { + return fmt.Errorf("mkdir base path on %s: %w", host, err) + } + + remoteTarPath := fmt.Sprintf("%s/zookeeper-%s.tar.gz", installPath, version) + if err := ssh.Copy(localTarPath, host, remoteTarPath); err != nil { + return fmt.Errorf("copy file to %s: %w", host, err) + } + + extractCmd := fmt.Sprintf("tar -xzf %s -C %s --strip-components=1", remoteTarPath, installPath) + if err := ssh.Run(host, extractCmd); err != nil { + return fmt.Errorf("extract file on %s: %w", host, err) + } + + return nil +} + +func configureServer(server topology.ZKServer, topo *topology.ZKTopology) error { + host := server.Host() + confDir := fmt.Sprintf("%s/conf_myid_%d", topo.Path, server.MyID) + dataDir := server.Config.DataDir + dataDirPath := fmt.Sprintf("%s/zk%d", dataDir, server.MyID) + + mkdirCmd := fmt.Sprintf("mkdir -p %s %s", confDir, dataDirPath) + if err := ssh.Run(host, mkdirCmd); err != nil { + return fmt.Errorf("mkdir on %s: %w", host, err) + } + + myidCmd := fmt.Sprintf("echo %d > %s/myid", server.MyID, dataDirPath) + if err := ssh.Run(host, myidCmd); err != nil { + return fmt.Errorf("write myid on %s: %w", host, err) + } + + config := buildConfig(server, topo) + if err := uploadFile(host, config, filepath.Join(confDir, "zoo.cfg")); err != nil { + return fmt.Errorf("upload zoo.cfg to %s: %w", host, err) + } + + dynamicConfig := buildDynamicConfig(topo) + if err := uploadFile(host, dynamicConfig, filepath.Join(confDir, "zoo.cfg.dynamic")); err != nil { + return fmt.Errorf("upload zoo.cfg.dynamic to %s: %w", host, err) + } + + return nil +} + +func uploadFile(host string, content string, remotePath string) error { + tmp, err := os.CreateTemp("", "arcusctl-*") + if err != nil { + return err + } + defer os.Remove(tmp.Name()) + + if _, err := tmp.WriteString(content); err != nil { + return err + } + + if err := tmp.Close(); err != nil { + return err + } + + return ssh.Copy(tmp.Name(), host, remotePath) +} diff --git a/zk-sample-topology.yml b/zk-sample-topology.yml new file mode 100644 index 0000000..49c220a --- /dev/null +++ b/zk-sample-topology.yml @@ -0,0 +1,39 @@ +# For more configuration options, see: +# https://zookeeper.apache.org/doc/r3.5.9/zookeeperAdmin.html#sc_configuration + +name: my-ensemble # required +path: /home/arcus/zookeeper # required + +servers: + - myid: 1 # required + address: zk1:2181:2888:3888 # required (host:clientPort:quorumPort:electionPort) + config: # optional - per-node override + data_log_dir: /data/zk1-txlog + + - myid: 2 + address: zk2:2181:2888:3888 + + - myid: 3 + address: zk3:2181:2888:3888 + +global_config: + # optional (default: 2000) + tick_time: 2000 + + # optional (default: 10) - ticks for follower to sync to leader + init_limit: 10 + + # optional (default: 5) - ticks between request and ack + sync_limit: 5 + + # required + data_dir: /var/lib/zk/data + + # optional (default: data_dir) - dedicated disk recommended + data_log_dir: /var/lib/zk/datalog + + # optional - other ZK options (string values only) + properties: + maxClientCnxns: "60" + autopurge.snapRetainCount: "10" + autopurge.purgeInterval: "24" \ No newline at end of file