mirror of
https://github.com/psviderski/uncloud.git
synced 2026-08-26 19:13:34 +00:00
split cluster Server into Machine and Server
This commit is contained in:
@@ -21,7 +21,7 @@ import (
|
|||||||
func InitCluster(dataDir, machineName string, netPrefix netip.Prefix, users []*pb.User) error {
|
func InitCluster(dataDir, machineName string, netPrefix netip.Prefix, users []*pb.User) error {
|
||||||
var err error
|
var err error
|
||||||
if machineName == "" {
|
if machineName == "" {
|
||||||
machineName, err = machine.NewRandomName()
|
machineName, err = cluster.NewRandomMachineName()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return fmt.Errorf("generate machine name: %w", err)
|
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))
|
state := cluster.NewState(cluster.StatePath(dataDir))
|
||||||
c := cluster.NewCluster(&cluster.Config{}, state)
|
c := cluster.NewServer(state)
|
||||||
if err = c.SetNetwork(netPrefix); err != nil {
|
if err = c.SetNetwork(netPrefix); err != nil {
|
||||||
return fmt.Errorf("set cluster network: %w", err)
|
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 {
|
type Daemon struct {
|
||||||
|
machine *machine.Machine
|
||||||
|
|
||||||
state *machine.State
|
state *machine.State
|
||||||
cluster *cluster.Server
|
cluster *cluster.Server
|
||||||
}
|
}
|
||||||
@@ -154,10 +156,13 @@ func New(dataDir string) (*Daemon, error) {
|
|||||||
state: mstate,
|
state: mstate,
|
||||||
}
|
}
|
||||||
if mstate.Network.IsConfigured() {
|
if mstate.Network.IsConfigured() {
|
||||||
config := &cluster.Config{
|
config := &machine.Config{
|
||||||
APIAddr: net.JoinHostPort(mstate.Network.ManagementIP.String(), strconv.Itoa(machine.APIPort)),
|
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
|
return d, nil
|
||||||
@@ -199,7 +204,7 @@ func (d *Daemon) Run(ctx context.Context) error {
|
|||||||
if d.cluster != nil {
|
if d.cluster != nil {
|
||||||
errGroup.Go(func() error {
|
errGroup.Go(func() error {
|
||||||
slog.Info("Starting cluster.")
|
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 fmt.Errorf("cluster failed: %w", err)
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
@@ -211,7 +216,7 @@ func (d *Daemon) Run(ctx context.Context) error {
|
|||||||
<-ctx.Done()
|
<-ctx.Done()
|
||||||
if d.cluster != nil {
|
if d.cluster != nil {
|
||||||
slog.Info("Stopping cluster.")
|
slog.Info("Stopping cluster.")
|
||||||
d.cluster.Stop()
|
d.machine.Stop()
|
||||||
slog.Info("Cluster server stopped.")
|
slog.Info("Cluster server stopped.")
|
||||||
}
|
}
|
||||||
return nil
|
return nil
|
||||||
|
|||||||
@@ -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
|
||||||
|
}
|
||||||
@@ -4,66 +4,24 @@ import (
|
|||||||
"bytes"
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"fmt"
|
"fmt"
|
||||||
"google.golang.org/grpc"
|
|
||||||
"google.golang.org/grpc/codes"
|
"google.golang.org/grpc/codes"
|
||||||
"google.golang.org/grpc/status"
|
"google.golang.org/grpc/status"
|
||||||
"google.golang.org/protobuf/proto"
|
"google.golang.org/protobuf/proto"
|
||||||
"log/slog"
|
|
||||||
"net"
|
|
||||||
"net/netip"
|
"net/netip"
|
||||||
"os"
|
|
||||||
"path/filepath"
|
|
||||||
"uncloud/internal/machine"
|
|
||||||
"uncloud/internal/machine/api/pb"
|
"uncloud/internal/machine/api/pb"
|
||||||
"uncloud/internal/machine/network"
|
"uncloud/internal/machine/network"
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
|
||||||
StateFile = "cluster.pb"
|
|
||||||
)
|
|
||||||
|
|
||||||
type Config struct {
|
|
||||||
APIAddr string
|
|
||||||
APISockPath string
|
|
||||||
}
|
|
||||||
|
|
||||||
type Server struct {
|
type Server struct {
|
||||||
// TODO: implement grpc Server
|
|
||||||
pb.UnimplementedClusterServer
|
pb.UnimplementedClusterServer
|
||||||
|
|
||||||
config Config
|
state *State
|
||||||
state *State
|
|
||||||
|
|
||||||
server *grpc.Server
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func NewCluster(config *Config, state *State) *Server {
|
func NewServer(state *State) *Server {
|
||||||
c := &Server{
|
return &Server{
|
||||||
config: *config,
|
state: state,
|
||||||
state: state,
|
|
||||||
server: grpc.NewServer(),
|
|
||||||
}
|
}
|
||||||
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) {
|
func (c *Server) Network() (netip.Prefix, error) {
|
||||||
@@ -109,7 +67,7 @@ func (c *Server) AddMachine(ctx context.Context, req *pb.AddMachineRequest) (*pb
|
|||||||
i++
|
i++
|
||||||
}
|
}
|
||||||
|
|
||||||
mid, err := machine.NewID()
|
mid, err := NewMachineID()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("generate machine ID: %w", err)
|
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 == "" {
|
if m.Name == "" {
|
||||||
m.Name, err = machine.NewRandomName()
|
m.Name, err = NewRandomMachineName()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("generate machine name: %w", err)
|
return nil, fmt.Errorf("generate machine name: %w", err)
|
||||||
}
|
}
|
||||||
@@ -186,41 +144,3 @@ type State struct {
|
|||||||
State *pb.State
|
State *pb.State
|
||||||
path string
|
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)
|
|
||||||
}
|
|
||||||
@@ -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)
|
||||||
|
}
|
||||||
@@ -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.")
|
||||||
|
}
|
||||||
@@ -1,15 +1,13 @@
|
|||||||
package machine
|
package machine
|
||||||
|
|
||||||
import (
|
import (
|
||||||
"crypto/rand"
|
|
||||||
"encoding/json"
|
"encoding/json"
|
||||||
"fmt"
|
"fmt"
|
||||||
"math/big"
|
|
||||||
"net/netip"
|
"net/netip"
|
||||||
"os"
|
"os"
|
||||||
"path/filepath"
|
"path/filepath"
|
||||||
|
"uncloud/internal/machine/cluster"
|
||||||
"uncloud/internal/machine/network"
|
"uncloud/internal/machine/network"
|
||||||
"uncloud/internal/secret"
|
|
||||||
)
|
)
|
||||||
|
|
||||||
const (
|
const (
|
||||||
@@ -85,33 +83,14 @@ func (c *State) Save() error {
|
|||||||
return os.WriteFile(c.path, data, 0600)
|
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.
|
// 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) {
|
func NewBootstrapConfig(name string, subnet netip.Prefix, peers ...network.PeerConfig) (*State, error) {
|
||||||
mid, err := NewID()
|
mid, err := cluster.NewMachineID()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("generate machine ID: %w", err)
|
return nil, fmt.Errorf("generate machine ID: %w", err)
|
||||||
}
|
}
|
||||||
if name == "" {
|
if name == "" {
|
||||||
name, err = NewRandomName()
|
name, err = cluster.NewRandomMachineName()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
return nil, fmt.Errorf("generate machine name: %w", err)
|
return nil, fmt.Errorf("generate machine name: %w", err)
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user