From afe9c13283610564be8dad866799741a831f356e Mon Sep 17 00:00:00 2001 From: Pavel Sviderski Date: Thu, 29 Aug 2024 10:17:16 +1000 Subject: [PATCH] Implement grpc server in the daemon to receive cluster requests --- Makefile | 9 + go.mod | 4 +- go.sum | 4 + internal/machine/api/pb/cluster.pb.go | 468 +++++++++++++++++++++ internal/machine/api/pb/cluster.proto | 35 ++ internal/machine/api/pb/cluster_grpc.pb.go | 121 ++++++ internal/machine/api/server.go | 20 + internal/machine/daemon/daemon.go | 45 +- 8 files changed, 703 insertions(+), 3 deletions(-) create mode 100644 Makefile create mode 100644 internal/machine/api/pb/cluster.pb.go create mode 100644 internal/machine/api/pb/cluster.proto create mode 100644 internal/machine/api/pb/cluster_grpc.pb.go create mode 100644 internal/machine/api/server.go diff --git a/Makefile b/Makefile new file mode 100644 index 00000000..18695922 --- /dev/null +++ b/Makefile @@ -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 diff --git a/go.mod b/go.mod index f8d1850a..69007d64 100644 --- a/go.mod +++ b/go.mod @@ -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 diff --git a/go.sum b/go.sum index 8eb35943..7152e6b2 100644 --- a/go.sum +++ b/go.sum @@ -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= diff --git a/internal/machine/api/pb/cluster.pb.go b/internal/machine/api/pb/cluster.pb.go new file mode 100644 index 00000000..92b53f28 --- /dev/null +++ b/internal/machine/api/pb/cluster.pb.go @@ -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 +} diff --git a/internal/machine/api/pb/cluster.proto b/internal/machine/api/pb/cluster.proto new file mode 100644 index 00000000..208415da --- /dev/null +++ b/internal/machine/api/pb/cluster.proto @@ -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; +} diff --git a/internal/machine/api/pb/cluster_grpc.pb.go b/internal/machine/api/pb/cluster_grpc.pb.go new file mode 100644 index 00000000..4da5b6ee --- /dev/null +++ b/internal/machine/api/pb/cluster_grpc.pb.go @@ -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", +} diff --git a/internal/machine/api/server.go b/internal/machine/api/server.go new file mode 100644 index 00000000..b319ea2a --- /dev/null +++ b/internal/machine/api/server.go @@ -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 +} diff --git a/internal/machine/daemon/daemon.go b/internal/machine/daemon/daemon.go index 7806f9e9..09beb17a 100644 --- a/internal/machine/daemon/daemon.go +++ b/internal/machine/daemon/daemon.go @@ -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() }