Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 6 additions & 1 deletion cmd/zk/delete.go
Original file line number Diff line number Diff line change
Expand Up @@ -10,8 +10,13 @@ var deleteCmd = &cobra.Command{
Short: "Delete a ZooKeeper ensemble",
Args: cobra.ExactArgs(1),
Run: func(cmd *cobra.Command, args []string) {
if err := zk.Delete(args[0]); err != nil {
purge, _ := cmd.Flags().GetBool("purge")
if err := zk.Delete(args[0], purge); err != nil {
panic(err)
}
},
}

func init() {
deleteCmd.Flags().Bool("purge", false, "also remove the installation directory on each server")
}
11 changes: 5 additions & 6 deletions internal/cluster/cluster.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,8 +33,7 @@ func loadCluster(serviceCode string) (
// installed there from same path and version.
func sharingCluster(
serviceCode string,
topoPath string,
version string,
installPath string,
host string,
) (string, error) {
registered, err := store.ListCluster()
Expand All @@ -57,8 +56,8 @@ func sharingCluster(
continue
}

if topo.Path == topoPath &&
meta.Version == version &&
otherInstallPath := memcachedInstallPath(topo.Path, meta.Version)
if installPath == otherInstallPath &&
hasServerOn(topo.Servers, host) {
return other, nil
}
Expand All @@ -68,10 +67,10 @@ func sharingCluster(
}

func hasServerOn(
topoServers []topology.CacheServer,
cacheServers []topology.CacheServer,
host string,
) bool {
for _, server := range topoServers {
for _, server := range cacheServers {
if server.Host() == host {
return true
}
Expand Down
2 changes: 1 addition & 1 deletion internal/cluster/delete.go
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ func removeInstallationDirs(
installPath := memcachedInstallPath(topo.Path, version)

for _, host := range distinctHosts(topo.Servers) {
other, err := sharingCluster(serviceCode, topo.Path, version, host)
other, err := sharingCluster(serviceCode, installPath, host)
if err != nil {
return err
}
Expand Down
2 changes: 1 addition & 1 deletion internal/store/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -190,7 +190,7 @@ func ListCluster() ([]string, error) {
}

func zkBaseDir() string {
return filepath.Join(internal.Config.Home, "clusters", "zk")
return filepath.Join(internal.Config.Home, "clusters", "zookeeper")
}

func zkDir(ensembleName string) string {
Expand Down
4 changes: 0 additions & 4 deletions internal/topology/zk.go
Original file line number Diff line number Diff line change
Expand Up @@ -63,10 +63,6 @@ func (topo *ZKTopology) Validate() error {
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
Expand Down
8 changes: 5 additions & 3 deletions internal/zk/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -17,13 +17,15 @@ 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)
dataDir := zkNodeDataPath(cfg.DataDir, server.MyID)
dataLogDir := zkNodeDataPath(cfg.DataLogDir, server.MyID)
dynamicConfigPath := zkDynamicConfigPath(topo.Path, topo.Name, 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, "dataDir=%s\n", dataDir)
fmt.Fprintf(&sb, "dataLogDir=%s\n", dataLogDir)
fmt.Fprintf(&sb, "dynamicConfigFile=%s\n", dynamicConfigPath)

sb.WriteString("standaloneEnabled=false\n")
Expand Down
84 changes: 72 additions & 12 deletions internal/zk/delete.go
Original file line number Diff line number Diff line change
Expand Up @@ -11,17 +11,19 @@ import (
"github.com/jam2in/arcusctl/internal/topology"
)

func Delete(ensembleName string) error {
_, topo, err := loadEnsemble(ensembleName)
const removeCommandTemplate = "rm -rf %s"

func Delete(ensembleName string, purge bool) error {
meta, topo, err := loadEnsemble(ensembleName)
if err != nil {
return err
}

if err := verifyTopology(topo.Servers, topo.Path); err != nil {
if err := verifyTopology(topo.Servers, topo.Path, topo.Name); err != nil {
return err
}

if err := verifyAllStopped(topo.Servers, topo.Path); err != nil {
if err := verifyAllStopped(topo.Servers, topo.Path, topo.Name); err != nil {
return err
}

Expand All @@ -34,11 +36,19 @@ func Delete(ensembleName string) error {
hostsMap := groupServersByHost(topo.Servers)
for host, servers := range hostsMap {
fmt.Printf("Removing files on %s...\n", host)
if err := removeHostFiles(host, servers, topo.Path); err != nil {
if err := removeHostFiles(host, servers, topo.Path, topo.Name); err != nil {
return fmt.Errorf("remove files on %s: %w", host, err)
}
}

// If user specified --purge, remove the installation directories on each server.
// If the installation directory is shared with another ensemble, it will not be removed.
if purge {
if err := removeInstallationDirs(ensembleName, topo, meta.Version); err != nil {
return err
}
}

if err := store.DeleteZK(ensembleName); err != nil {
return fmt.Errorf("delete metadata: %w", err)
}
Expand All @@ -47,20 +57,29 @@ func Delete(ensembleName string) error {
return nil
}

func verifyTopology(servers []topology.ZKServer, topoPath string) error {
func verifyTopology(
servers []topology.ZKServer,
topoPath string,
ensembleName string,
) error {
for _, server := range servers {
confDir := fmt.Sprintf("%s/conf_myid_%d", topoPath, server.MyID)
confDir := zkConfigDir(topoPath, ensembleName, server.MyID)

if err := ssh.Run(server.Host(), fmt.Sprintf("test -d %s", confDir)); err != nil {
return fmt.Errorf("topology mismatch: %s not found on %s", confDir, server.Host())
}
}
return nil
}

func verifyAllStopped(servers []topology.ZKServer, topoPath string) error {
func verifyAllStopped(
servers []topology.ZKServer,
topoPath string,
ensembleName string,
) error {
for _, server := range servers {
confPath := zkConfigPath(topoPath, server.MyID)
cmd := fmt.Sprintf("pgrep -f '[Q]uorumPeerMain.*%s' > /dev/null 2>&1", confPath)
confDir := zkConfigDir(topoPath, ensembleName, server.MyID)
cmd := fmt.Sprintf("pgrep -f '[Q]uorumPeerMain.*%s' > /dev/null 2>&1", confDir)
if err := ssh.Run(server.Host(), cmd); err == nil {
return fmt.Errorf("server %s (myid=%d) is still running. stop the ensemble before delete",
server.Host(), server.MyID)
Expand All @@ -77,15 +96,56 @@ func groupServersByHost(servers []topology.ZKServer) map[string][]topology.ZKSer
return hosts
}

func removeHostFiles(host string, servers []topology.ZKServer, topoPath string) error {
removePaths := []string{topoPath}
func removeHostFiles(
host string,
servers []topology.ZKServer,
topoPath string,
ensembleName string,
) error {
var removePaths []string

for _, server := range servers {
// Remove conf/<ensembleName>/zk<myid>.cfg and data/log directories
removePaths = append(
removePaths,
zkConfigDir(topoPath, ensembleName, server.MyID),
)
removePaths = append(removePaths, nodeDataPaths(server)...)
}

return ssh.Run(host, "rm -rf "+strings.Join(removePaths, " "))
}

func removeInstallationDirs(
ensembleName string,
topo *topology.ZKTopology,
version string,
) error {
installPath := zkInstallPath(topo.Path, version)

for host := range groupServersByHost(topo.Servers) {
other, err := sharingEnsemble(ensembleName, installPath, host)
if err != nil {
return err
}

if other != "" {
fmt.Printf(
"Skip removing directory on %s: install path is shared with ensemble %q\n",
host, other,
)
continue
}

fmt.Printf("Removing installation directory on %s...\n", host)
if err := ssh.Run(host, fmt.Sprintf(removeCommandTemplate, installPath)); err != nil {
return fmt.Errorf("remove installation on %s: %w", host, err)
}
}

return nil
}

func nodeDataPaths(server topology.ZKServer) []string {
nodeDirName := fmt.Sprintf("zk%d", server.MyID)

Expand Down
14 changes: 11 additions & 3 deletions internal/zk/deploy.go
Original file line number Diff line number Diff line change
Expand Up @@ -63,14 +63,15 @@ func prepareTopology(topoPath string) (*topology.ZKTopology, []byte, error) {

func printPlan(topo *topology.ZKTopology, version string) {
fmt.Printf("ZooKeeper ensemble %q will be deployed (version: %s)\n\n", topo.Name, version)
installPath := zkInstallPath(topo.Path, version)

w := tabwriter.NewWriter(os.Stdout, 0, 0, 2, ' ', 0)
fmt.Fprintln(w, "ROLE\tHOST\tPORTS\tDIRECTORIES")
fmt.Fprintln(w, "MYID\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)
fmt.Fprintf(w, "%d\t%s\t%s\t%s\n", s.MyID, host, ports, installPath)
}
w.Flush()

Expand All @@ -91,5 +92,12 @@ func printRecoveryGuide(deployed []topology.ZKServer, topo *topology.ZKTopology)
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_<myid>\n", topo.Path)
for _, server := range deployed {
confDir := zkConfigDir(topo.Path, topo.Name, server.MyID)
dataDir := zkNodeDataPath(server.Config.DataDir, server.MyID)
dataLogDir := zkNodeDataPath(server.Config.DataLogDir, server.MyID)

fmt.Printf(" %s:\n", server.Host())
fmt.Printf(" rm -rf %s %s %s\n", confDir, dataDir, dataLogDir)
}
}
59 changes: 58 additions & 1 deletion internal/zk/ensemble.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,7 @@ package zk

import (
"fmt"
"path"

"github.com/jam2in/arcusctl/internal/store"
"github.com/jam2in/arcusctl/internal/topology"
Expand Down Expand Up @@ -32,8 +33,64 @@ func loadEnsemble(name string) (*store.ZKMeta, *topology.ZKTopology, error) {
}

func mergeServerConfigs(topo *topology.ZKTopology) {
globalConfig := topo.GlobalConfig

if topo.GlobalConfig.DataDir == "" {
globalConfig.DataDir = path.Join(topo.Path, "data", topo.Name)
}

for i := range topo.Servers {
merged := mergeConfig(topo.GlobalConfig, topo.Servers[i].Config)
merged := mergeConfig(globalConfig, topo.Servers[i].Config)
topo.Servers[i].Config = &merged
}
}

// sharingEnsemble returns the name of another ensemble
// that shares the same installation path and has a server on the given host.
func sharingEnsemble(
ensembleName string,
installPath string,
host string,
) (string, error) {
registered, err := store.ListZK()
if err != nil {
return "", err
}

for _, other := range registered {
if other == ensembleName {
continue
}

meta, err := store.LoadZKMeta(other)
if err != nil {
continue
}

topo, err := store.LoadZKTopology(other)
if err != nil {
continue
}

otherInstallPath := zkInstallPath(topo.Path, meta.Version)
if otherInstallPath == installPath &&
hasServerOn(topo.Servers, host) {
return other, nil
}
}

return "", nil
}

func hasServerOn(
zkServers []topology.ZKServer,
host string,
) bool {
for _, server := range zkServers {
if server.Host() == host {
return true
}
}

return false
}
Loading
Loading