From 7e93abb9722602a6e1d36668199b6222838a71ab Mon Sep 17 00:00:00 2001 From: Pavel Sviderski Date: Fri, 6 Sep 2024 13:03:45 +1000 Subject: [PATCH] move InitCluster to Machine --- cmd/uncloud/machine/init.go | 9 ++- internal/daemon/daemon.go | 94 ----------------------- internal/machine/cluster/state.go | 6 ++ internal/machine/machine.go | 123 ++++++++++++++++++++++++++---- internal/machine/state.go | 19 +++-- 5 files changed, 131 insertions(+), 120 deletions(-) diff --git a/cmd/uncloud/machine/init.go b/cmd/uncloud/machine/init.go index b6f49035..99456b13 100644 --- a/cmd/uncloud/machine/init.go +++ b/cmd/uncloud/machine/init.go @@ -6,7 +6,6 @@ import ( "net/netip" "uncloud/internal/machine" "uncloud/internal/machine/api/pb" - "uncloud/internal/machine/daemon" "uncloud/internal/machine/network" "uncloud/internal/secret" ) @@ -44,7 +43,13 @@ func NewInitCommand() *cobra.Command { users = append(users, user) } - if err = daemon.InitCluster(opts.dataDir, opts.name, netPrefix, users); err != nil { + // TODO: ideally this should be an RPC call to the machine API via unix socket. + config := &machine.Config{DataDir: opts.dataDir} + mach, err := machine.NewMachine(config) + if err != nil { + return fmt.Errorf("init machine: %w", err) + } + if err = mach.InitCluster(opts.name, netPrefix, users); err != nil { return fmt.Errorf("initialise cluster: %w", err) } return nil diff --git a/internal/daemon/daemon.go b/internal/daemon/daemon.go index d2a3436e..68723314 100644 --- a/internal/daemon/daemon.go +++ b/internal/daemon/daemon.go @@ -4,103 +4,9 @@ import ( "context" "fmt" "log/slog" - "net/netip" "uncloud/internal/machine" - "uncloud/internal/machine/api/pb" - "uncloud/internal/machine/cluster" - "uncloud/internal/machine/network" ) -// InitCluster resets the local machine and initialises a new cluster with it. -// TODO: ideally, this should be an RPC call to the daemon API to correctly handle the leave request and reconfiguration. -func InitCluster(dataDir, machineName string, netPrefix netip.Prefix, users []*pb.User) error { - var err error - if machineName == "" { - machineName, err = cluster.NewRandomMachineName() - if err != nil { - return fmt.Errorf("generate machine name: %w", err) - } - } - privKey, pubKey, err := network.NewMachineKeys() - if err != nil { - return fmt.Errorf("generate machine keys: %w", err) - } - - state := cluster.NewState(cluster.StatePath(dataDir)) - c := cluster.NewServer(state) - if err = c.SetNetwork(netPrefix); err != nil { - return fmt.Errorf("set cluster network: %w", err) - } - - // Use all routable addresses as endpoints. - addrs, err := network.ListRoutableIPs() - if err != nil { - return fmt.Errorf("list routable addresses: %w", err) - } - endpoints := make([]*pb.IPPort, len(addrs)) - for i, addr := range addrs { - addrPort := netip.AddrPortFrom(addr, network.WireGuardPort) - endpoints[i] = pb.NewIPPort(addrPort) - } - // Register the new machine in the cluster to populate the state and get its ID and subnet. - req := &pb.AddMachineRequest{ - Name: machineName, - Network: &pb.NetworkConfig{ - Endpoints: endpoints, - PublicKey: pubKey, - }, - } - resp, err := c.AddMachine(context.Background(), req) - if err != nil { - return fmt.Errorf("add machine to cluster: %w", err) - } - - m := resp.Machine - subnet, err := m.Network.Subnet.ToPrefix() - if err != nil { - return err - } - manageIP, err := m.Network.ManagementIp.ToAddr() - if err != nil { - return err - } - mcfg := &machine.State{ - ID: m.Id, - Name: m.Name, - Network: &network.Config{ - Subnet: subnet, - ManagementIP: manageIP, - PrivateKey: privKey, - PublicKey: pubKey, - }, - } - - // Add users to the cluster and build peers config from them. - peers := make([]network.PeerConfig, len(users)) - for i, u := range users { - if err = c.AddUser(u); err != nil { - return fmt.Errorf("add user to cluster: %w", err) - } - userManageIP, uErr := u.Network.ManagementIp.ToAddr() - if uErr != nil { - return uErr - } - peers[i] = network.PeerConfig{ - ManagementIP: userManageIP, - PublicKey: u.Network.PublicKey, - } - } - mcfg.Network.Peers = peers - - mcfg.SetPath(machine.StatePath(dataDir)) - if err = mcfg.Save(); err != nil { - return fmt.Errorf("save machine config: %w", err) - } - - fmt.Printf("Cluster initialised with machine %q\n", m.Name) - return nil -} - type Daemon struct { machine *machine.Machine } diff --git a/internal/machine/cluster/state.go b/internal/machine/cluster/state.go index 3a3490b8..d5682123 100644 --- a/internal/machine/cluster/state.go +++ b/internal/machine/cluster/state.go @@ -32,6 +32,12 @@ func (s *State) Load() error { if err = proto.Unmarshal(data, s.State); err != nil { return fmt.Errorf("parse state file %q: %w", s.path, err) } + if s.State.Machines == nil { + s.State.Machines = make(map[string]*pb.Machine) + } + if s.State.Endpoints == nil { + s.State.Endpoints = make(map[string]*pb.MachineEndpoints) + } return nil } diff --git a/internal/machine/machine.go b/internal/machine/machine.go index 71593421..c4ec9606 100644 --- a/internal/machine/machine.go +++ b/internal/machine/machine.go @@ -27,6 +27,7 @@ type Machine struct { state *State networkServer *grpc.Server + clusterState *cluster.State cluster *cluster.Server // TODO: create localServer for unix socket. } @@ -59,28 +60,25 @@ func NewMachine(config *Config) (*Machine, error) { } } + m := &Machine{ + config: *config, + state: state, + networkServer: grpc.NewServer(), + } + clusterStatePath := cluster.StatePath(config.DataDir) clusterState := cluster.NewState(clusterStatePath) if err = clusterState.Load(); err != nil { if !errors.Is(err, os.ErrNotExist) { return nil, fmt.Errorf("load cluster state: %w", err) } - slog.Info("Cluster state file not found, creating a new one.", "path", clusterStatePath) - if err = clusterState.Save(); err != nil { - return nil, fmt.Errorf("save cluster state: %w", err) - } + } else { + // Cluster state is successfully loaded, start the cluster server. + m.cluster = cluster.NewServer(clusterState) + pb.RegisterClusterServer(m.networkServer, m.cluster) } - clusterServer := cluster.NewServer(clusterState) - networkServer := grpc.NewServer() - pb.RegisterClusterServer(networkServer, clusterServer) - - return &Machine{ - config: *config, - state: state, - networkServer: networkServer, - cluster: clusterServer, - }, nil + return m, nil } func (m *Machine) Run(ctx context.Context) error { @@ -145,3 +143,100 @@ func (m *Machine) Run(ctx context.Context) error { return errGroup.Wait() } + +// InitCluster resets the local machine and initialises a new cluster with it. +// TODO: ideally, this should be an RPC call to the machine API to correctly handle the leave request and reconfiguration. +func (m *Machine) InitCluster(machineName string, netPrefix netip.Prefix, users []*pb.User) error { + var err error + if machineName == "" { + machineName, err = cluster.NewRandomMachineName() + if err != nil { + return fmt.Errorf("generate machine name: %w", err) + } + } + + // TODO: a proper cluster leave mechanism and machine reset should be implemented later. + // For now assume the cluster server is not running. + clusterStatePath := cluster.StatePath(m.config.DataDir) + clusterState := cluster.NewState(clusterStatePath) + if err = clusterState.Save(); err != nil { + return fmt.Errorf("save cluster state: %w", err) + } + clusterServer := cluster.NewServer(clusterState) + // TODO: register and start the cluster server when this becomes an RPC call. + + if err = clusterServer.SetNetwork(netPrefix); err != nil { + return fmt.Errorf("set cluster network: %w", err) + } + + // Use the public and all routable IPs as endpoints. + ips, err := network.ListRoutableIPs() + if err != nil { + return fmt.Errorf("list routable addresses: %w", err) + } + publicIP, err := network.GetPublicIP() + // Ignore the error if failed to get the public IP using API services. + if err == nil { + ips = append([]netip.Addr{publicIP}, ips...) + } + endpoints := make([]*pb.IPPort, len(ips)) + for i, addr := range ips { + addrPort := netip.AddrPortFrom(addr, network.WireGuardPort) + endpoints[i] = pb.NewIPPort(addrPort) + } + + // Register the new machine in the cluster to populate the state and get its ID and subnet. + // Public and private keys have already been initialised in the machine state when it was created. + req := &pb.AddMachineRequest{ + Name: machineName, + Network: &pb.NetworkConfig{ + Endpoints: endpoints, + PublicKey: m.state.Network.PublicKey, + }, + } + resp, err := clusterServer.AddMachine(context.Background(), req) + if err != nil { + return fmt.Errorf("add machine to cluster: %w", err) + } + + subnet, err := resp.Machine.Network.Subnet.ToPrefix() + if err != nil { + return err + } + manageIP, err := resp.Machine.Network.ManagementIp.ToAddr() + if err != nil { + return err + } + // Add users to the cluster and build peers config from them. + peers := make([]network.PeerConfig, len(users)) + for i, u := range users { + if err = clusterServer.AddUser(u); err != nil { + return fmt.Errorf("add user to cluster: %w", err) + } + userManageIP, uErr := u.Network.ManagementIp.ToAddr() + if uErr != nil { + return uErr + } + peers[i] = network.PeerConfig{ + ManagementIP: userManageIP, + PublicKey: u.Network.PublicKey, + } + } + + // Update the machine state with the new cluster configuration. + m.state.ID = resp.Machine.Id + m.state.Name = resp.Machine.Name + m.state.Network = &network.Config{ + Subnet: subnet, + ManagementIP: manageIP, + PrivateKey: m.state.Network.PrivateKey, + PublicKey: m.state.Network.PublicKey, + Peers: peers, + } + if err = m.state.Save(); err != nil { + return fmt.Errorf("save machine state: %w", err) + } + + fmt.Printf("Cluster initialised with machine %q\n", m.state.Name) + return nil +} diff --git a/internal/machine/state.go b/internal/machine/state.go index 74ba4265..544d0952 100644 --- a/internal/machine/state.go +++ b/internal/machine/state.go @@ -39,17 +39,16 @@ func StatePath(dataDir string) string { func ParseState(path string) (*State, error) { data, err := os.ReadFile(path) if err != nil { - return nil, fmt.Errorf("read config file: %w", err) + return nil, fmt.Errorf("read state file: %w", err) } - var config State - if err = json.Unmarshal(data, &config); err != nil { - return nil, fmt.Errorf("parse config file %q: %w", path, err) + state := State{path: path} + if err = json.Unmarshal(data, &state); err != nil { + return nil, fmt.Errorf("parse state file %q: %w", path, err) } - - if config.Network == nil { - return nil, fmt.Errorf("missing network configuration in config file %q", path) + if state.Network == nil { + return nil, fmt.Errorf("missing network configuration in state file %q", path) } - return &config, nil + return &state, nil } // SetPath sets the file path the state can be saved to. @@ -69,11 +68,11 @@ func (c *State) Encode() ([]byte, error) { // Save writes the state data to the file at the given path. func (c *State) Save() error { if c.path == "" { - return fmt.Errorf("config path not set") + return fmt.Errorf("state path not set") } dir, _ := filepath.Split(c.path) if err := os.MkdirAll(dir, 0700); err != nil { - return fmt.Errorf("create config directory %q: %w", dir, err) + return fmt.Errorf("create state directory %q: %w", dir, err) } data, err := c.Encode()