diff --git a/internal/daemon/daemon.go b/internal/daemon/daemon.go index ae783589..1dd59755 100644 --- a/internal/daemon/daemon.go +++ b/internal/daemon/daemon.go @@ -21,7 +21,7 @@ import ( func InitCluster(dataDir, machineName string, netPrefix netip.Prefix, users []*pb.User) error { var err error if machineName == "" { - machineName, err = machine.NewRandomName() + machineName, err = cluster.NewRandomMachineName() if err != nil { return fmt.Errorf("generate machine name: %w", err) } @@ -32,7 +32,7 @@ func InitCluster(dataDir, machineName string, netPrefix netip.Prefix, users []*p } state := cluster.NewState(cluster.StatePath(dataDir)) - c := cluster.NewCluster(&cluster.Config{}, state) + c := cluster.NewServer(state) if err = c.SetNetwork(netPrefix); err != nil { return fmt.Errorf("set cluster network: %w", err) } @@ -107,6 +107,8 @@ func InitCluster(dataDir, machineName string, netPrefix netip.Prefix, users []*p } type Daemon struct { + machine *machine.Machine + state *machine.State cluster *cluster.Server } @@ -154,10 +156,13 @@ func New(dataDir string) (*Daemon, error) { state: mstate, } if mstate.Network.IsConfigured() { - config := &cluster.Config{ + config := &machine.Config{ APIAddr: net.JoinHostPort(mstate.Network.ManagementIP.String(), strconv.Itoa(machine.APIPort)), } - d.cluster = cluster.NewCluster(config, cstate) + d.machine, err = machine.NewMachine(config) + if err != nil { + return nil, fmt.Errorf("init machine: %w", err) + } } return d, nil @@ -199,7 +204,7 @@ func (d *Daemon) Run(ctx context.Context) error { if d.cluster != nil { errGroup.Go(func() error { slog.Info("Starting cluster.") - if err := d.cluster.Run(); err != nil { + if err := d.machine.Run(); err != nil { return fmt.Errorf("cluster failed: %w", err) } return nil @@ -211,7 +216,7 @@ func (d *Daemon) Run(ctx context.Context) error { <-ctx.Done() if d.cluster != nil { slog.Info("Stopping cluster.") - d.cluster.Stop() + d.machine.Stop() slog.Info("Cluster server stopped.") } return nil diff --git a/internal/machine/cluster/machine.go b/internal/machine/cluster/machine.go new file mode 100644 index 00000000..56802a64 --- /dev/null +++ b/internal/machine/cluster/machine.go @@ -0,0 +1,27 @@ +package cluster + +import ( + "crypto/rand" + "fmt" + "math/big" + "uncloud/internal/secret" +) + +// NewMachineID generates a new unique machine ID. +func NewMachineID() (string, error) { + return secret.NewID() +} + +// NewRandomMachineName generates a random machine name in the format "machine-xxxx". +func NewRandomMachineName() (string, error) { + const charset = "abcdefghijklmnopqrstuvwxyz0123456789" + suffix := make([]byte, 4) + for i := range suffix { + randIdx, err := rand.Int(rand.Reader, big.NewInt(int64(len(charset)))) + if err != nil { + return "", fmt.Errorf("get random number: %w", err) + } + suffix[i] = charset[randIdx.Int64()] + } + return "machine-" + string(suffix), nil +} diff --git a/internal/machine/cluster/cluster.go b/internal/machine/cluster/server.go similarity index 67% rename from internal/machine/cluster/cluster.go rename to internal/machine/cluster/server.go index be325144..2f93c401 100644 --- a/internal/machine/cluster/cluster.go +++ b/internal/machine/cluster/server.go @@ -4,66 +4,24 @@ import ( "bytes" "context" "fmt" - "google.golang.org/grpc" "google.golang.org/grpc/codes" "google.golang.org/grpc/status" "google.golang.org/protobuf/proto" - "log/slog" - "net" "net/netip" - "os" - "path/filepath" - "uncloud/internal/machine" "uncloud/internal/machine/api/pb" "uncloud/internal/machine/network" ) -const ( - StateFile = "cluster.pb" -) - -type Config struct { - APIAddr string - APISockPath string -} - type Server struct { - // TODO: implement grpc Server pb.UnimplementedClusterServer - config Config - state *State - - server *grpc.Server + state *State } -func NewCluster(config *Config, state *State) *Server { - c := &Server{ - config: *config, - state: state, - server: grpc.NewServer(), +func NewServer(state *State) *Server { + return &Server{ + state: state, } - pb.RegisterClusterServer(c.server, c) - return c -} - -func (c *Server) Run() error { - listener, err := net.Listen("tcp", c.config.APIAddr) - if err != nil { - return fmt.Errorf("listen API port: %w", err) - } - slog.Info("Starting API server.", "addr", c.config.APIAddr) - if err = c.server.Serve(listener); err != nil { - return fmt.Errorf("API server failed: %w", err) - } - return nil -} - -func (c *Server) Stop() { - slog.Info("Stopping API server.") - // TODO: implement timeout for graceful shutdown. - c.server.GracefulStop() - slog.Info("API server stopped.") } func (c *Server) Network() (netip.Prefix, error) { @@ -109,7 +67,7 @@ func (c *Server) AddMachine(ctx context.Context, req *pb.AddMachineRequest) (*pb i++ } - mid, err := machine.NewID() + mid, err := NewMachineID() if err != nil { return nil, fmt.Errorf("generate machine ID: %w", err) } @@ -126,7 +84,7 @@ func (c *Server) AddMachine(ctx context.Context, req *pb.AddMachineRequest) (*pb }, } if m.Name == "" { - m.Name, err = machine.NewRandomName() + m.Name, err = NewRandomMachineName() if err != nil { return nil, fmt.Errorf("generate machine name: %w", err) } @@ -186,41 +144,3 @@ type State struct { State *pb.State path string } - -func StatePath(dataDir string) string { - return filepath.Join(dataDir, StateFile) -} - -func NewState(path string) *State { - return &State{ - State: &pb.State{ - Machines: make(map[string]*pb.Machine), - Endpoints: make(map[string]*pb.MachineEndpoints), - }, - path: path, - } -} - -func (s *State) Load() error { - data, err := os.ReadFile(s.path) - if err != nil { - return fmt.Errorf("read state file %q: %w", s.path, err) - } - if err = proto.Unmarshal(data, s.State); err != nil { - return fmt.Errorf("parse state file %q: %w", s.path, err) - } - return nil -} - -func (s *State) Save() error { - dir, _ := filepath.Split(s.path) - if err := os.MkdirAll(dir, 0700); err != nil { - return fmt.Errorf("create state directory %q: %w", dir, err) - } - - data, err := proto.Marshal(s.State) - if err != nil { - return fmt.Errorf("marshal state: %w", err) - } - return os.WriteFile(s.path, data, 0600) -} diff --git a/internal/machine/cluster/state.go b/internal/machine/cluster/state.go new file mode 100644 index 00000000..3a3490b8 --- /dev/null +++ b/internal/machine/cluster/state.go @@ -0,0 +1,49 @@ +package cluster + +import ( + "fmt" + "google.golang.org/protobuf/proto" + "os" + "path/filepath" + "uncloud/internal/machine/api/pb" +) + +const StateFile = "cluster.pb" + +func StatePath(dataDir string) string { + return filepath.Join(dataDir, StateFile) +} + +func NewState(path string) *State { + return &State{ + State: &pb.State{ + Machines: make(map[string]*pb.Machine), + Endpoints: make(map[string]*pb.MachineEndpoints), + }, + path: path, + } +} + +func (s *State) Load() error { + data, err := os.ReadFile(s.path) + if err != nil { + return fmt.Errorf("read state file %q: %w", s.path, err) + } + if err = proto.Unmarshal(data, s.State); err != nil { + return fmt.Errorf("parse state file %q: %w", s.path, err) + } + return nil +} + +func (s *State) Save() error { + dir, _ := filepath.Split(s.path) + if err := os.MkdirAll(dir, 0700); err != nil { + return fmt.Errorf("create state directory %q: %w", dir, err) + } + + data, err := proto.Marshal(s.State) + if err != nil { + return fmt.Errorf("marshal state: %w", err) + } + return os.WriteFile(s.path, data, 0600) +} diff --git a/internal/machine/machine.go b/internal/machine/machine.go new file mode 100644 index 00000000..c8321b3a --- /dev/null +++ b/internal/machine/machine.go @@ -0,0 +1,54 @@ +package machine + +import ( + "fmt" + "google.golang.org/grpc" + "log/slog" + "net" + "uncloud/internal/machine/api/pb" + "uncloud/internal/machine/cluster" +) + +type Config struct { + // DataDir is the directory where the machine stores its persistent state. + DataDir string + APIAddr string + APISockPath string +} + +type Machine struct { + config Config + server *grpc.Server +} + +func NewMachine(config *Config) (*Machine, error) { + m := &Machine{ + config: *config, + server: grpc.NewServer(), + } + + clusterState := cluster.NewState(cluster.StatePath(config.DataDir)) + clusterServer := cluster.NewServer(clusterState) + pb.RegisterClusterServer(m.server, clusterServer) + + return m, nil +} + +func (m *Machine) Run() error { + listener, err := net.Listen("tcp", m.config.APIAddr) + if err != nil { + return fmt.Errorf("listen API port: %w", err) + } + slog.Info("Starting API server.", "addr", m.config.APIAddr) + if err = m.server.Serve(listener); err != nil { + return fmt.Errorf("API server failed: %w", err) + } + return nil +} + +func (m *Machine) Stop() { + slog.Info("Stopping API server.") + // TODO: implement timeout for graceful shutdown. + m.server.GracefulStop() + slog.Info("API server stopped.") +} diff --git a/internal/machine/state.go b/internal/machine/state.go index 03ef4103..74ba4265 100644 --- a/internal/machine/state.go +++ b/internal/machine/state.go @@ -1,15 +1,13 @@ package machine import ( - "crypto/rand" "encoding/json" "fmt" - "math/big" "net/netip" "os" "path/filepath" + "uncloud/internal/machine/cluster" "uncloud/internal/machine/network" - "uncloud/internal/secret" ) const ( @@ -85,33 +83,14 @@ func (c *State) Save() error { return os.WriteFile(c.path, data, 0600) } -// NewID generates a new unique machine ID. -func NewID() (string, error) { - return secret.NewID() -} - -// NewRandomName generates a random machine name in the format "machine-xxxx". -func NewRandomName() (string, error) { - const charset = "abcdefghijklmnopqrstuvwxyz0123456789" - suffix := make([]byte, 4) - for i := range suffix { - randIdx, err := rand.Int(rand.Reader, big.NewInt(int64(len(charset)))) - if err != nil { - return "", fmt.Errorf("get random number: %w", err) - } - suffix[i] = charset[randIdx.Int64()] - } - return "machine-" + string(suffix), nil -} - // NewBootstrapConfig returns a new machine configuration that should be applied to the first machine in a cluster. func NewBootstrapConfig(name string, subnet netip.Prefix, peers ...network.PeerConfig) (*State, error) { - mid, err := NewID() + mid, err := cluster.NewMachineID() if err != nil { return nil, fmt.Errorf("generate machine ID: %w", err) } if name == "" { - name, err = NewRandomName() + name, err = cluster.NewRandomMachineName() if err != nil { return nil, fmt.Errorf("generate machine name: %w", err) }