Implement grpc server in the daemon to receive cluster requests

This commit is contained in:
Pavel Sviderski
2024-08-29 10:17:16 +10:00
parent e074b6855f
commit afe9c13283
8 changed files with 703 additions and 3 deletions
+9
View File
@@ -0,0 +1,9 @@
.PHONY: build
uncloudd-dev1:
GOOS=linux GOARCH=amd64 go build -o uncloudd-linux-amd64 cmd/uncloudd/main.go && \
scp uncloudd-linux-amd64 spy@192.168.40.243:~/ && \
ssh spy@192.168.40.243 sudo install ./uncloudd-linux-amd64 /usr/local/bin/uncloudd
.PHONY: proto
proto:
protoc --go_out=. --go_opt=paths=source_relative --go-grpc_out=. --go-grpc_opt=paths=source_relative internal/machine/cluster/pb/cluster.proto
+2 -2
View File
@@ -131,8 +131,8 @@ require (
golang.org/x/text v0.17.0 // indirect
golang.org/x/tools v0.24.0 // indirect
golang.zx2c4.com/wireguard v0.0.0-20231211153847-12269c276173 // indirect
google.golang.org/genproto/googleapis/rpc v0.0.0-20240617180043-68d350f18fd4 // indirect
google.golang.org/grpc v1.64.0 // indirect
google.golang.org/genproto/googleapis/rpc v0.0.0-20240826202546-f6391c0de4c7 // indirect
google.golang.org/grpc v1.65.0 // indirect
google.golang.org/protobuf v1.34.2 // indirect
gotest.tools/v3 v3.5.1 // indirect
lukechampine.com/blake3 v1.3.0 // indirect
+4
View File
@@ -732,6 +732,8 @@ google.golang.org/genproto v0.0.0-20190819201941-24fa4b261c55/go.mod h1:DMBHOl98
google.golang.org/genproto v0.0.0-20200526211855-cb27e3aa2013/go.mod h1:NbSheEEYHJ7i3ixzK3sjbqSGDJWnxyFXZblF3eUsNvo=
google.golang.org/genproto/googleapis/rpc v0.0.0-20240617180043-68d350f18fd4 h1:Di6ANFilr+S60a4S61ZM00vLdw0IrQOSMS2/6mrnOU0=
google.golang.org/genproto/googleapis/rpc v0.0.0-20240617180043-68d350f18fd4/go.mod h1:Ue6ibwXGpU+dqIcODieyLOcgj7z8+IcskoNIgZxtrFY=
google.golang.org/genproto/googleapis/rpc v0.0.0-20240826202546-f6391c0de4c7 h1:2035KHhUv+EpyB+hWgJnaWKJOdX1E95w2S8Rr4uWKTs=
google.golang.org/genproto/googleapis/rpc v0.0.0-20240826202546-f6391c0de4c7/go.mod h1:UqMtugtsSgubUsoxbuAoiCXvqvErP7Gf0so0mK9tHxU=
google.golang.org/grpc v1.19.0/go.mod h1:mqu4LbDTu4XGKhr4mRzUsmM4RtVoemTSY81AxZiDr8c=
google.golang.org/grpc v1.20.1/go.mod h1:10oTOabMzJvdu6/UiuZezV6QK5dSlG84ov/aaiqXj38=
google.golang.org/grpc v1.23.0/go.mod h1:Y5yQAOtifL1yxbo5wqy6BxZv8vAUGQwXBOALyacEbxg=
@@ -740,6 +742,8 @@ google.golang.org/grpc v1.27.0/go.mod h1:qbnxyOmOxrQa7FizSgH+ReBfzJrCY1pSN7KXBS8
google.golang.org/grpc v1.33.2/go.mod h1:JMHMWHQWaTccqQQlmk3MJZS+GWXOdAesneDmEnv2fbc=
google.golang.org/grpc v1.64.0 h1:KH3VH9y/MgNQg1dE7b3XfVK0GsPSIzJwdF617gUSbvY=
google.golang.org/grpc v1.64.0/go.mod h1:oxjF8E3FBnjp+/gVFYdWacaLDx9na1aqy9oovLpxQYg=
google.golang.org/grpc v1.65.0 h1:bs/cUb4lp1G5iImFFd3u5ixQzweKizoZJAwBNLR42lc=
google.golang.org/grpc v1.65.0/go.mod h1:WgYC2ypjlB0EiQi6wdKixMqukr6lBc0Vo+oOgjrM5ZQ=
google.golang.org/protobuf v0.0.0-20200109180630-ec00e32a8dfd/go.mod h1:DFci5gLYBciE7Vtevhsrf46CRTquxDuWsQurQQe4oz8=
google.golang.org/protobuf v0.0.0-20200221191635-4d8936d0db64/go.mod h1:kwYJMbMJ01Woi6D6+Kah6886xMZcty6N08ah7+eCXa0=
google.golang.org/protobuf v0.0.0-20200228230310-ab0ca4ff8a60/go.mod h1:cfTl7dwQJ+fmap5saPgwCLgHXTUD7jkjRqWcaiX5VyM=
+468
View File
@@ -0,0 +1,468 @@
// Code generated by protoc-gen-go. DO NOT EDIT.
// versions:
// protoc-gen-go v1.34.2
// protoc v5.27.3
// source: internal/machine/cluster/pb/cluster.proto
package pb
import (
protoreflect "google.golang.org/protobuf/reflect/protoreflect"
protoimpl "google.golang.org/protobuf/runtime/protoimpl"
reflect "reflect"
sync "sync"
)
const (
// Verify that this generated code is sufficiently up-to-date.
_ = protoimpl.EnforceVersion(20 - protoimpl.MinVersion)
// Verify that runtime/protoimpl is sufficiently up-to-date.
_ = protoimpl.EnforceVersion(protoimpl.MaxVersion - 20)
)
type MachineInfo struct {
state protoimpl.MessageState
sizeCache protoimpl.SizeCache
unknownFields protoimpl.UnknownFields
Id []byte `protobuf:"bytes,1,opt,name=id,proto3" json:"id,omitempty"`
Name string `protobuf:"bytes,2,opt,name=name,proto3" json:"name,omitempty"`
Subnet *IPPrefix `protobuf:"bytes,3,opt,name=subnet,proto3" json:"subnet,omitempty"`
Endpoints []*IPPort `protobuf:"bytes,4,rep,name=endpoints,proto3" json:"endpoints,omitempty"`
PublicKey []byte `protobuf:"bytes,5,opt,name=publicKey,proto3" json:"publicKey,omitempty"`
}
func (x *MachineInfo) Reset() {
*x = MachineInfo{}
if protoimpl.UnsafeEnabled {
mi := &file_internal_machine_cluster_pb_cluster_proto_msgTypes[0]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
}
func (x *MachineInfo) String() string {
return protoimpl.X.MessageStringOf(x)
}
func (*MachineInfo) ProtoMessage() {}
func (x *MachineInfo) ProtoReflect() protoreflect.Message {
mi := &file_internal_machine_cluster_pb_cluster_proto_msgTypes[0]
if protoimpl.UnsafeEnabled && x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
ms.StoreMessageInfo(mi)
}
return ms
}
return mi.MessageOf(x)
}
// Deprecated: Use MachineInfo.ProtoReflect.Descriptor instead.
func (*MachineInfo) Descriptor() ([]byte, []int) {
return file_internal_machine_cluster_pb_cluster_proto_rawDescGZIP(), []int{0}
}
func (x *MachineInfo) GetId() []byte {
if x != nil {
return x.Id
}
return nil
}
func (x *MachineInfo) GetName() string {
if x != nil {
return x.Name
}
return ""
}
func (x *MachineInfo) GetSubnet() *IPPrefix {
if x != nil {
return x.Subnet
}
return nil
}
func (x *MachineInfo) GetEndpoints() []*IPPort {
if x != nil {
return x.Endpoints
}
return nil
}
func (x *MachineInfo) GetPublicKey() []byte {
if x != nil {
return x.PublicKey
}
return nil
}
type AddMachineRequest struct {
state protoimpl.MessageState
sizeCache protoimpl.SizeCache
unknownFields protoimpl.UnknownFields
Machine *MachineInfo `protobuf:"bytes,1,opt,name=machine,proto3" json:"machine,omitempty"`
}
func (x *AddMachineRequest) Reset() {
*x = AddMachineRequest{}
if protoimpl.UnsafeEnabled {
mi := &file_internal_machine_cluster_pb_cluster_proto_msgTypes[1]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
}
func (x *AddMachineRequest) String() string {
return protoimpl.X.MessageStringOf(x)
}
func (*AddMachineRequest) ProtoMessage() {}
func (x *AddMachineRequest) ProtoReflect() protoreflect.Message {
mi := &file_internal_machine_cluster_pb_cluster_proto_msgTypes[1]
if protoimpl.UnsafeEnabled && x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
ms.StoreMessageInfo(mi)
}
return ms
}
return mi.MessageOf(x)
}
// Deprecated: Use AddMachineRequest.ProtoReflect.Descriptor instead.
func (*AddMachineRequest) Descriptor() ([]byte, []int) {
return file_internal_machine_cluster_pb_cluster_proto_rawDescGZIP(), []int{1}
}
func (x *AddMachineRequest) GetMachine() *MachineInfo {
if x != nil {
return x.Machine
}
return nil
}
type AddMachineResponse struct {
state protoimpl.MessageState
sizeCache protoimpl.SizeCache
unknownFields protoimpl.UnknownFields
Machine *MachineInfo `protobuf:"bytes,1,opt,name=machine,proto3" json:"machine,omitempty"`
}
func (x *AddMachineResponse) Reset() {
*x = AddMachineResponse{}
if protoimpl.UnsafeEnabled {
mi := &file_internal_machine_cluster_pb_cluster_proto_msgTypes[2]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
}
func (x *AddMachineResponse) String() string {
return protoimpl.X.MessageStringOf(x)
}
func (*AddMachineResponse) ProtoMessage() {}
func (x *AddMachineResponse) ProtoReflect() protoreflect.Message {
mi := &file_internal_machine_cluster_pb_cluster_proto_msgTypes[2]
if protoimpl.UnsafeEnabled && x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
ms.StoreMessageInfo(mi)
}
return ms
}
return mi.MessageOf(x)
}
// Deprecated: Use AddMachineResponse.ProtoReflect.Descriptor instead.
func (*AddMachineResponse) Descriptor() ([]byte, []int) {
return file_internal_machine_cluster_pb_cluster_proto_rawDescGZIP(), []int{2}
}
func (x *AddMachineResponse) GetMachine() *MachineInfo {
if x != nil {
return x.Machine
}
return nil
}
type IPPort struct {
state protoimpl.MessageState
sizeCache protoimpl.SizeCache
unknownFields protoimpl.UnknownFields
Ip []byte `protobuf:"bytes,1,opt,name=ip,proto3" json:"ip,omitempty"`
Port int32 `protobuf:"varint,2,opt,name=port,proto3" json:"port,omitempty"`
}
func (x *IPPort) Reset() {
*x = IPPort{}
if protoimpl.UnsafeEnabled {
mi := &file_internal_machine_cluster_pb_cluster_proto_msgTypes[3]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
}
func (x *IPPort) String() string {
return protoimpl.X.MessageStringOf(x)
}
func (*IPPort) ProtoMessage() {}
func (x *IPPort) ProtoReflect() protoreflect.Message {
mi := &file_internal_machine_cluster_pb_cluster_proto_msgTypes[3]
if protoimpl.UnsafeEnabled && x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
ms.StoreMessageInfo(mi)
}
return ms
}
return mi.MessageOf(x)
}
// Deprecated: Use IPPort.ProtoReflect.Descriptor instead.
func (*IPPort) Descriptor() ([]byte, []int) {
return file_internal_machine_cluster_pb_cluster_proto_rawDescGZIP(), []int{3}
}
func (x *IPPort) GetIp() []byte {
if x != nil {
return x.Ip
}
return nil
}
func (x *IPPort) GetPort() int32 {
if x != nil {
return x.Port
}
return 0
}
type IPPrefix struct {
state protoimpl.MessageState
sizeCache protoimpl.SizeCache
unknownFields protoimpl.UnknownFields
Ip []byte `protobuf:"bytes,1,opt,name=ip,proto3" json:"ip,omitempty"`
Bits int32 `protobuf:"varint,2,opt,name=bits,proto3" json:"bits,omitempty"`
}
func (x *IPPrefix) Reset() {
*x = IPPrefix{}
if protoimpl.UnsafeEnabled {
mi := &file_internal_machine_cluster_pb_cluster_proto_msgTypes[4]
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
ms.StoreMessageInfo(mi)
}
}
func (x *IPPrefix) String() string {
return protoimpl.X.MessageStringOf(x)
}
func (*IPPrefix) ProtoMessage() {}
func (x *IPPrefix) ProtoReflect() protoreflect.Message {
mi := &file_internal_machine_cluster_pb_cluster_proto_msgTypes[4]
if protoimpl.UnsafeEnabled && x != nil {
ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x))
if ms.LoadMessageInfo() == nil {
ms.StoreMessageInfo(mi)
}
return ms
}
return mi.MessageOf(x)
}
// Deprecated: Use IPPrefix.ProtoReflect.Descriptor instead.
func (*IPPrefix) Descriptor() ([]byte, []int) {
return file_internal_machine_cluster_pb_cluster_proto_rawDescGZIP(), []int{4}
}
func (x *IPPrefix) GetIp() []byte {
if x != nil {
return x.Ip
}
return nil
}
func (x *IPPrefix) GetBits() int32 {
if x != nil {
return x.Bits
}
return 0
}
var File_internal_machine_cluster_pb_cluster_proto protoreflect.FileDescriptor
var file_internal_machine_cluster_pb_cluster_proto_rawDesc = []byte{
0x0a, 0x29, 0x69, 0x6e, 0x74, 0x65, 0x72, 0x6e, 0x61, 0x6c, 0x2f, 0x6d, 0x61, 0x63, 0x68, 0x69,
0x6e, 0x65, 0x2f, 0x63, 0x6c, 0x75, 0x73, 0x74, 0x65, 0x72, 0x2f, 0x70, 0x62, 0x2f, 0x63, 0x6c,
0x75, 0x73, 0x74, 0x65, 0x72, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x12, 0x07, 0x63, 0x6c, 0x75,
0x73, 0x74, 0x65, 0x72, 0x22, 0xa9, 0x01, 0x0a, 0x0b, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65,
0x49, 0x6e, 0x66, 0x6f, 0x12, 0x0e, 0x0a, 0x02, 0x69, 0x64, 0x18, 0x01, 0x20, 0x01, 0x28, 0x0c,
0x52, 0x02, 0x69, 0x64, 0x12, 0x12, 0x0a, 0x04, 0x6e, 0x61, 0x6d, 0x65, 0x18, 0x02, 0x20, 0x01,
0x28, 0x09, 0x52, 0x04, 0x6e, 0x61, 0x6d, 0x65, 0x12, 0x29, 0x0a, 0x06, 0x73, 0x75, 0x62, 0x6e,
0x65, 0x74, 0x18, 0x03, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x11, 0x2e, 0x63, 0x6c, 0x75, 0x73, 0x74,
0x65, 0x72, 0x2e, 0x49, 0x50, 0x50, 0x72, 0x65, 0x66, 0x69, 0x78, 0x52, 0x06, 0x73, 0x75, 0x62,
0x6e, 0x65, 0x74, 0x12, 0x2d, 0x0a, 0x09, 0x65, 0x6e, 0x64, 0x70, 0x6f, 0x69, 0x6e, 0x74, 0x73,
0x18, 0x04, 0x20, 0x03, 0x28, 0x0b, 0x32, 0x0f, 0x2e, 0x63, 0x6c, 0x75, 0x73, 0x74, 0x65, 0x72,
0x2e, 0x49, 0x50, 0x50, 0x6f, 0x72, 0x74, 0x52, 0x09, 0x65, 0x6e, 0x64, 0x70, 0x6f, 0x69, 0x6e,
0x74, 0x73, 0x12, 0x1c, 0x0a, 0x09, 0x70, 0x75, 0x62, 0x6c, 0x69, 0x63, 0x4b, 0x65, 0x79, 0x18,
0x05, 0x20, 0x01, 0x28, 0x0c, 0x52, 0x09, 0x70, 0x75, 0x62, 0x6c, 0x69, 0x63, 0x4b, 0x65, 0x79,
0x22, 0x43, 0x0a, 0x11, 0x41, 0x64, 0x64, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x52, 0x65,
0x71, 0x75, 0x65, 0x73, 0x74, 0x12, 0x2e, 0x0a, 0x07, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65,
0x18, 0x01, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x14, 0x2e, 0x63, 0x6c, 0x75, 0x73, 0x74, 0x65, 0x72,
0x2e, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x49, 0x6e, 0x66, 0x6f, 0x52, 0x07, 0x6d, 0x61,
0x63, 0x68, 0x69, 0x6e, 0x65, 0x22, 0x44, 0x0a, 0x12, 0x41, 0x64, 0x64, 0x4d, 0x61, 0x63, 0x68,
0x69, 0x6e, 0x65, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x12, 0x2e, 0x0a, 0x07, 0x6d,
0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x18, 0x01, 0x20, 0x01, 0x28, 0x0b, 0x32, 0x14, 0x2e, 0x63,
0x6c, 0x75, 0x73, 0x74, 0x65, 0x72, 0x2e, 0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x49, 0x6e,
0x66, 0x6f, 0x52, 0x07, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x22, 0x2c, 0x0a, 0x06, 0x49,
0x50, 0x50, 0x6f, 0x72, 0x74, 0x12, 0x0e, 0x0a, 0x02, 0x69, 0x70, 0x18, 0x01, 0x20, 0x01, 0x28,
0x0c, 0x52, 0x02, 0x69, 0x70, 0x12, 0x12, 0x0a, 0x04, 0x70, 0x6f, 0x72, 0x74, 0x18, 0x02, 0x20,
0x01, 0x28, 0x05, 0x52, 0x04, 0x70, 0x6f, 0x72, 0x74, 0x22, 0x2e, 0x0a, 0x08, 0x49, 0x50, 0x50,
0x72, 0x65, 0x66, 0x69, 0x78, 0x12, 0x0e, 0x0a, 0x02, 0x69, 0x70, 0x18, 0x01, 0x20, 0x01, 0x28,
0x0c, 0x52, 0x02, 0x69, 0x70, 0x12, 0x12, 0x0a, 0x04, 0x62, 0x69, 0x74, 0x73, 0x18, 0x02, 0x20,
0x01, 0x28, 0x05, 0x52, 0x04, 0x62, 0x69, 0x74, 0x73, 0x32, 0x50, 0x0a, 0x07, 0x43, 0x6c, 0x75,
0x73, 0x74, 0x65, 0x72, 0x12, 0x45, 0x0a, 0x0a, 0x41, 0x64, 0x64, 0x4d, 0x61, 0x63, 0x68, 0x69,
0x6e, 0x65, 0x12, 0x1a, 0x2e, 0x63, 0x6c, 0x75, 0x73, 0x74, 0x65, 0x72, 0x2e, 0x41, 0x64, 0x64,
0x4d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, 0x1b,
0x2e, 0x63, 0x6c, 0x75, 0x73, 0x74, 0x65, 0x72, 0x2e, 0x41, 0x64, 0x64, 0x4d, 0x61, 0x63, 0x68,
0x69, 0x6e, 0x65, 0x52, 0x65, 0x73, 0x70, 0x6f, 0x6e, 0x73, 0x65, 0x42, 0x3b, 0x5a, 0x39, 0x67,
0x69, 0x74, 0x68, 0x75, 0x62, 0x2e, 0x63, 0x6f, 0x6d, 0x2f, 0x70, 0x73, 0x76, 0x69, 0x64, 0x65,
0x72, 0x73, 0x6b, 0x69, 0x2f, 0x75, 0x6e, 0x63, 0x6c, 0x6f, 0x75, 0x64, 0x2f, 0x69, 0x6e, 0x74,
0x65, 0x72, 0x6e, 0x61, 0x6c, 0x2f, 0x6d, 0x61, 0x63, 0x68, 0x69, 0x6e, 0x65, 0x2f, 0x63, 0x6c,
0x75, 0x73, 0x74, 0x65, 0x72, 0x2f, 0x70, 0x62, 0x62, 0x06, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x33,
}
var (
file_internal_machine_cluster_pb_cluster_proto_rawDescOnce sync.Once
file_internal_machine_cluster_pb_cluster_proto_rawDescData = file_internal_machine_cluster_pb_cluster_proto_rawDesc
)
func file_internal_machine_cluster_pb_cluster_proto_rawDescGZIP() []byte {
file_internal_machine_cluster_pb_cluster_proto_rawDescOnce.Do(func() {
file_internal_machine_cluster_pb_cluster_proto_rawDescData = protoimpl.X.CompressGZIP(file_internal_machine_cluster_pb_cluster_proto_rawDescData)
})
return file_internal_machine_cluster_pb_cluster_proto_rawDescData
}
var file_internal_machine_cluster_pb_cluster_proto_msgTypes = make([]protoimpl.MessageInfo, 5)
var file_internal_machine_cluster_pb_cluster_proto_goTypes = []any{
(*MachineInfo)(nil), // 0: cluster.MachineInfo
(*AddMachineRequest)(nil), // 1: cluster.AddMachineRequest
(*AddMachineResponse)(nil), // 2: cluster.AddMachineResponse
(*IPPort)(nil), // 3: cluster.IPPort
(*IPPrefix)(nil), // 4: cluster.IPPrefix
}
var file_internal_machine_cluster_pb_cluster_proto_depIdxs = []int32{
4, // 0: cluster.MachineInfo.subnet:type_name -> cluster.IPPrefix
3, // 1: cluster.MachineInfo.endpoints:type_name -> cluster.IPPort
0, // 2: cluster.AddMachineRequest.machine:type_name -> cluster.MachineInfo
0, // 3: cluster.AddMachineResponse.machine:type_name -> cluster.MachineInfo
1, // 4: cluster.Cluster.AddMachine:input_type -> cluster.AddMachineRequest
2, // 5: cluster.Cluster.AddMachine:output_type -> cluster.AddMachineResponse
5, // [5:6] is the sub-list for method output_type
4, // [4:5] is the sub-list for method input_type
4, // [4:4] is the sub-list for extension type_name
4, // [4:4] is the sub-list for extension extendee
0, // [0:4] is the sub-list for field type_name
}
func init() { file_internal_machine_cluster_pb_cluster_proto_init() }
func file_internal_machine_cluster_pb_cluster_proto_init() {
if File_internal_machine_cluster_pb_cluster_proto != nil {
return
}
if !protoimpl.UnsafeEnabled {
file_internal_machine_cluster_pb_cluster_proto_msgTypes[0].Exporter = func(v any, i int) any {
switch v := v.(*MachineInfo); i {
case 0:
return &v.state
case 1:
return &v.sizeCache
case 2:
return &v.unknownFields
default:
return nil
}
}
file_internal_machine_cluster_pb_cluster_proto_msgTypes[1].Exporter = func(v any, i int) any {
switch v := v.(*AddMachineRequest); i {
case 0:
return &v.state
case 1:
return &v.sizeCache
case 2:
return &v.unknownFields
default:
return nil
}
}
file_internal_machine_cluster_pb_cluster_proto_msgTypes[2].Exporter = func(v any, i int) any {
switch v := v.(*AddMachineResponse); i {
case 0:
return &v.state
case 1:
return &v.sizeCache
case 2:
return &v.unknownFields
default:
return nil
}
}
file_internal_machine_cluster_pb_cluster_proto_msgTypes[3].Exporter = func(v any, i int) any {
switch v := v.(*IPPort); i {
case 0:
return &v.state
case 1:
return &v.sizeCache
case 2:
return &v.unknownFields
default:
return nil
}
}
file_internal_machine_cluster_pb_cluster_proto_msgTypes[4].Exporter = func(v any, i int) any {
switch v := v.(*IPPrefix); i {
case 0:
return &v.state
case 1:
return &v.sizeCache
case 2:
return &v.unknownFields
default:
return nil
}
}
}
type x struct{}
out := protoimpl.TypeBuilder{
File: protoimpl.DescBuilder{
GoPackagePath: reflect.TypeOf(x{}).PkgPath(),
RawDescriptor: file_internal_machine_cluster_pb_cluster_proto_rawDesc,
NumEnums: 0,
NumMessages: 5,
NumExtensions: 0,
NumServices: 1,
},
GoTypes: file_internal_machine_cluster_pb_cluster_proto_goTypes,
DependencyIndexes: file_internal_machine_cluster_pb_cluster_proto_depIdxs,
MessageInfos: file_internal_machine_cluster_pb_cluster_proto_msgTypes,
}.Build()
File_internal_machine_cluster_pb_cluster_proto = out.File
file_internal_machine_cluster_pb_cluster_proto_rawDesc = nil
file_internal_machine_cluster_pb_cluster_proto_goTypes = nil
file_internal_machine_cluster_pb_cluster_proto_depIdxs = nil
}
+35
View File
@@ -0,0 +1,35 @@
syntax = "proto3";
package cluster;
option go_package = "github.com/psviderski/uncloud/internal/machine/cluster/pb";
service Cluster {
rpc AddMachine(AddMachineRequest) returns (AddMachineResponse);
}
message MachineInfo {
bytes id = 1;
string name = 2;
IPPrefix subnet = 3;
repeated IPPort endpoints = 4;
bytes publicKey = 5;
}
message AddMachineRequest {
MachineInfo machine = 1;
}
message AddMachineResponse {
MachineInfo machine = 1;
}
message IPPort {
bytes ip = 1;
int32 port = 2;
}
message IPPrefix {
bytes ip = 1;
int32 bits = 2;
}
+121
View File
@@ -0,0 +1,121 @@
// Code generated by protoc-gen-go-grpc. DO NOT EDIT.
// versions:
// - protoc-gen-go-grpc v1.5.1
// - protoc v5.27.3
// source: internal/machine/cluster/pb/cluster.proto
package pb
import (
context "context"
grpc "google.golang.org/grpc"
codes "google.golang.org/grpc/codes"
status "google.golang.org/grpc/status"
)
// This is a compile-time assertion to ensure that this generated file
// is compatible with the grpc package it is being compiled against.
// Requires gRPC-Go v1.64.0 or later.
const _ = grpc.SupportPackageIsVersion9
const (
Cluster_AddMachine_FullMethodName = "/cluster.Cluster/AddMachine"
)
// ClusterClient is the client API for Cluster service.
//
// For semantics around ctx use and closing/ending streaming RPCs, please refer to https://pkg.go.dev/google.golang.org/grpc/?tab=doc#ClientConn.NewStream.
type ClusterClient interface {
AddMachine(ctx context.Context, in *AddMachineRequest, opts ...grpc.CallOption) (*AddMachineResponse, error)
}
type clusterClient struct {
cc grpc.ClientConnInterface
}
func NewClusterClient(cc grpc.ClientConnInterface) ClusterClient {
return &clusterClient{cc}
}
func (c *clusterClient) AddMachine(ctx context.Context, in *AddMachineRequest, opts ...grpc.CallOption) (*AddMachineResponse, error) {
cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...)
out := new(AddMachineResponse)
err := c.cc.Invoke(ctx, Cluster_AddMachine_FullMethodName, in, out, cOpts...)
if err != nil {
return nil, err
}
return out, nil
}
// ClusterServer is the server API for Cluster service.
// All implementations must embed UnimplementedClusterServer
// for forward compatibility.
type ClusterServer interface {
AddMachine(context.Context, *AddMachineRequest) (*AddMachineResponse, error)
mustEmbedUnimplementedClusterServer()
}
// UnimplementedClusterServer must be embedded to have
// forward compatible implementations.
//
// NOTE: this should be embedded by value instead of pointer to avoid a nil
// pointer dereference when methods are called.
type UnimplementedClusterServer struct{}
func (UnimplementedClusterServer) AddMachine(context.Context, *AddMachineRequest) (*AddMachineResponse, error) {
return nil, status.Errorf(codes.Unimplemented, "method AddMachine not implemented")
}
func (UnimplementedClusterServer) mustEmbedUnimplementedClusterServer() {}
func (UnimplementedClusterServer) testEmbeddedByValue() {}
// UnsafeClusterServer may be embedded to opt out of forward compatibility for this service.
// Use of this interface is not recommended, as added methods to ClusterServer will
// result in compilation errors.
type UnsafeClusterServer interface {
mustEmbedUnimplementedClusterServer()
}
func RegisterClusterServer(s grpc.ServiceRegistrar, srv ClusterServer) {
// If the following call pancis, it indicates UnimplementedClusterServer was
// embedded by pointer and is nil. This will cause panics if an
// unimplemented method is ever invoked, so we test this at initialization
// time to prevent it from happening at runtime later due to I/O.
if t, ok := srv.(interface{ testEmbeddedByValue() }); ok {
t.testEmbeddedByValue()
}
s.RegisterService(&Cluster_ServiceDesc, srv)
}
func _Cluster_AddMachine_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) {
in := new(AddMachineRequest)
if err := dec(in); err != nil {
return nil, err
}
if interceptor == nil {
return srv.(ClusterServer).AddMachine(ctx, in)
}
info := &grpc.UnaryServerInfo{
Server: srv,
FullMethod: Cluster_AddMachine_FullMethodName,
}
handler := func(ctx context.Context, req interface{}) (interface{}, error) {
return srv.(ClusterServer).AddMachine(ctx, req.(*AddMachineRequest))
}
return interceptor(ctx, in, info, handler)
}
// Cluster_ServiceDesc is the grpc.ServiceDesc for Cluster service.
// It's only intended for direct use with grpc.RegisterService,
// and not to be introspected or modified (even as a copy)
var Cluster_ServiceDesc = grpc.ServiceDesc{
ServiceName: "cluster.Cluster",
HandlerType: (*ClusterServer)(nil),
Methods: []grpc.MethodDesc{
{
MethodName: "AddMachine",
Handler: _Cluster_AddMachine_Handler,
},
},
Streams: []grpc.StreamDesc{},
Metadata: "internal/machine/cluster/pb/cluster.proto",
}
+20
View File
@@ -0,0 +1,20 @@
package api
import (
"context"
pb2 "uncloud/internal/machine/api/pb"
)
// Server is the gRPC server for the Cluster service.
type Server struct {
pb2.UnimplementedClusterServer
}
func NewServer() *Server {
return &Server{}
}
// AddMachine adds a machine to the cluster.
func (s *Server) AddMachine(ctx context.Context, req *pb2.AddMachineRequest) (*pb2.AddMachineResponse, error) {
return &pb2.AddMachineResponse{}, nil
}
+44 -1
View File
@@ -3,10 +3,20 @@ package daemon
import (
"context"
"fmt"
"golang.org/x/sync/errgroup"
"google.golang.org/grpc"
"log/slog"
"net"
"uncloud/internal/machine"
"uncloud/internal/machine/api"
"uncloud/internal/machine/api/pb"
"uncloud/internal/machine/network"
)
const (
MachineAPIPort = 51000
)
func Run(ctx context.Context, dataDir string) error {
cfg, err := machine.ParseConfig(machine.ConfigPath(dataDir))
if err != nil {
@@ -30,5 +40,38 @@ func Run(ctx context.Context, dataDir string) error {
//}
//fmt.Println("Addresses:", addrs)
return wgnet.Run(ctx)
addr := fmt.Sprintf("127.0.0.1:%d", MachineAPIPort)
listener, err := net.Listen("tcp", addr)
if err != nil {
return fmt.Errorf("listen API port: %w", err)
}
grpcServer := grpc.NewServer()
pb.RegisterClusterServer(grpcServer, api.NewServer())
// Use an errgroup to coordinate error handling and graceful shutdown of multiple daemon components.
errGroup, ctx := errgroup.WithContext(ctx)
errGroup.Go(func() error {
slog.Info("Starting API server.", "addr", addr)
if sErr := grpcServer.Serve(listener); sErr != nil {
return fmt.Errorf("API server failed: %w", sErr)
}
return nil
})
errGroup.Go(func() error {
if err = wgnet.Run(ctx); err != nil {
return fmt.Errorf("WireGuard network failed: %w", err)
}
return nil
})
// Shutdown goroutine.
errGroup.Go(func() error {
<-ctx.Done()
slog.Info("Stopping API server.")
// TODO: implement timeout for graceful shutdown.
grpcServer.GracefulStop()
slog.Info("API server stopped.")
return nil
})
return errGroup.Wait()
}