diff --git a/persys-automation/internal/config/config.go b/persys-automation/internal/config/config.go index 77928f3..0c5133a 100644 --- a/persys-automation/internal/config/config.go +++ b/persys-automation/internal/config/config.go @@ -80,7 +80,7 @@ func Load() (*Config, error) { ), VaultPKIMount: envOr("AUTOMATION_VAULT_PKI_MOUNT", "pki"), VaultPKIRole: envOr("AUTOMATION_VAULT_PKI_ROLE", "persys-automation"), - VaultCertTTL: envDurationOr("AUTOMATION_VAULT_CERT_TTL", 1*time.Hour), + VaultCertTTL: envDurationOr("AUTOMATION_VAULT_CERT_TTL", 24*time.Hour), VaultServiceName: envOr("AUTOMATION_VAULT_SERVICE_NAME", "persys-automation"), VaultServiceDomain: strings.TrimSpace(os.Getenv("AUTOMATION_VAULT_SERVICE_DOMAIN")), VaultRetryInterval: envDurationOr("AUTOMATION_VAULT_RETRY_INTERVAL", 30*time.Second), diff --git a/persys-gateway/README.md b/persys-gateway/README.md index 54f1534..1139459 100644 --- a/persys-gateway/README.md +++ b/persys-gateway/README.md @@ -111,6 +111,12 @@ mTLS API: - `GET /clusters/:cluster_id/workloads` - `GET /clusters/:cluster_id/nodes` - `GET /clusters/:cluster_id/cluster/metrics` +- `GET/POST /clusters/:cluster_id/disks` — standalone block disks (AgentControl) +- `GET/DELETE /clusters/:cluster_id/disks/:id` +- `GET/POST /clusters/:cluster_id/buckets` — object storage (RGW via scheduler) +- `GET/DELETE /clusters/:cluster_id/buckets/:id` +- `GET /clusters/:cluster_id/buckets/:id/access` — Vault-backed S3 credentials +- `GET /clusters/:cluster_id/buckets/:id/objects` — list objects (prefix/pagination) - `POST /clusters/:cluster_id/forgery/projects/upsert` - `POST /clusters/:cluster_id/forgery/builds/trigger` - `POST /clusters/:cluster_id/forgery/webhooks/test` diff --git a/persys-gateway/cmd/main.go b/persys-gateway/cmd/main.go index 9041565..85beaa9 100755 --- a/persys-gateway/cmd/main.go +++ b/persys-gateway/cmd/main.go @@ -4,7 +4,6 @@ import ( "context" "crypto/tls" "crypto/x509" - "fmt" "log" "net/http" "net/http/pprof" @@ -171,19 +170,28 @@ func main() { app.authService = services.NewAuthService(app.db, ctx, jwtSecret) app.githubService = services.NewGithubService(cnf, webhookTLS) + if gs, ok := app.githubService.(interface{ SetCertManager(*certmanager.Manager) }); ok { + gs.SetCertManager(vaultCertManager) + } app.clusterControl = services.NewClusterControlService(cnf) + app.clusterControl.SetCertManager(vaultCertManager) app.clusterControl.Start(ctx) app.forgeryService = services.NewForgeryService(cnf, webhookTLS) + app.forgeryService.SetCertManager(vaultCertManager) app.webhookService, err = services.NewWebhookService(cnf, webhookTLS, app.db) if err != nil { log.Fatalf("failed to initialize webhook service: %v", err) } + if ws, ok := app.webhookService.(interface{ SetCertManager(*certmanager.Manager) }); ok { + ws.SetCertManager(vaultCertManager) + } app.webhookService.Start(ctx) app.automationService, err = services.NewAutomationService(cnf, webhookTLS) if err != nil { log.Fatalf("failed to initialize automation service: %v", err) } + app.automationService.SetCertManager(vaultCertManager) app.authController = controllers.NewAuthController( app.authService, ctx, app.githubService, app.db, @@ -207,12 +215,12 @@ func main() { mtlsRouter := gin.New() nonMTLSRouter := gin.New() - mtlsRouter.Use(gin.Logger()) + mtlsRouter.Use(middleware.AccessLogger()) mtlsRouter.Use(cors.New(corsConfig)) mtlsRouter.Use(gootelgin.Middleware("persys-gateway-mtls")) mtlsRouter.Use(middleware.ServiceIdentityHeader("persys-gateway")) - nonMTLSRouter.Use(gin.Logger()) + nonMTLSRouter.Use(middleware.AccessLogger()) nonMTLSRouter.Use(cors.New(corsConfig)) nonMTLSRouter.Use(gootelgin.Middleware("persys-gateway-public")) nonMTLSRouter.Use(middleware.ServiceIdentityHeader("persys-gateway")) @@ -318,7 +326,6 @@ func main() { automationGroup.Use(gwRouter.Resolve(catalog.AuthUser)) automationRouteController := routes.NewAutomationRouteController(app.automationController) automationRouteController.AutomationRoute(automationGroup) - // Webhook stays unauthenticated at the gateway level by design — its // own HMAC signature verification (X-Hub-Signature-256) IS its auth // mechanism, checked inside webhook.service.go. Mounted on the @@ -359,14 +366,27 @@ func main() { mtlsServer := &http.Server{Addr: cnf.App.HTTPAddr, Handler: mtlsRouter, TLSConfig: tlsConfig} nonMTLSServer := &http.Server{Addr: cnf.App.HTTPAddrPublic, Handler: nonMTLSRouter} + // goroutine, heap, allocs, block, mutex, threadcreate are served via pprof.Index + // through /debug/pprof/{profile-name} automatically once Index is registered debugMux := http.NewServeMux() debugMux.HandleFunc("/debug/pprof/", pprof.Index) debugMux.HandleFunc("/debug/pprof/cmdline", pprof.Cmdline) debugMux.HandleFunc("/debug/pprof/profile", pprof.Profile) debugMux.HandleFunc("/debug/pprof/symbol", pprof.Symbol) + debugMux.HandleFunc("/debug/pprof/trace", pprof.Trace) - // goroutine, heap, allocs, block, mutex, threadcreate are served via pprof.Index - // through /debug/pprof/{profile-name} automatically once Index is registered + // debugMux.HandleFunc("/debug/force-rotate", func(w http.ResponseWriter, r *http.Request) { + // if r.Method != http.MethodPost { + // http.Error(w, "POST only", http.StatusMethodNotAllowed) + // return + // } + // if err := vaultCertManager.ForceRotate(r.Context()); err != nil { + // http.Error(w, err.Error(), http.StatusInternalServerError) + // return + // } + // w.Header().Set("Content-Type", "application/json") + // _, _ = w.Write([]byte(`{"status":"rotated"}`)) + // }) debugServer := &http.Server{ Addr: "0.0.0.0:6060", @@ -414,17 +434,6 @@ func main() { } func buildMTLSClientConfig(cnf *config.Config) (*tls.Config, error) { - cert, err := tls.LoadX509KeyPair(cnf.TLS.CertPath, cnf.TLS.KeyPath) - if err != nil { - return nil, err - } - caCert, err := os.ReadFile(cnf.TLS.CAPath) - if err != nil { - return nil, err - } - caPool := x509.NewCertPool() - if !caPool.AppendCertsFromPEM(caCert) { - return nil, fmt.Errorf("invalid CA bundle") - } - return &tls.Config{Certificates: []tls.Certificate{cert}, RootCAs: caPool}, nil + // Live config: GetClientCertificate reloads from disk after certmanager rotates. + return services.LiveClientTLSConfig(cnf.TLS.CertPath, cnf.TLS.KeyPath, cnf.TLS.CAPath) } diff --git a/persys-gateway/config/config.go b/persys-gateway/config/config.go index 4ab839b..4118222 100755 --- a/persys-gateway/config/config.go +++ b/persys-gateway/config/config.go @@ -343,7 +343,7 @@ func (c *Config) applyDefaults() { c.Vault.ManagerAddr = "vault-manager:50069" } if c.Vault.CertTTL == time.Duration(0) { - c.Vault.CertTTL = 1 * time.Hour + c.Vault.CertTTL = 24 * time.Hour } if c.Vault.RetryInterval == time.Duration(0) { c.Vault.RetryInterval = 30 * time.Second diff --git a/persys-gateway/internal/controlv1/control.pb.go b/persys-gateway/internal/controlv1/control.pb.go index 8d4807d..583bcbf 100644 --- a/persys-gateway/internal/controlv1/control.pb.go +++ b/persys-gateway/internal/controlv1/control.pb.go @@ -4294,6 +4294,1502 @@ func (*ControlMessage_Apply) isControlMessage_Message() {} func (*ControlMessage_Delete) isControlMessage_Message() {} +type CreateDiskRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Name string `protobuf:"bytes,1,opt,name=name,proto3" json:"name,omitempty"` + Driver string `protobuf:"bytes,2,opt,name=driver,proto3" json:"driver,omitempty"` // local | ceph-rbd | nfs + SizeGb int64 `protobuf:"varint,3,opt,name=size_gb,json=sizeGb,proto3" json:"size_gb,omitempty"` + FsType string `protobuf:"bytes,4,opt,name=fs_type,json=fsType,proto3" json:"fs_type,omitempty"` + AccessMode string `protobuf:"bytes,5,opt,name=access_mode,json=accessMode,proto3" json:"access_mode,omitempty"` + RetainPolicy string `protobuf:"bytes,6,opt,name=retain_policy,json=retainPolicy,proto3" json:"retain_policy,omitempty"` // Delete | Retain + NodeId string `protobuf:"bytes,7,opt,name=node_id,json=nodeId,proto3" json:"node_id,omitempty"` // optional pre-pin for local + MountPath string `protobuf:"bytes,8,opt,name=mount_path,json=mountPath,proto3" json:"mount_path,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *CreateDiskRequest) Reset() { + *x = CreateDiskRequest{} + mi := &file_control_proto_msgTypes[60] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *CreateDiskRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*CreateDiskRequest) ProtoMessage() {} + +func (x *CreateDiskRequest) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[60] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use CreateDiskRequest.ProtoReflect.Descriptor instead. +func (*CreateDiskRequest) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{60} +} + +func (x *CreateDiskRequest) GetName() string { + if x != nil { + return x.Name + } + return "" +} + +func (x *CreateDiskRequest) GetDriver() string { + if x != nil { + return x.Driver + } + return "" +} + +func (x *CreateDiskRequest) GetSizeGb() int64 { + if x != nil { + return x.SizeGb + } + return 0 +} + +func (x *CreateDiskRequest) GetFsType() string { + if x != nil { + return x.FsType + } + return "" +} + +func (x *CreateDiskRequest) GetAccessMode() string { + if x != nil { + return x.AccessMode + } + return "" +} + +func (x *CreateDiskRequest) GetRetainPolicy() string { + if x != nil { + return x.RetainPolicy + } + return "" +} + +func (x *CreateDiskRequest) GetNodeId() string { + if x != nil { + return x.NodeId + } + return "" +} + +func (x *CreateDiskRequest) GetMountPath() string { + if x != nil { + return x.MountPath + } + return "" +} + +type CreateDiskResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Disk *DiskView `protobuf:"bytes,1,opt,name=disk,proto3" json:"disk,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *CreateDiskResponse) Reset() { + *x = CreateDiskResponse{} + mi := &file_control_proto_msgTypes[61] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *CreateDiskResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*CreateDiskResponse) ProtoMessage() {} + +func (x *CreateDiskResponse) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[61] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use CreateDiskResponse.ProtoReflect.Descriptor instead. +func (*CreateDiskResponse) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{61} +} + +func (x *CreateDiskResponse) GetDisk() *DiskView { + if x != nil { + return x.Disk + } + return nil +} + +type ListDisksRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ListDisksRequest) Reset() { + *x = ListDisksRequest{} + mi := &file_control_proto_msgTypes[62] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ListDisksRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ListDisksRequest) ProtoMessage() {} + +func (x *ListDisksRequest) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[62] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ListDisksRequest.ProtoReflect.Descriptor instead. +func (*ListDisksRequest) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{62} +} + +type ListDisksResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Disks []*DiskView `protobuf:"bytes,1,rep,name=disks,proto3" json:"disks,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ListDisksResponse) Reset() { + *x = ListDisksResponse{} + mi := &file_control_proto_msgTypes[63] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ListDisksResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ListDisksResponse) ProtoMessage() {} + +func (x *ListDisksResponse) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[63] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ListDisksResponse.ProtoReflect.Descriptor instead. +func (*ListDisksResponse) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{63} +} + +func (x *ListDisksResponse) GetDisks() []*DiskView { + if x != nil { + return x.Disks + } + return nil +} + +type GetDiskRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + DiskId string `protobuf:"bytes,1,opt,name=disk_id,json=diskId,proto3" json:"disk_id,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *GetDiskRequest) Reset() { + *x = GetDiskRequest{} + mi := &file_control_proto_msgTypes[64] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *GetDiskRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*GetDiskRequest) ProtoMessage() {} + +func (x *GetDiskRequest) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[64] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use GetDiskRequest.ProtoReflect.Descriptor instead. +func (*GetDiskRequest) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{64} +} + +func (x *GetDiskRequest) GetDiskId() string { + if x != nil { + return x.DiskId + } + return "" +} + +type GetDiskResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Disk *DiskView `protobuf:"bytes,1,opt,name=disk,proto3" json:"disk,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *GetDiskResponse) Reset() { + *x = GetDiskResponse{} + mi := &file_control_proto_msgTypes[65] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *GetDiskResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*GetDiskResponse) ProtoMessage() {} + +func (x *GetDiskResponse) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[65] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use GetDiskResponse.ProtoReflect.Descriptor instead. +func (*GetDiskResponse) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{65} +} + +func (x *GetDiskResponse) GetDisk() *DiskView { + if x != nil { + return x.Disk + } + return nil +} + +type DeleteDiskRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + DiskId string `protobuf:"bytes,1,opt,name=disk_id,json=diskId,proto3" json:"disk_id,omitempty"` + Force bool `protobuf:"varint,2,opt,name=force,proto3" json:"force,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *DeleteDiskRequest) Reset() { + *x = DeleteDiskRequest{} + mi := &file_control_proto_msgTypes[66] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *DeleteDiskRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*DeleteDiskRequest) ProtoMessage() {} + +func (x *DeleteDiskRequest) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[66] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use DeleteDiskRequest.ProtoReflect.Descriptor instead. +func (*DeleteDiskRequest) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{66} +} + +func (x *DeleteDiskRequest) GetDiskId() string { + if x != nil { + return x.DiskId + } + return "" +} + +func (x *DeleteDiskRequest) GetForce() bool { + if x != nil { + return x.Force + } + return false +} + +type DeleteDiskResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Success bool `protobuf:"varint,1,opt,name=success,proto3" json:"success,omitempty"` + ErrorMessage string `protobuf:"bytes,2,opt,name=error_message,json=errorMessage,proto3" json:"error_message,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *DeleteDiskResponse) Reset() { + *x = DeleteDiskResponse{} + mi := &file_control_proto_msgTypes[67] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *DeleteDiskResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*DeleteDiskResponse) ProtoMessage() {} + +func (x *DeleteDiskResponse) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[67] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use DeleteDiskResponse.ProtoReflect.Descriptor instead. +func (*DeleteDiskResponse) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{67} +} + +func (x *DeleteDiskResponse) GetSuccess() bool { + if x != nil { + return x.Success + } + return false +} + +func (x *DeleteDiskResponse) GetErrorMessage() string { + if x != nil { + return x.ErrorMessage + } + return "" +} + +type DiskView struct { + state protoimpl.MessageState `protogen:"open.v1"` + Id string `protobuf:"bytes,1,opt,name=id,proto3" json:"id,omitempty"` + Name string `protobuf:"bytes,2,opt,name=name,proto3" json:"name,omitempty"` + Driver string `protobuf:"bytes,3,opt,name=driver,proto3" json:"driver,omitempty"` + SizeGb int64 `protobuf:"varint,4,opt,name=size_gb,json=sizeGb,proto3" json:"size_gb,omitempty"` + FsType string `protobuf:"bytes,5,opt,name=fs_type,json=fsType,proto3" json:"fs_type,omitempty"` + AccessMode string `protobuf:"bytes,6,opt,name=access_mode,json=accessMode,proto3" json:"access_mode,omitempty"` + RetainPolicy string `protobuf:"bytes,7,opt,name=retain_policy,json=retainPolicy,proto3" json:"retain_policy,omitempty"` + Phase string `protobuf:"bytes,8,opt,name=phase,proto3" json:"phase,omitempty"` + LastError string `protobuf:"bytes,9,opt,name=last_error,json=lastError,proto3" json:"last_error,omitempty"` + NodeId string `protobuf:"bytes,10,opt,name=node_id,json=nodeId,proto3" json:"node_id,omitempty"` + Device string `protobuf:"bytes,11,opt,name=device,proto3" json:"device,omitempty"` + Standalone bool `protobuf:"varint,12,opt,name=standalone,proto3" json:"standalone,omitempty"` + MountPath string `protobuf:"bytes,13,opt,name=mount_path,json=mountPath,proto3" json:"mount_path,omitempty"` + WorkloadRefs []string `protobuf:"bytes,14,rep,name=workload_refs,json=workloadRefs,proto3" json:"workload_refs,omitempty"` + AttachedNodes []string `protobuf:"bytes,15,rep,name=attached_nodes,json=attachedNodes,proto3" json:"attached_nodes,omitempty"` + CreatedAt *timestamppb.Timestamp `protobuf:"bytes,16,opt,name=created_at,json=createdAt,proto3" json:"created_at,omitempty"` + UpdatedAt *timestamppb.Timestamp `protobuf:"bytes,17,opt,name=updated_at,json=updatedAt,proto3" json:"updated_at,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *DiskView) Reset() { + *x = DiskView{} + mi := &file_control_proto_msgTypes[68] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *DiskView) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*DiskView) ProtoMessage() {} + +func (x *DiskView) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[68] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use DiskView.ProtoReflect.Descriptor instead. +func (*DiskView) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{68} +} + +func (x *DiskView) GetId() string { + if x != nil { + return x.Id + } + return "" +} + +func (x *DiskView) GetName() string { + if x != nil { + return x.Name + } + return "" +} + +func (x *DiskView) GetDriver() string { + if x != nil { + return x.Driver + } + return "" +} + +func (x *DiskView) GetSizeGb() int64 { + if x != nil { + return x.SizeGb + } + return 0 +} + +func (x *DiskView) GetFsType() string { + if x != nil { + return x.FsType + } + return "" +} + +func (x *DiskView) GetAccessMode() string { + if x != nil { + return x.AccessMode + } + return "" +} + +func (x *DiskView) GetRetainPolicy() string { + if x != nil { + return x.RetainPolicy + } + return "" +} + +func (x *DiskView) GetPhase() string { + if x != nil { + return x.Phase + } + return "" +} + +func (x *DiskView) GetLastError() string { + if x != nil { + return x.LastError + } + return "" +} + +func (x *DiskView) GetNodeId() string { + if x != nil { + return x.NodeId + } + return "" +} + +func (x *DiskView) GetDevice() string { + if x != nil { + return x.Device + } + return "" +} + +func (x *DiskView) GetStandalone() bool { + if x != nil { + return x.Standalone + } + return false +} + +func (x *DiskView) GetMountPath() string { + if x != nil { + return x.MountPath + } + return "" +} + +func (x *DiskView) GetWorkloadRefs() []string { + if x != nil { + return x.WorkloadRefs + } + return nil +} + +func (x *DiskView) GetAttachedNodes() []string { + if x != nil { + return x.AttachedNodes + } + return nil +} + +func (x *DiskView) GetCreatedAt() *timestamppb.Timestamp { + if x != nil { + return x.CreatedAt + } + return nil +} + +func (x *DiskView) GetUpdatedAt() *timestamppb.Timestamp { + if x != nil { + return x.UpdatedAt + } + return nil +} + +type CreateBucketRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Name string `protobuf:"bytes,1,opt,name=name,proto3" json:"name,omitempty"` + Region string `protobuf:"bytes,2,opt,name=region,proto3" json:"region,omitempty"` + Versioning bool `protobuf:"varint,3,opt,name=versioning,proto3" json:"versioning,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *CreateBucketRequest) Reset() { + *x = CreateBucketRequest{} + mi := &file_control_proto_msgTypes[69] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *CreateBucketRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*CreateBucketRequest) ProtoMessage() {} + +func (x *CreateBucketRequest) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[69] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use CreateBucketRequest.ProtoReflect.Descriptor instead. +func (*CreateBucketRequest) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{69} +} + +func (x *CreateBucketRequest) GetName() string { + if x != nil { + return x.Name + } + return "" +} + +func (x *CreateBucketRequest) GetRegion() string { + if x != nil { + return x.Region + } + return "" +} + +func (x *CreateBucketRequest) GetVersioning() bool { + if x != nil { + return x.Versioning + } + return false +} + +type CreateBucketResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Bucket *BucketView `protobuf:"bytes,1,opt,name=bucket,proto3" json:"bucket,omitempty"` + Access *BucketAccess `protobuf:"bytes,2,opt,name=access,proto3" json:"access,omitempty"` // credentials returned once on create + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *CreateBucketResponse) Reset() { + *x = CreateBucketResponse{} + mi := &file_control_proto_msgTypes[70] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *CreateBucketResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*CreateBucketResponse) ProtoMessage() {} + +func (x *CreateBucketResponse) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[70] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use CreateBucketResponse.ProtoReflect.Descriptor instead. +func (*CreateBucketResponse) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{70} +} + +func (x *CreateBucketResponse) GetBucket() *BucketView { + if x != nil { + return x.Bucket + } + return nil +} + +func (x *CreateBucketResponse) GetAccess() *BucketAccess { + if x != nil { + return x.Access + } + return nil +} + +type ListBucketsRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ListBucketsRequest) Reset() { + *x = ListBucketsRequest{} + mi := &file_control_proto_msgTypes[71] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ListBucketsRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ListBucketsRequest) ProtoMessage() {} + +func (x *ListBucketsRequest) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[71] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ListBucketsRequest.ProtoReflect.Descriptor instead. +func (*ListBucketsRequest) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{71} +} + +type ListBucketsResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Buckets []*BucketView `protobuf:"bytes,1,rep,name=buckets,proto3" json:"buckets,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ListBucketsResponse) Reset() { + *x = ListBucketsResponse{} + mi := &file_control_proto_msgTypes[72] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ListBucketsResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ListBucketsResponse) ProtoMessage() {} + +func (x *ListBucketsResponse) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[72] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ListBucketsResponse.ProtoReflect.Descriptor instead. +func (*ListBucketsResponse) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{72} +} + +func (x *ListBucketsResponse) GetBuckets() []*BucketView { + if x != nil { + return x.Buckets + } + return nil +} + +type GetBucketRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + BucketId string `protobuf:"bytes,1,opt,name=bucket_id,json=bucketId,proto3" json:"bucket_id,omitempty"` // id or name + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *GetBucketRequest) Reset() { + *x = GetBucketRequest{} + mi := &file_control_proto_msgTypes[73] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *GetBucketRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*GetBucketRequest) ProtoMessage() {} + +func (x *GetBucketRequest) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[73] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use GetBucketRequest.ProtoReflect.Descriptor instead. +func (*GetBucketRequest) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{73} +} + +func (x *GetBucketRequest) GetBucketId() string { + if x != nil { + return x.BucketId + } + return "" +} + +type GetBucketResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Bucket *BucketView `protobuf:"bytes,1,opt,name=bucket,proto3" json:"bucket,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *GetBucketResponse) Reset() { + *x = GetBucketResponse{} + mi := &file_control_proto_msgTypes[74] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *GetBucketResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*GetBucketResponse) ProtoMessage() {} + +func (x *GetBucketResponse) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[74] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use GetBucketResponse.ProtoReflect.Descriptor instead. +func (*GetBucketResponse) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{74} +} + +func (x *GetBucketResponse) GetBucket() *BucketView { + if x != nil { + return x.Bucket + } + return nil +} + +type DeleteBucketRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + BucketId string `protobuf:"bytes,1,opt,name=bucket_id,json=bucketId,proto3" json:"bucket_id,omitempty"` + Force bool `protobuf:"varint,2,opt,name=force,proto3" json:"force,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *DeleteBucketRequest) Reset() { + *x = DeleteBucketRequest{} + mi := &file_control_proto_msgTypes[75] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *DeleteBucketRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*DeleteBucketRequest) ProtoMessage() {} + +func (x *DeleteBucketRequest) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[75] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use DeleteBucketRequest.ProtoReflect.Descriptor instead. +func (*DeleteBucketRequest) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{75} +} + +func (x *DeleteBucketRequest) GetBucketId() string { + if x != nil { + return x.BucketId + } + return "" +} + +func (x *DeleteBucketRequest) GetForce() bool { + if x != nil { + return x.Force + } + return false +} + +type DeleteBucketResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Success bool `protobuf:"varint,1,opt,name=success,proto3" json:"success,omitempty"` + ErrorMessage string `protobuf:"bytes,2,opt,name=error_message,json=errorMessage,proto3" json:"error_message,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *DeleteBucketResponse) Reset() { + *x = DeleteBucketResponse{} + mi := &file_control_proto_msgTypes[76] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *DeleteBucketResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*DeleteBucketResponse) ProtoMessage() {} + +func (x *DeleteBucketResponse) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[76] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use DeleteBucketResponse.ProtoReflect.Descriptor instead. +func (*DeleteBucketResponse) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{76} +} + +func (x *DeleteBucketResponse) GetSuccess() bool { + if x != nil { + return x.Success + } + return false +} + +func (x *DeleteBucketResponse) GetErrorMessage() string { + if x != nil { + return x.ErrorMessage + } + return "" +} + +type GetBucketAccessRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + BucketId string `protobuf:"bytes,1,opt,name=bucket_id,json=bucketId,proto3" json:"bucket_id,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *GetBucketAccessRequest) Reset() { + *x = GetBucketAccessRequest{} + mi := &file_control_proto_msgTypes[77] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *GetBucketAccessRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*GetBucketAccessRequest) ProtoMessage() {} + +func (x *GetBucketAccessRequest) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[77] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use GetBucketAccessRequest.ProtoReflect.Descriptor instead. +func (*GetBucketAccessRequest) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{77} +} + +func (x *GetBucketAccessRequest) GetBucketId() string { + if x != nil { + return x.BucketId + } + return "" +} + +type GetBucketAccessResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Access *BucketAccess `protobuf:"bytes,1,opt,name=access,proto3" json:"access,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *GetBucketAccessResponse) Reset() { + *x = GetBucketAccessResponse{} + mi := &file_control_proto_msgTypes[78] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *GetBucketAccessResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*GetBucketAccessResponse) ProtoMessage() {} + +func (x *GetBucketAccessResponse) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[78] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use GetBucketAccessResponse.ProtoReflect.Descriptor instead. +func (*GetBucketAccessResponse) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{78} +} + +func (x *GetBucketAccessResponse) GetAccess() *BucketAccess { + if x != nil { + return x.Access + } + return nil +} + +type ListBucketObjectsRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + BucketId string `protobuf:"bytes,1,opt,name=bucket_id,json=bucketId,proto3" json:"bucket_id,omitempty"` + Prefix string `protobuf:"bytes,2,opt,name=prefix,proto3" json:"prefix,omitempty"` + ContinuationToken string `protobuf:"bytes,3,opt,name=continuation_token,json=continuationToken,proto3" json:"continuation_token,omitempty"` + MaxKeys int32 `protobuf:"varint,4,opt,name=max_keys,json=maxKeys,proto3" json:"max_keys,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ListBucketObjectsRequest) Reset() { + *x = ListBucketObjectsRequest{} + mi := &file_control_proto_msgTypes[79] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ListBucketObjectsRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ListBucketObjectsRequest) ProtoMessage() {} + +func (x *ListBucketObjectsRequest) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[79] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ListBucketObjectsRequest.ProtoReflect.Descriptor instead. +func (*ListBucketObjectsRequest) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{79} +} + +func (x *ListBucketObjectsRequest) GetBucketId() string { + if x != nil { + return x.BucketId + } + return "" +} + +func (x *ListBucketObjectsRequest) GetPrefix() string { + if x != nil { + return x.Prefix + } + return "" +} + +func (x *ListBucketObjectsRequest) GetContinuationToken() string { + if x != nil { + return x.ContinuationToken + } + return "" +} + +func (x *ListBucketObjectsRequest) GetMaxKeys() int32 { + if x != nil { + return x.MaxKeys + } + return 0 +} + +type ListBucketObjectsResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Objects []*ObjectInfo `protobuf:"bytes,1,rep,name=objects,proto3" json:"objects,omitempty"` + NextContinuationToken string `protobuf:"bytes,2,opt,name=next_continuation_token,json=nextContinuationToken,proto3" json:"next_continuation_token,omitempty"` + IsTruncated bool `protobuf:"varint,3,opt,name=is_truncated,json=isTruncated,proto3" json:"is_truncated,omitempty"` + Prefix string `protobuf:"bytes,4,opt,name=prefix,proto3" json:"prefix,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ListBucketObjectsResponse) Reset() { + *x = ListBucketObjectsResponse{} + mi := &file_control_proto_msgTypes[80] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ListBucketObjectsResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ListBucketObjectsResponse) ProtoMessage() {} + +func (x *ListBucketObjectsResponse) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[80] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ListBucketObjectsResponse.ProtoReflect.Descriptor instead. +func (*ListBucketObjectsResponse) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{80} +} + +func (x *ListBucketObjectsResponse) GetObjects() []*ObjectInfo { + if x != nil { + return x.Objects + } + return nil +} + +func (x *ListBucketObjectsResponse) GetNextContinuationToken() string { + if x != nil { + return x.NextContinuationToken + } + return "" +} + +func (x *ListBucketObjectsResponse) GetIsTruncated() bool { + if x != nil { + return x.IsTruncated + } + return false +} + +func (x *ListBucketObjectsResponse) GetPrefix() string { + if x != nil { + return x.Prefix + } + return "" +} + +type BucketView struct { + state protoimpl.MessageState `protogen:"open.v1"` + Id string `protobuf:"bytes,1,opt,name=id,proto3" json:"id,omitempty"` + Name string `protobuf:"bytes,2,opt,name=name,proto3" json:"name,omitempty"` + Region string `protobuf:"bytes,3,opt,name=region,proto3" json:"region,omitempty"` + Owner string `protobuf:"bytes,4,opt,name=owner,proto3" json:"owner,omitempty"` + Endpoint string `protobuf:"bytes,5,opt,name=endpoint,proto3" json:"endpoint,omitempty"` + Versioning bool `protobuf:"varint,6,opt,name=versioning,proto3" json:"versioning,omitempty"` + ObjectCount int64 `protobuf:"varint,7,opt,name=object_count,json=objectCount,proto3" json:"object_count,omitempty"` + SizeBytes int64 `protobuf:"varint,8,opt,name=size_bytes,json=sizeBytes,proto3" json:"size_bytes,omitempty"` + Phase string `protobuf:"bytes,9,opt,name=phase,proto3" json:"phase,omitempty"` + LastError string `protobuf:"bytes,10,opt,name=last_error,json=lastError,proto3" json:"last_error,omitempty"` + CreatedAt *timestamppb.Timestamp `protobuf:"bytes,11,opt,name=created_at,json=createdAt,proto3" json:"created_at,omitempty"` + UpdatedAt *timestamppb.Timestamp `protobuf:"bytes,12,opt,name=updated_at,json=updatedAt,proto3" json:"updated_at,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *BucketView) Reset() { + *x = BucketView{} + mi := &file_control_proto_msgTypes[81] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *BucketView) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*BucketView) ProtoMessage() {} + +func (x *BucketView) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[81] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use BucketView.ProtoReflect.Descriptor instead. +func (*BucketView) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{81} +} + +func (x *BucketView) GetId() string { + if x != nil { + return x.Id + } + return "" +} + +func (x *BucketView) GetName() string { + if x != nil { + return x.Name + } + return "" +} + +func (x *BucketView) GetRegion() string { + if x != nil { + return x.Region + } + return "" +} + +func (x *BucketView) GetOwner() string { + if x != nil { + return x.Owner + } + return "" +} + +func (x *BucketView) GetEndpoint() string { + if x != nil { + return x.Endpoint + } + return "" +} + +func (x *BucketView) GetVersioning() bool { + if x != nil { + return x.Versioning + } + return false +} + +func (x *BucketView) GetObjectCount() int64 { + if x != nil { + return x.ObjectCount + } + return 0 +} + +func (x *BucketView) GetSizeBytes() int64 { + if x != nil { + return x.SizeBytes + } + return 0 +} + +func (x *BucketView) GetPhase() string { + if x != nil { + return x.Phase + } + return "" +} + +func (x *BucketView) GetLastError() string { + if x != nil { + return x.LastError + } + return "" +} + +func (x *BucketView) GetCreatedAt() *timestamppb.Timestamp { + if x != nil { + return x.CreatedAt + } + return nil +} + +func (x *BucketView) GetUpdatedAt() *timestamppb.Timestamp { + if x != nil { + return x.UpdatedAt + } + return nil +} + +type BucketAccess struct { + state protoimpl.MessageState `protogen:"open.v1"` + Endpoint string `protobuf:"bytes,1,opt,name=endpoint,proto3" json:"endpoint,omitempty"` + Region string `protobuf:"bytes,2,opt,name=region,proto3" json:"region,omitempty"` + Bucket string `protobuf:"bytes,3,opt,name=bucket,proto3" json:"bucket,omitempty"` + AccessKey string `protobuf:"bytes,4,opt,name=access_key,json=accessKey,proto3" json:"access_key,omitempty"` + SecretKey string `protobuf:"bytes,5,opt,name=secret_key,json=secretKey,proto3" json:"secret_key,omitempty"` + VaultPath string `protobuf:"bytes,6,opt,name=vault_path,json=vaultPath,proto3" json:"vault_path,omitempty"` // when secrets live in Vault + S3Url string `protobuf:"bytes,7,opt,name=s3_url,json=s3Url,proto3" json:"s3_url,omitempty"` // e.g. s3://bucket + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *BucketAccess) Reset() { + *x = BucketAccess{} + mi := &file_control_proto_msgTypes[82] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *BucketAccess) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*BucketAccess) ProtoMessage() {} + +func (x *BucketAccess) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[82] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use BucketAccess.ProtoReflect.Descriptor instead. +func (*BucketAccess) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{82} +} + +func (x *BucketAccess) GetEndpoint() string { + if x != nil { + return x.Endpoint + } + return "" +} + +func (x *BucketAccess) GetRegion() string { + if x != nil { + return x.Region + } + return "" +} + +func (x *BucketAccess) GetBucket() string { + if x != nil { + return x.Bucket + } + return "" +} + +func (x *BucketAccess) GetAccessKey() string { + if x != nil { + return x.AccessKey + } + return "" +} + +func (x *BucketAccess) GetSecretKey() string { + if x != nil { + return x.SecretKey + } + return "" +} + +func (x *BucketAccess) GetVaultPath() string { + if x != nil { + return x.VaultPath + } + return "" +} + +func (x *BucketAccess) GetS3Url() string { + if x != nil { + return x.S3Url + } + return "" +} + +type ObjectInfo struct { + state protoimpl.MessageState `protogen:"open.v1"` + Key string `protobuf:"bytes,1,opt,name=key,proto3" json:"key,omitempty"` + SizeBytes int64 `protobuf:"varint,2,opt,name=size_bytes,json=sizeBytes,proto3" json:"size_bytes,omitempty"` + Etag string `protobuf:"bytes,3,opt,name=etag,proto3" json:"etag,omitempty"` + LastModified string `protobuf:"bytes,4,opt,name=last_modified,json=lastModified,proto3" json:"last_modified,omitempty"` + StorageClass string `protobuf:"bytes,5,opt,name=storage_class,json=storageClass,proto3" json:"storage_class,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ObjectInfo) Reset() { + *x = ObjectInfo{} + mi := &file_control_proto_msgTypes[83] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ObjectInfo) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ObjectInfo) ProtoMessage() {} + +func (x *ObjectInfo) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[83] + if x != nil { + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + if ms.LoadMessageInfo() == nil { + ms.StoreMessageInfo(mi) + } + return ms + } + return mi.MessageOf(x) +} + +// Deprecated: Use ObjectInfo.ProtoReflect.Descriptor instead. +func (*ObjectInfo) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{83} +} + +func (x *ObjectInfo) GetKey() string { + if x != nil { + return x.Key + } + return "" +} + +func (x *ObjectInfo) GetSizeBytes() int64 { + if x != nil { + return x.SizeBytes + } + return 0 +} + +func (x *ObjectInfo) GetEtag() string { + if x != nil { + return x.Etag + } + return "" +} + +func (x *ObjectInfo) GetLastModified() string { + if x != nil { + return x.LastModified + } + return "" +} + +func (x *ObjectInfo) GetStorageClass() string { + if x != nil { + return x.StorageClass + } + return "" +} + var File_control_proto protoreflect.FileDescriptor const file_control_proto_rawDesc = "" + @@ -4658,7 +6154,135 @@ const file_control_proto_rawDesc = "" + "\theartbeat\x18\x02 \x01(\v2#.persys.control.v1.HeartbeatRequestH\x00R\theartbeat\x12?\n" + "\x05apply\x18\x03 \x01(\v2'.persys.control.v1.ApplyWorkloadRequestH\x00R\x05apply\x12B\n" + "\x06delete\x18\x04 \x01(\v2(.persys.control.v1.DeleteWorkloadRequestH\x00R\x06deleteB\t\n" + - "\amessage*\xda\x01\n" + + "\amessage\"\xef\x01\n" + + "\x11CreateDiskRequest\x12\x12\n" + + "\x04name\x18\x01 \x01(\tR\x04name\x12\x16\n" + + "\x06driver\x18\x02 \x01(\tR\x06driver\x12\x17\n" + + "\asize_gb\x18\x03 \x01(\x03R\x06sizeGb\x12\x17\n" + + "\afs_type\x18\x04 \x01(\tR\x06fsType\x12\x1f\n" + + "\vaccess_mode\x18\x05 \x01(\tR\n" + + "accessMode\x12#\n" + + "\rretain_policy\x18\x06 \x01(\tR\fretainPolicy\x12\x17\n" + + "\anode_id\x18\a \x01(\tR\x06nodeId\x12\x1d\n" + + "\n" + + "mount_path\x18\b \x01(\tR\tmountPath\"E\n" + + "\x12CreateDiskResponse\x12/\n" + + "\x04disk\x18\x01 \x01(\v2\x1b.persys.control.v1.DiskViewR\x04disk\"\x12\n" + + "\x10ListDisksRequest\"F\n" + + "\x11ListDisksResponse\x121\n" + + "\x05disks\x18\x01 \x03(\v2\x1b.persys.control.v1.DiskViewR\x05disks\")\n" + + "\x0eGetDiskRequest\x12\x17\n" + + "\adisk_id\x18\x01 \x01(\tR\x06diskId\"B\n" + + "\x0fGetDiskResponse\x12/\n" + + "\x04disk\x18\x01 \x01(\v2\x1b.persys.control.v1.DiskViewR\x04disk\"B\n" + + "\x11DeleteDiskRequest\x12\x17\n" + + "\adisk_id\x18\x01 \x01(\tR\x06diskId\x12\x14\n" + + "\x05force\x18\x02 \x01(\bR\x05force\"S\n" + + "\x12DeleteDiskResponse\x12\x18\n" + + "\asuccess\x18\x01 \x01(\bR\asuccess\x12#\n" + + "\rerror_message\x18\x02 \x01(\tR\ferrorMessage\"\xa5\x04\n" + + "\bDiskView\x12\x0e\n" + + "\x02id\x18\x01 \x01(\tR\x02id\x12\x12\n" + + "\x04name\x18\x02 \x01(\tR\x04name\x12\x16\n" + + "\x06driver\x18\x03 \x01(\tR\x06driver\x12\x17\n" + + "\asize_gb\x18\x04 \x01(\x03R\x06sizeGb\x12\x17\n" + + "\afs_type\x18\x05 \x01(\tR\x06fsType\x12\x1f\n" + + "\vaccess_mode\x18\x06 \x01(\tR\n" + + "accessMode\x12#\n" + + "\rretain_policy\x18\a \x01(\tR\fretainPolicy\x12\x14\n" + + "\x05phase\x18\b \x01(\tR\x05phase\x12\x1d\n" + + "\n" + + "last_error\x18\t \x01(\tR\tlastError\x12\x17\n" + + "\anode_id\x18\n" + + " \x01(\tR\x06nodeId\x12\x16\n" + + "\x06device\x18\v \x01(\tR\x06device\x12\x1e\n" + + "\n" + + "standalone\x18\f \x01(\bR\n" + + "standalone\x12\x1d\n" + + "\n" + + "mount_path\x18\r \x01(\tR\tmountPath\x12#\n" + + "\rworkload_refs\x18\x0e \x03(\tR\fworkloadRefs\x12%\n" + + "\x0eattached_nodes\x18\x0f \x03(\tR\rattachedNodes\x129\n" + + "\n" + + "created_at\x18\x10 \x01(\v2\x1a.google.protobuf.TimestampR\tcreatedAt\x129\n" + + "\n" + + "updated_at\x18\x11 \x01(\v2\x1a.google.protobuf.TimestampR\tupdatedAt\"a\n" + + "\x13CreateBucketRequest\x12\x12\n" + + "\x04name\x18\x01 \x01(\tR\x04name\x12\x16\n" + + "\x06region\x18\x02 \x01(\tR\x06region\x12\x1e\n" + + "\n" + + "versioning\x18\x03 \x01(\bR\n" + + "versioning\"\x86\x01\n" + + "\x14CreateBucketResponse\x125\n" + + "\x06bucket\x18\x01 \x01(\v2\x1d.persys.control.v1.BucketViewR\x06bucket\x127\n" + + "\x06access\x18\x02 \x01(\v2\x1f.persys.control.v1.BucketAccessR\x06access\"\x14\n" + + "\x12ListBucketsRequest\"N\n" + + "\x13ListBucketsResponse\x127\n" + + "\abuckets\x18\x01 \x03(\v2\x1d.persys.control.v1.BucketViewR\abuckets\"/\n" + + "\x10GetBucketRequest\x12\x1b\n" + + "\tbucket_id\x18\x01 \x01(\tR\bbucketId\"J\n" + + "\x11GetBucketResponse\x125\n" + + "\x06bucket\x18\x01 \x01(\v2\x1d.persys.control.v1.BucketViewR\x06bucket\"H\n" + + "\x13DeleteBucketRequest\x12\x1b\n" + + "\tbucket_id\x18\x01 \x01(\tR\bbucketId\x12\x14\n" + + "\x05force\x18\x02 \x01(\bR\x05force\"U\n" + + "\x14DeleteBucketResponse\x12\x18\n" + + "\asuccess\x18\x01 \x01(\bR\asuccess\x12#\n" + + "\rerror_message\x18\x02 \x01(\tR\ferrorMessage\"5\n" + + "\x16GetBucketAccessRequest\x12\x1b\n" + + "\tbucket_id\x18\x01 \x01(\tR\bbucketId\"R\n" + + "\x17GetBucketAccessResponse\x127\n" + + "\x06access\x18\x01 \x01(\v2\x1f.persys.control.v1.BucketAccessR\x06access\"\x99\x01\n" + + "\x18ListBucketObjectsRequest\x12\x1b\n" + + "\tbucket_id\x18\x01 \x01(\tR\bbucketId\x12\x16\n" + + "\x06prefix\x18\x02 \x01(\tR\x06prefix\x12-\n" + + "\x12continuation_token\x18\x03 \x01(\tR\x11continuationToken\x12\x19\n" + + "\bmax_keys\x18\x04 \x01(\x05R\amaxKeys\"\xc7\x01\n" + + "\x19ListBucketObjectsResponse\x127\n" + + "\aobjects\x18\x01 \x03(\v2\x1d.persys.control.v1.ObjectInfoR\aobjects\x126\n" + + "\x17next_continuation_token\x18\x02 \x01(\tR\x15nextContinuationToken\x12!\n" + + "\fis_truncated\x18\x03 \x01(\bR\visTruncated\x12\x16\n" + + "\x06prefix\x18\x04 \x01(\tR\x06prefix\"\x87\x03\n" + + "\n" + + "BucketView\x12\x0e\n" + + "\x02id\x18\x01 \x01(\tR\x02id\x12\x12\n" + + "\x04name\x18\x02 \x01(\tR\x04name\x12\x16\n" + + "\x06region\x18\x03 \x01(\tR\x06region\x12\x14\n" + + "\x05owner\x18\x04 \x01(\tR\x05owner\x12\x1a\n" + + "\bendpoint\x18\x05 \x01(\tR\bendpoint\x12\x1e\n" + + "\n" + + "versioning\x18\x06 \x01(\bR\n" + + "versioning\x12!\n" + + "\fobject_count\x18\a \x01(\x03R\vobjectCount\x12\x1d\n" + + "\n" + + "size_bytes\x18\b \x01(\x03R\tsizeBytes\x12\x14\n" + + "\x05phase\x18\t \x01(\tR\x05phase\x12\x1d\n" + + "\n" + + "last_error\x18\n" + + " \x01(\tR\tlastError\x129\n" + + "\n" + + "created_at\x18\v \x01(\v2\x1a.google.protobuf.TimestampR\tcreatedAt\x129\n" + + "\n" + + "updated_at\x18\f \x01(\v2\x1a.google.protobuf.TimestampR\tupdatedAt\"\xce\x01\n" + + "\fBucketAccess\x12\x1a\n" + + "\bendpoint\x18\x01 \x01(\tR\bendpoint\x12\x16\n" + + "\x06region\x18\x02 \x01(\tR\x06region\x12\x16\n" + + "\x06bucket\x18\x03 \x01(\tR\x06bucket\x12\x1d\n" + + "\n" + + "access_key\x18\x04 \x01(\tR\taccessKey\x12\x1d\n" + + "\n" + + "secret_key\x18\x05 \x01(\tR\tsecretKey\x12\x1d\n" + + "\n" + + "vault_path\x18\x06 \x01(\tR\tvaultPath\x12\x15\n" + + "\x06s3_url\x18\a \x01(\tR\x05s3Url\"\x9b\x01\n" + + "\n" + + "ObjectInfo\x12\x10\n" + + "\x03key\x18\x01 \x01(\tR\x03key\x12\x1d\n" + + "\n" + + "size_bytes\x18\x02 \x01(\x03R\tsizeBytes\x12\x12\n" + + "\x04etag\x18\x03 \x01(\tR\x04etag\x12#\n" + + "\rlast_modified\x18\x04 \x01(\tR\flastModified\x12#\n" + + "\rstorage_class\x18\x05 \x01(\tR\fstorageClass*\xda\x01\n" + "\x14AutomationActionType\x12&\n" + "\"AUTOMATION_ACTION_TYPE_UNSPECIFIED\x10\x00\x12'\n" + "#AUTOMATION_ACTION_SET_DESIRED_STATE\x10\x01\x12$\n" + @@ -4674,7 +6298,7 @@ const file_control_proto_rawDesc = "" + "\rRUNTIME_ERROR\x10\x05\x12\x11\n" + "\rNETWORK_ERROR\x10\x06\x12\x11\n" + "\rSTORAGE_ERROR\x10\a\x12\x12\n" + - "\x0eVM_BOOT_FAILED\x10\b2\xaa\x0f\n" + + "\x0eVM_BOOT_FAILED\x10\b2\xdc\x16\n" + "\fAgentControl\x12_\n" + "\fRegisterNode\x12&.persys.control.v1.RegisterNodeRequest\x1a'.persys.control.v1.RegisterNodeResponse\x12V\n" + "\tHeartbeat\x12#.persys.control.v1.HeartbeatRequest\x1a$.persys.control.v1.HeartbeatResponse\x12b\n" + @@ -4696,7 +6320,19 @@ const file_control_proto_rawDesc = "" + "\n" + "ListEvents\x12$.persys.control.v1.ListEventsRequest\x1a%.persys.control.v1.ListEventsResponse\x12]\n" + "\vWatchEvents\x12%.persys.control.v1.WatchEventsRequest\x1a%.persys.control.v1.SchedulerEventView0\x01\x12Y\n" + - "\rControlStream\x12!.persys.control.v1.ControlMessage\x1a!.persys.control.v1.ControlMessage(\x010\x01B7Z5github.com/persys-dev/persys/api/control/v1;controlv1b\x06proto3" + "\rControlStream\x12!.persys.control.v1.ControlMessage\x1a!.persys.control.v1.ControlMessage(\x010\x01\x12Y\n" + + "\n" + + "CreateDisk\x12$.persys.control.v1.CreateDiskRequest\x1a%.persys.control.v1.CreateDiskResponse\x12V\n" + + "\tListDisks\x12#.persys.control.v1.ListDisksRequest\x1a$.persys.control.v1.ListDisksResponse\x12P\n" + + "\aGetDisk\x12!.persys.control.v1.GetDiskRequest\x1a\".persys.control.v1.GetDiskResponse\x12Y\n" + + "\n" + + "DeleteDisk\x12$.persys.control.v1.DeleteDiskRequest\x1a%.persys.control.v1.DeleteDiskResponse\x12_\n" + + "\fCreateBucket\x12&.persys.control.v1.CreateBucketRequest\x1a'.persys.control.v1.CreateBucketResponse\x12\\\n" + + "\vListBuckets\x12%.persys.control.v1.ListBucketsRequest\x1a&.persys.control.v1.ListBucketsResponse\x12V\n" + + "\tGetBucket\x12#.persys.control.v1.GetBucketRequest\x1a$.persys.control.v1.GetBucketResponse\x12_\n" + + "\fDeleteBucket\x12&.persys.control.v1.DeleteBucketRequest\x1a'.persys.control.v1.DeleteBucketResponse\x12h\n" + + "\x0fGetBucketAccess\x12).persys.control.v1.GetBucketAccessRequest\x1a*.persys.control.v1.GetBucketAccessResponse\x12n\n" + + "\x11ListBucketObjects\x12+.persys.control.v1.ListBucketObjectsRequest\x1a,.persys.control.v1.ListBucketObjectsResponseB7Z5github.com/persys-dev/persys/api/control/v1;controlv1b\x06proto3" var ( file_control_proto_rawDescOnce sync.Once @@ -4711,7 +6347,7 @@ func file_control_proto_rawDescGZIP() []byte { } var file_control_proto_enumTypes = make([]protoimpl.EnumInfo, 2) -var file_control_proto_msgTypes = make([]protoimpl.MessageInfo, 66) +var file_control_proto_msgTypes = make([]protoimpl.MessageInfo, 90) var file_control_proto_goTypes = []any{ (AutomationActionType)(0), // 0: persys.control.v1.AutomationActionType (FailureReason)(0), // 1: persys.control.v1.FailureReason @@ -4775,125 +6411,182 @@ var file_control_proto_goTypes = []any{ (*GetClusterSummaryRequest)(nil), // 59: persys.control.v1.GetClusterSummaryRequest (*GetClusterSummaryResponse)(nil), // 60: persys.control.v1.GetClusterSummaryResponse (*ControlMessage)(nil), // 61: persys.control.v1.ControlMessage - nil, // 62: persys.control.v1.RegisterNodeRequest.LabelsEntry - nil, // 63: persys.control.v1.WorkloadSpec.MetadataEntry - nil, // 64: persys.control.v1.ContainerSpec.EnvEntry - nil, // 65: persys.control.v1.ComposeSpec.EnvEntry - nil, // 66: persys.control.v1.NodeView.LabelsEntry - nil, // 67: persys.control.v1.SchedulerEventView.DetailsEntry - (*timestamppb.Timestamp)(nil), // 68: google.protobuf.Timestamp + (*CreateDiskRequest)(nil), // 62: persys.control.v1.CreateDiskRequest + (*CreateDiskResponse)(nil), // 63: persys.control.v1.CreateDiskResponse + (*ListDisksRequest)(nil), // 64: persys.control.v1.ListDisksRequest + (*ListDisksResponse)(nil), // 65: persys.control.v1.ListDisksResponse + (*GetDiskRequest)(nil), // 66: persys.control.v1.GetDiskRequest + (*GetDiskResponse)(nil), // 67: persys.control.v1.GetDiskResponse + (*DeleteDiskRequest)(nil), // 68: persys.control.v1.DeleteDiskRequest + (*DeleteDiskResponse)(nil), // 69: persys.control.v1.DeleteDiskResponse + (*DiskView)(nil), // 70: persys.control.v1.DiskView + (*CreateBucketRequest)(nil), // 71: persys.control.v1.CreateBucketRequest + (*CreateBucketResponse)(nil), // 72: persys.control.v1.CreateBucketResponse + (*ListBucketsRequest)(nil), // 73: persys.control.v1.ListBucketsRequest + (*ListBucketsResponse)(nil), // 74: persys.control.v1.ListBucketsResponse + (*GetBucketRequest)(nil), // 75: persys.control.v1.GetBucketRequest + (*GetBucketResponse)(nil), // 76: persys.control.v1.GetBucketResponse + (*DeleteBucketRequest)(nil), // 77: persys.control.v1.DeleteBucketRequest + (*DeleteBucketResponse)(nil), // 78: persys.control.v1.DeleteBucketResponse + (*GetBucketAccessRequest)(nil), // 79: persys.control.v1.GetBucketAccessRequest + (*GetBucketAccessResponse)(nil), // 80: persys.control.v1.GetBucketAccessResponse + (*ListBucketObjectsRequest)(nil), // 81: persys.control.v1.ListBucketObjectsRequest + (*ListBucketObjectsResponse)(nil), // 82: persys.control.v1.ListBucketObjectsResponse + (*BucketView)(nil), // 83: persys.control.v1.BucketView + (*BucketAccess)(nil), // 84: persys.control.v1.BucketAccess + (*ObjectInfo)(nil), // 85: persys.control.v1.ObjectInfo + nil, // 86: persys.control.v1.RegisterNodeRequest.LabelsEntry + nil, // 87: persys.control.v1.WorkloadSpec.MetadataEntry + nil, // 88: persys.control.v1.ContainerSpec.EnvEntry + nil, // 89: persys.control.v1.ComposeSpec.EnvEntry + nil, // 90: persys.control.v1.NodeView.LabelsEntry + nil, // 91: persys.control.v1.SchedulerEventView.DetailsEntry + (*timestamppb.Timestamp)(nil), // 92: google.protobuf.Timestamp } var file_control_proto_depIdxs = []int32{ - 0, // 0: persys.control.v1.AutomationSuggestion.action_type:type_name -> persys.control.v1.AutomationActionType - 68, // 1: persys.control.v1.AutomationSuggestion.suggested_at:type_name -> google.protobuf.Timestamp - 2, // 2: persys.control.v1.SubmitAutomationSuggestionRequest.suggestion:type_name -> persys.control.v1.AutomationSuggestion - 68, // 3: persys.control.v1.SubmitAutomationSuggestionResponse.decided_at:type_name -> google.protobuf.Timestamp - 6, // 4: persys.control.v1.RegisterNodeRequest.capabilities:type_name -> persys.control.v1.NodeCapabilities - 62, // 5: persys.control.v1.RegisterNodeRequest.labels:type_name -> persys.control.v1.RegisterNodeRequest.LabelsEntry - 68, // 6: persys.control.v1.RegisterNodeRequest.timestamp:type_name -> google.protobuf.Timestamp - 7, // 7: persys.control.v1.NodeCapabilities.storage_pools:type_name -> persys.control.v1.StoragePool - 68, // 8: persys.control.v1.RegisterNodeResponse.lease_expires_at:type_name -> google.protobuf.Timestamp - 10, // 9: persys.control.v1.HeartbeatRequest.usage:type_name -> persys.control.v1.NodeUsage - 29, // 10: persys.control.v1.HeartbeatRequest.workload_statuses:type_name -> persys.control.v1.WorkloadStatus - 68, // 11: persys.control.v1.HeartbeatRequest.timestamp:type_name -> google.protobuf.Timestamp - 27, // 12: persys.control.v1.HeartbeatRequest.workload_usage:type_name -> persys.control.v1.WorkloadUsageSnapshot - 68, // 13: persys.control.v1.HeartbeatResponse.lease_expires_at:type_name -> google.protobuf.Timestamp - 16, // 14: persys.control.v1.ApplyWorkloadRequest.spec:type_name -> persys.control.v1.WorkloadSpec - 1, // 15: persys.control.v1.ApplyWorkloadResponse.failure_reason:type_name -> persys.control.v1.FailureReason - 17, // 16: persys.control.v1.WorkloadSpec.resources:type_name -> persys.control.v1.ResourceRequirements - 18, // 17: persys.control.v1.WorkloadSpec.container:type_name -> persys.control.v1.ContainerSpec - 21, // 18: persys.control.v1.WorkloadSpec.compose:type_name -> persys.control.v1.ComposeSpec - 22, // 19: persys.control.v1.WorkloadSpec.vm:type_name -> persys.control.v1.VMSpec - 63, // 20: persys.control.v1.WorkloadSpec.metadata:type_name -> persys.control.v1.WorkloadSpec.MetadataEntry - 64, // 21: persys.control.v1.ContainerSpec.env:type_name -> persys.control.v1.ContainerSpec.EnvEntry - 19, // 22: persys.control.v1.ContainerSpec.volumes:type_name -> persys.control.v1.VolumeMount - 20, // 23: persys.control.v1.ContainerSpec.ports:type_name -> persys.control.v1.Port - 26, // 24: persys.control.v1.ContainerSpec.managed_volumes:type_name -> persys.control.v1.ManagedVolumeSpec - 65, // 25: persys.control.v1.ComposeSpec.env:type_name -> persys.control.v1.ComposeSpec.EnvEntry - 23, // 26: persys.control.v1.VMSpec.disks:type_name -> persys.control.v1.DiskConfig - 24, // 27: persys.control.v1.VMSpec.networks:type_name -> persys.control.v1.NetworkConfig - 25, // 28: persys.control.v1.VMSpec.cloud_init:type_name -> persys.control.v1.CloudInitConfig - 26, // 29: persys.control.v1.VMSpec.managed_volumes:type_name -> persys.control.v1.ManagedVolumeSpec - 68, // 30: persys.control.v1.WorkloadUsageSnapshot.collected_at:type_name -> google.protobuf.Timestamp - 68, // 31: persys.control.v1.ReasonDetail.last_transition:type_name -> google.protobuf.Timestamp - 68, // 32: persys.control.v1.ReasonDetail.next_retry_at:type_name -> google.protobuf.Timestamp - 1, // 33: persys.control.v1.WorkloadStatus.failure_reason:type_name -> persys.control.v1.FailureReason - 68, // 34: persys.control.v1.WorkloadStatus.last_transition:type_name -> google.protobuf.Timestamp - 28, // 35: persys.control.v1.WorkloadStatus.reason:type_name -> persys.control.v1.ReasonDetail - 27, // 36: persys.control.v1.WorkloadStatus.usage:type_name -> persys.control.v1.WorkloadUsageSnapshot - 49, // 37: persys.control.v1.DrainNodeResponse.node:type_name -> persys.control.v1.NodeView - 49, // 38: persys.control.v1.UndrainNodeResponse.node:type_name -> persys.control.v1.NodeView - 36, // 39: persys.control.v1.TaintNodeRequest.taint:type_name -> persys.control.v1.NodeTaint - 49, // 40: persys.control.v1.TaintNodeResponse.node:type_name -> persys.control.v1.NodeView - 49, // 41: persys.control.v1.UntaintNodeResponse.node:type_name -> persys.control.v1.NodeView - 49, // 42: persys.control.v1.SetNodeLabelResponse.node:type_name -> persys.control.v1.NodeView - 49, // 43: persys.control.v1.DeleteNodeLabelResponse.node:type_name -> persys.control.v1.NodeView - 49, // 44: persys.control.v1.ListNodesResponse.nodes:type_name -> persys.control.v1.NodeView - 49, // 45: persys.control.v1.GetNodeResponse.node:type_name -> persys.control.v1.NodeView - 68, // 46: persys.control.v1.NodeView.status_updated_at:type_name -> google.protobuf.Timestamp - 68, // 47: persys.control.v1.NodeView.last_heartbeat:type_name -> google.protobuf.Timestamp - 66, // 48: persys.control.v1.NodeView.labels:type_name -> persys.control.v1.NodeView.LabelsEntry - 36, // 49: persys.control.v1.NodeView.taints:type_name -> persys.control.v1.NodeTaint - 68, // 50: persys.control.v1.SchedulerEventView.timestamp:type_name -> google.protobuf.Timestamp - 67, // 51: persys.control.v1.SchedulerEventView.details:type_name -> persys.control.v1.SchedulerEventView.DetailsEntry - 51, // 52: persys.control.v1.ListEventsResponse.events:type_name -> persys.control.v1.SchedulerEventView - 58, // 53: persys.control.v1.ListWorkloadsResponse.workloads:type_name -> persys.control.v1.WorkloadView - 58, // 54: persys.control.v1.GetWorkloadResponse.workload:type_name -> persys.control.v1.WorkloadView - 68, // 55: persys.control.v1.WorkloadView.retry_next_at:type_name -> google.protobuf.Timestamp - 68, // 56: persys.control.v1.WorkloadView.last_updated:type_name -> google.protobuf.Timestamp - 28, // 57: persys.control.v1.WorkloadView.reason:type_name -> persys.control.v1.ReasonDetail - 27, // 58: persys.control.v1.WorkloadView.usage:type_name -> persys.control.v1.WorkloadUsageSnapshot - 68, // 59: persys.control.v1.WorkloadView.created_at:type_name -> google.protobuf.Timestamp - 68, // 60: persys.control.v1.GetClusterSummaryResponse.generated_at:type_name -> google.protobuf.Timestamp - 5, // 61: persys.control.v1.ControlMessage.register:type_name -> persys.control.v1.RegisterNodeRequest - 9, // 62: persys.control.v1.ControlMessage.heartbeat:type_name -> persys.control.v1.HeartbeatRequest - 12, // 63: persys.control.v1.ControlMessage.apply:type_name -> persys.control.v1.ApplyWorkloadRequest - 14, // 64: persys.control.v1.ControlMessage.delete:type_name -> persys.control.v1.DeleteWorkloadRequest - 5, // 65: persys.control.v1.AgentControl.RegisterNode:input_type -> persys.control.v1.RegisterNodeRequest - 9, // 66: persys.control.v1.AgentControl.Heartbeat:input_type -> persys.control.v1.HeartbeatRequest - 12, // 67: persys.control.v1.AgentControl.ApplyWorkload:input_type -> persys.control.v1.ApplyWorkloadRequest - 14, // 68: persys.control.v1.AgentControl.DeleteWorkload:input_type -> persys.control.v1.DeleteWorkloadRequest - 30, // 69: persys.control.v1.AgentControl.RetryWorkload:input_type -> persys.control.v1.RetryWorkloadRequest - 32, // 70: persys.control.v1.AgentControl.DrainNode:input_type -> persys.control.v1.DrainNodeRequest - 34, // 71: persys.control.v1.AgentControl.UndrainNode:input_type -> persys.control.v1.UndrainNodeRequest - 37, // 72: persys.control.v1.AgentControl.TaintNode:input_type -> persys.control.v1.TaintNodeRequest - 39, // 73: persys.control.v1.AgentControl.UntaintNode:input_type -> persys.control.v1.UntaintNodeRequest - 41, // 74: persys.control.v1.AgentControl.SetNodeLabel:input_type -> persys.control.v1.SetNodeLabelRequest - 43, // 75: persys.control.v1.AgentControl.DeleteNodeLabel:input_type -> persys.control.v1.DeleteNodeLabelRequest - 3, // 76: persys.control.v1.AgentControl.SubmitAutomationSuggestion:input_type -> persys.control.v1.SubmitAutomationSuggestionRequest - 45, // 77: persys.control.v1.AgentControl.ListNodes:input_type -> persys.control.v1.ListNodesRequest - 46, // 78: persys.control.v1.AgentControl.GetNode:input_type -> persys.control.v1.GetNodeRequest - 50, // 79: persys.control.v1.AgentControl.ListWorkloads:input_type -> persys.control.v1.ListWorkloadsRequest - 55, // 80: persys.control.v1.AgentControl.GetWorkload:input_type -> persys.control.v1.GetWorkloadRequest - 59, // 81: persys.control.v1.AgentControl.GetClusterSummary:input_type -> persys.control.v1.GetClusterSummaryRequest - 52, // 82: persys.control.v1.AgentControl.ListEvents:input_type -> persys.control.v1.ListEventsRequest - 54, // 83: persys.control.v1.AgentControl.WatchEvents:input_type -> persys.control.v1.WatchEventsRequest - 61, // 84: persys.control.v1.AgentControl.ControlStream:input_type -> persys.control.v1.ControlMessage - 8, // 85: persys.control.v1.AgentControl.RegisterNode:output_type -> persys.control.v1.RegisterNodeResponse - 11, // 86: persys.control.v1.AgentControl.Heartbeat:output_type -> persys.control.v1.HeartbeatResponse - 13, // 87: persys.control.v1.AgentControl.ApplyWorkload:output_type -> persys.control.v1.ApplyWorkloadResponse - 15, // 88: persys.control.v1.AgentControl.DeleteWorkload:output_type -> persys.control.v1.DeleteWorkloadResponse - 31, // 89: persys.control.v1.AgentControl.RetryWorkload:output_type -> persys.control.v1.RetryWorkloadResponse - 33, // 90: persys.control.v1.AgentControl.DrainNode:output_type -> persys.control.v1.DrainNodeResponse - 35, // 91: persys.control.v1.AgentControl.UndrainNode:output_type -> persys.control.v1.UndrainNodeResponse - 38, // 92: persys.control.v1.AgentControl.TaintNode:output_type -> persys.control.v1.TaintNodeResponse - 40, // 93: persys.control.v1.AgentControl.UntaintNode:output_type -> persys.control.v1.UntaintNodeResponse - 42, // 94: persys.control.v1.AgentControl.SetNodeLabel:output_type -> persys.control.v1.SetNodeLabelResponse - 44, // 95: persys.control.v1.AgentControl.DeleteNodeLabel:output_type -> persys.control.v1.DeleteNodeLabelResponse - 4, // 96: persys.control.v1.AgentControl.SubmitAutomationSuggestion:output_type -> persys.control.v1.SubmitAutomationSuggestionResponse - 47, // 97: persys.control.v1.AgentControl.ListNodes:output_type -> persys.control.v1.ListNodesResponse - 48, // 98: persys.control.v1.AgentControl.GetNode:output_type -> persys.control.v1.GetNodeResponse - 56, // 99: persys.control.v1.AgentControl.ListWorkloads:output_type -> persys.control.v1.ListWorkloadsResponse - 57, // 100: persys.control.v1.AgentControl.GetWorkload:output_type -> persys.control.v1.GetWorkloadResponse - 60, // 101: persys.control.v1.AgentControl.GetClusterSummary:output_type -> persys.control.v1.GetClusterSummaryResponse - 53, // 102: persys.control.v1.AgentControl.ListEvents:output_type -> persys.control.v1.ListEventsResponse - 51, // 103: persys.control.v1.AgentControl.WatchEvents:output_type -> persys.control.v1.SchedulerEventView - 61, // 104: persys.control.v1.AgentControl.ControlStream:output_type -> persys.control.v1.ControlMessage - 85, // [85:105] is the sub-list for method output_type - 65, // [65:85] is the sub-list for method input_type - 65, // [65:65] is the sub-list for extension type_name - 65, // [65:65] is the sub-list for extension extendee - 0, // [0:65] is the sub-list for field type_name + 0, // 0: persys.control.v1.AutomationSuggestion.action_type:type_name -> persys.control.v1.AutomationActionType + 92, // 1: persys.control.v1.AutomationSuggestion.suggested_at:type_name -> google.protobuf.Timestamp + 2, // 2: persys.control.v1.SubmitAutomationSuggestionRequest.suggestion:type_name -> persys.control.v1.AutomationSuggestion + 92, // 3: persys.control.v1.SubmitAutomationSuggestionResponse.decided_at:type_name -> google.protobuf.Timestamp + 6, // 4: persys.control.v1.RegisterNodeRequest.capabilities:type_name -> persys.control.v1.NodeCapabilities + 86, // 5: persys.control.v1.RegisterNodeRequest.labels:type_name -> persys.control.v1.RegisterNodeRequest.LabelsEntry + 92, // 6: persys.control.v1.RegisterNodeRequest.timestamp:type_name -> google.protobuf.Timestamp + 7, // 7: persys.control.v1.NodeCapabilities.storage_pools:type_name -> persys.control.v1.StoragePool + 92, // 8: persys.control.v1.RegisterNodeResponse.lease_expires_at:type_name -> google.protobuf.Timestamp + 10, // 9: persys.control.v1.HeartbeatRequest.usage:type_name -> persys.control.v1.NodeUsage + 29, // 10: persys.control.v1.HeartbeatRequest.workload_statuses:type_name -> persys.control.v1.WorkloadStatus + 92, // 11: persys.control.v1.HeartbeatRequest.timestamp:type_name -> google.protobuf.Timestamp + 27, // 12: persys.control.v1.HeartbeatRequest.workload_usage:type_name -> persys.control.v1.WorkloadUsageSnapshot + 92, // 13: persys.control.v1.HeartbeatResponse.lease_expires_at:type_name -> google.protobuf.Timestamp + 16, // 14: persys.control.v1.ApplyWorkloadRequest.spec:type_name -> persys.control.v1.WorkloadSpec + 1, // 15: persys.control.v1.ApplyWorkloadResponse.failure_reason:type_name -> persys.control.v1.FailureReason + 17, // 16: persys.control.v1.WorkloadSpec.resources:type_name -> persys.control.v1.ResourceRequirements + 18, // 17: persys.control.v1.WorkloadSpec.container:type_name -> persys.control.v1.ContainerSpec + 21, // 18: persys.control.v1.WorkloadSpec.compose:type_name -> persys.control.v1.ComposeSpec + 22, // 19: persys.control.v1.WorkloadSpec.vm:type_name -> persys.control.v1.VMSpec + 87, // 20: persys.control.v1.WorkloadSpec.metadata:type_name -> persys.control.v1.WorkloadSpec.MetadataEntry + 88, // 21: persys.control.v1.ContainerSpec.env:type_name -> persys.control.v1.ContainerSpec.EnvEntry + 19, // 22: persys.control.v1.ContainerSpec.volumes:type_name -> persys.control.v1.VolumeMount + 20, // 23: persys.control.v1.ContainerSpec.ports:type_name -> persys.control.v1.Port + 26, // 24: persys.control.v1.ContainerSpec.managed_volumes:type_name -> persys.control.v1.ManagedVolumeSpec + 89, // 25: persys.control.v1.ComposeSpec.env:type_name -> persys.control.v1.ComposeSpec.EnvEntry + 23, // 26: persys.control.v1.VMSpec.disks:type_name -> persys.control.v1.DiskConfig + 24, // 27: persys.control.v1.VMSpec.networks:type_name -> persys.control.v1.NetworkConfig + 25, // 28: persys.control.v1.VMSpec.cloud_init:type_name -> persys.control.v1.CloudInitConfig + 26, // 29: persys.control.v1.VMSpec.managed_volumes:type_name -> persys.control.v1.ManagedVolumeSpec + 92, // 30: persys.control.v1.WorkloadUsageSnapshot.collected_at:type_name -> google.protobuf.Timestamp + 92, // 31: persys.control.v1.ReasonDetail.last_transition:type_name -> google.protobuf.Timestamp + 92, // 32: persys.control.v1.ReasonDetail.next_retry_at:type_name -> google.protobuf.Timestamp + 1, // 33: persys.control.v1.WorkloadStatus.failure_reason:type_name -> persys.control.v1.FailureReason + 92, // 34: persys.control.v1.WorkloadStatus.last_transition:type_name -> google.protobuf.Timestamp + 28, // 35: persys.control.v1.WorkloadStatus.reason:type_name -> persys.control.v1.ReasonDetail + 27, // 36: persys.control.v1.WorkloadStatus.usage:type_name -> persys.control.v1.WorkloadUsageSnapshot + 49, // 37: persys.control.v1.DrainNodeResponse.node:type_name -> persys.control.v1.NodeView + 49, // 38: persys.control.v1.UndrainNodeResponse.node:type_name -> persys.control.v1.NodeView + 36, // 39: persys.control.v1.TaintNodeRequest.taint:type_name -> persys.control.v1.NodeTaint + 49, // 40: persys.control.v1.TaintNodeResponse.node:type_name -> persys.control.v1.NodeView + 49, // 41: persys.control.v1.UntaintNodeResponse.node:type_name -> persys.control.v1.NodeView + 49, // 42: persys.control.v1.SetNodeLabelResponse.node:type_name -> persys.control.v1.NodeView + 49, // 43: persys.control.v1.DeleteNodeLabelResponse.node:type_name -> persys.control.v1.NodeView + 49, // 44: persys.control.v1.ListNodesResponse.nodes:type_name -> persys.control.v1.NodeView + 49, // 45: persys.control.v1.GetNodeResponse.node:type_name -> persys.control.v1.NodeView + 92, // 46: persys.control.v1.NodeView.status_updated_at:type_name -> google.protobuf.Timestamp + 92, // 47: persys.control.v1.NodeView.last_heartbeat:type_name -> google.protobuf.Timestamp + 90, // 48: persys.control.v1.NodeView.labels:type_name -> persys.control.v1.NodeView.LabelsEntry + 36, // 49: persys.control.v1.NodeView.taints:type_name -> persys.control.v1.NodeTaint + 92, // 50: persys.control.v1.SchedulerEventView.timestamp:type_name -> google.protobuf.Timestamp + 91, // 51: persys.control.v1.SchedulerEventView.details:type_name -> persys.control.v1.SchedulerEventView.DetailsEntry + 51, // 52: persys.control.v1.ListEventsResponse.events:type_name -> persys.control.v1.SchedulerEventView + 58, // 53: persys.control.v1.ListWorkloadsResponse.workloads:type_name -> persys.control.v1.WorkloadView + 58, // 54: persys.control.v1.GetWorkloadResponse.workload:type_name -> persys.control.v1.WorkloadView + 92, // 55: persys.control.v1.WorkloadView.retry_next_at:type_name -> google.protobuf.Timestamp + 92, // 56: persys.control.v1.WorkloadView.last_updated:type_name -> google.protobuf.Timestamp + 28, // 57: persys.control.v1.WorkloadView.reason:type_name -> persys.control.v1.ReasonDetail + 27, // 58: persys.control.v1.WorkloadView.usage:type_name -> persys.control.v1.WorkloadUsageSnapshot + 92, // 59: persys.control.v1.WorkloadView.created_at:type_name -> google.protobuf.Timestamp + 92, // 60: persys.control.v1.GetClusterSummaryResponse.generated_at:type_name -> google.protobuf.Timestamp + 5, // 61: persys.control.v1.ControlMessage.register:type_name -> persys.control.v1.RegisterNodeRequest + 9, // 62: persys.control.v1.ControlMessage.heartbeat:type_name -> persys.control.v1.HeartbeatRequest + 12, // 63: persys.control.v1.ControlMessage.apply:type_name -> persys.control.v1.ApplyWorkloadRequest + 14, // 64: persys.control.v1.ControlMessage.delete:type_name -> persys.control.v1.DeleteWorkloadRequest + 70, // 65: persys.control.v1.CreateDiskResponse.disk:type_name -> persys.control.v1.DiskView + 70, // 66: persys.control.v1.ListDisksResponse.disks:type_name -> persys.control.v1.DiskView + 70, // 67: persys.control.v1.GetDiskResponse.disk:type_name -> persys.control.v1.DiskView + 92, // 68: persys.control.v1.DiskView.created_at:type_name -> google.protobuf.Timestamp + 92, // 69: persys.control.v1.DiskView.updated_at:type_name -> google.protobuf.Timestamp + 83, // 70: persys.control.v1.CreateBucketResponse.bucket:type_name -> persys.control.v1.BucketView + 84, // 71: persys.control.v1.CreateBucketResponse.access:type_name -> persys.control.v1.BucketAccess + 83, // 72: persys.control.v1.ListBucketsResponse.buckets:type_name -> persys.control.v1.BucketView + 83, // 73: persys.control.v1.GetBucketResponse.bucket:type_name -> persys.control.v1.BucketView + 84, // 74: persys.control.v1.GetBucketAccessResponse.access:type_name -> persys.control.v1.BucketAccess + 85, // 75: persys.control.v1.ListBucketObjectsResponse.objects:type_name -> persys.control.v1.ObjectInfo + 92, // 76: persys.control.v1.BucketView.created_at:type_name -> google.protobuf.Timestamp + 92, // 77: persys.control.v1.BucketView.updated_at:type_name -> google.protobuf.Timestamp + 5, // 78: persys.control.v1.AgentControl.RegisterNode:input_type -> persys.control.v1.RegisterNodeRequest + 9, // 79: persys.control.v1.AgentControl.Heartbeat:input_type -> persys.control.v1.HeartbeatRequest + 12, // 80: persys.control.v1.AgentControl.ApplyWorkload:input_type -> persys.control.v1.ApplyWorkloadRequest + 14, // 81: persys.control.v1.AgentControl.DeleteWorkload:input_type -> persys.control.v1.DeleteWorkloadRequest + 30, // 82: persys.control.v1.AgentControl.RetryWorkload:input_type -> persys.control.v1.RetryWorkloadRequest + 32, // 83: persys.control.v1.AgentControl.DrainNode:input_type -> persys.control.v1.DrainNodeRequest + 34, // 84: persys.control.v1.AgentControl.UndrainNode:input_type -> persys.control.v1.UndrainNodeRequest + 37, // 85: persys.control.v1.AgentControl.TaintNode:input_type -> persys.control.v1.TaintNodeRequest + 39, // 86: persys.control.v1.AgentControl.UntaintNode:input_type -> persys.control.v1.UntaintNodeRequest + 41, // 87: persys.control.v1.AgentControl.SetNodeLabel:input_type -> persys.control.v1.SetNodeLabelRequest + 43, // 88: persys.control.v1.AgentControl.DeleteNodeLabel:input_type -> persys.control.v1.DeleteNodeLabelRequest + 3, // 89: persys.control.v1.AgentControl.SubmitAutomationSuggestion:input_type -> persys.control.v1.SubmitAutomationSuggestionRequest + 45, // 90: persys.control.v1.AgentControl.ListNodes:input_type -> persys.control.v1.ListNodesRequest + 46, // 91: persys.control.v1.AgentControl.GetNode:input_type -> persys.control.v1.GetNodeRequest + 50, // 92: persys.control.v1.AgentControl.ListWorkloads:input_type -> persys.control.v1.ListWorkloadsRequest + 55, // 93: persys.control.v1.AgentControl.GetWorkload:input_type -> persys.control.v1.GetWorkloadRequest + 59, // 94: persys.control.v1.AgentControl.GetClusterSummary:input_type -> persys.control.v1.GetClusterSummaryRequest + 52, // 95: persys.control.v1.AgentControl.ListEvents:input_type -> persys.control.v1.ListEventsRequest + 54, // 96: persys.control.v1.AgentControl.WatchEvents:input_type -> persys.control.v1.WatchEventsRequest + 61, // 97: persys.control.v1.AgentControl.ControlStream:input_type -> persys.control.v1.ControlMessage + 62, // 98: persys.control.v1.AgentControl.CreateDisk:input_type -> persys.control.v1.CreateDiskRequest + 64, // 99: persys.control.v1.AgentControl.ListDisks:input_type -> persys.control.v1.ListDisksRequest + 66, // 100: persys.control.v1.AgentControl.GetDisk:input_type -> persys.control.v1.GetDiskRequest + 68, // 101: persys.control.v1.AgentControl.DeleteDisk:input_type -> persys.control.v1.DeleteDiskRequest + 71, // 102: persys.control.v1.AgentControl.CreateBucket:input_type -> persys.control.v1.CreateBucketRequest + 73, // 103: persys.control.v1.AgentControl.ListBuckets:input_type -> persys.control.v1.ListBucketsRequest + 75, // 104: persys.control.v1.AgentControl.GetBucket:input_type -> persys.control.v1.GetBucketRequest + 77, // 105: persys.control.v1.AgentControl.DeleteBucket:input_type -> persys.control.v1.DeleteBucketRequest + 79, // 106: persys.control.v1.AgentControl.GetBucketAccess:input_type -> persys.control.v1.GetBucketAccessRequest + 81, // 107: persys.control.v1.AgentControl.ListBucketObjects:input_type -> persys.control.v1.ListBucketObjectsRequest + 8, // 108: persys.control.v1.AgentControl.RegisterNode:output_type -> persys.control.v1.RegisterNodeResponse + 11, // 109: persys.control.v1.AgentControl.Heartbeat:output_type -> persys.control.v1.HeartbeatResponse + 13, // 110: persys.control.v1.AgentControl.ApplyWorkload:output_type -> persys.control.v1.ApplyWorkloadResponse + 15, // 111: persys.control.v1.AgentControl.DeleteWorkload:output_type -> persys.control.v1.DeleteWorkloadResponse + 31, // 112: persys.control.v1.AgentControl.RetryWorkload:output_type -> persys.control.v1.RetryWorkloadResponse + 33, // 113: persys.control.v1.AgentControl.DrainNode:output_type -> persys.control.v1.DrainNodeResponse + 35, // 114: persys.control.v1.AgentControl.UndrainNode:output_type -> persys.control.v1.UndrainNodeResponse + 38, // 115: persys.control.v1.AgentControl.TaintNode:output_type -> persys.control.v1.TaintNodeResponse + 40, // 116: persys.control.v1.AgentControl.UntaintNode:output_type -> persys.control.v1.UntaintNodeResponse + 42, // 117: persys.control.v1.AgentControl.SetNodeLabel:output_type -> persys.control.v1.SetNodeLabelResponse + 44, // 118: persys.control.v1.AgentControl.DeleteNodeLabel:output_type -> persys.control.v1.DeleteNodeLabelResponse + 4, // 119: persys.control.v1.AgentControl.SubmitAutomationSuggestion:output_type -> persys.control.v1.SubmitAutomationSuggestionResponse + 47, // 120: persys.control.v1.AgentControl.ListNodes:output_type -> persys.control.v1.ListNodesResponse + 48, // 121: persys.control.v1.AgentControl.GetNode:output_type -> persys.control.v1.GetNodeResponse + 56, // 122: persys.control.v1.AgentControl.ListWorkloads:output_type -> persys.control.v1.ListWorkloadsResponse + 57, // 123: persys.control.v1.AgentControl.GetWorkload:output_type -> persys.control.v1.GetWorkloadResponse + 60, // 124: persys.control.v1.AgentControl.GetClusterSummary:output_type -> persys.control.v1.GetClusterSummaryResponse + 53, // 125: persys.control.v1.AgentControl.ListEvents:output_type -> persys.control.v1.ListEventsResponse + 51, // 126: persys.control.v1.AgentControl.WatchEvents:output_type -> persys.control.v1.SchedulerEventView + 61, // 127: persys.control.v1.AgentControl.ControlStream:output_type -> persys.control.v1.ControlMessage + 63, // 128: persys.control.v1.AgentControl.CreateDisk:output_type -> persys.control.v1.CreateDiskResponse + 65, // 129: persys.control.v1.AgentControl.ListDisks:output_type -> persys.control.v1.ListDisksResponse + 67, // 130: persys.control.v1.AgentControl.GetDisk:output_type -> persys.control.v1.GetDiskResponse + 69, // 131: persys.control.v1.AgentControl.DeleteDisk:output_type -> persys.control.v1.DeleteDiskResponse + 72, // 132: persys.control.v1.AgentControl.CreateBucket:output_type -> persys.control.v1.CreateBucketResponse + 74, // 133: persys.control.v1.AgentControl.ListBuckets:output_type -> persys.control.v1.ListBucketsResponse + 76, // 134: persys.control.v1.AgentControl.GetBucket:output_type -> persys.control.v1.GetBucketResponse + 78, // 135: persys.control.v1.AgentControl.DeleteBucket:output_type -> persys.control.v1.DeleteBucketResponse + 80, // 136: persys.control.v1.AgentControl.GetBucketAccess:output_type -> persys.control.v1.GetBucketAccessResponse + 82, // 137: persys.control.v1.AgentControl.ListBucketObjects:output_type -> persys.control.v1.ListBucketObjectsResponse + 108, // [108:138] is the sub-list for method output_type + 78, // [78:108] is the sub-list for method input_type + 78, // [78:78] is the sub-list for extension type_name + 78, // [78:78] is the sub-list for extension extendee + 0, // [0:78] is the sub-list for field type_name } func init() { file_control_proto_init() } @@ -4918,7 +6611,7 @@ func file_control_proto_init() { GoPackagePath: reflect.TypeOf(x{}).PkgPath(), RawDescriptor: unsafe.Slice(unsafe.StringData(file_control_proto_rawDesc), len(file_control_proto_rawDesc)), NumEnums: 2, - NumMessages: 66, + NumMessages: 90, NumExtensions: 0, NumServices: 1, }, diff --git a/persys-gateway/internal/controlv1/control_grpc.pb.go b/persys-gateway/internal/controlv1/control_grpc.pb.go index 050608d..1ccea93 100644 --- a/persys-gateway/internal/controlv1/control_grpc.pb.go +++ b/persys-gateway/internal/controlv1/control_grpc.pb.go @@ -39,6 +39,16 @@ const ( AgentControl_ListEvents_FullMethodName = "/persys.control.v1.AgentControl/ListEvents" AgentControl_WatchEvents_FullMethodName = "/persys.control.v1.AgentControl/WatchEvents" AgentControl_ControlStream_FullMethodName = "/persys.control.v1.AgentControl/ControlStream" + AgentControl_CreateDisk_FullMethodName = "/persys.control.v1.AgentControl/CreateDisk" + AgentControl_ListDisks_FullMethodName = "/persys.control.v1.AgentControl/ListDisks" + AgentControl_GetDisk_FullMethodName = "/persys.control.v1.AgentControl/GetDisk" + AgentControl_DeleteDisk_FullMethodName = "/persys.control.v1.AgentControl/DeleteDisk" + AgentControl_CreateBucket_FullMethodName = "/persys.control.v1.AgentControl/CreateBucket" + AgentControl_ListBuckets_FullMethodName = "/persys.control.v1.AgentControl/ListBuckets" + AgentControl_GetBucket_FullMethodName = "/persys.control.v1.AgentControl/GetBucket" + AgentControl_DeleteBucket_FullMethodName = "/persys.control.v1.AgentControl/DeleteBucket" + AgentControl_GetBucketAccess_FullMethodName = "/persys.control.v1.AgentControl/GetBucketAccess" + AgentControl_ListBucketObjects_FullMethodName = "/persys.control.v1.AgentControl/ListBucketObjects" ) // AgentControlClient is the client API for AgentControl service. @@ -68,8 +78,9 @@ type AgentControlClient interface { ListWorkloads(ctx context.Context, in *ListWorkloadsRequest, opts ...grpc.CallOption) (*ListWorkloadsResponse, error) GetWorkload(ctx context.Context, in *GetWorkloadRequest, opts ...grpc.CallOption) (*GetWorkloadResponse, error) GetClusterSummary(ctx context.Context, in *GetClusterSummaryRequest, opts ...grpc.CallOption) (*GetClusterSummaryResponse, error) - // Cluster-wide events: node lost, workload scheduled, drift detected, - // etc (see internal/scheduler/events.go for producers). ListEvents is a + // Cluster-wide events: node joined, node lost, node left, workload + // scheduled, drift detected, retries, reschedules, etc (see + // internal/scheduler/events.go for producers). ListEvents is a // plain unary call (auto-bridged to REST by persys-gateway's // reflection-based grpcbridge, no gateway changes needed). WatchEvents // is a server-streaming call — grpcbridge explicitly does not bridge @@ -81,6 +92,18 @@ type AgentControlClient interface { WatchEvents(ctx context.Context, in *WatchEventsRequest, opts ...grpc.CallOption) (grpc.ServerStreamingClient[SchedulerEventView], error) // Optional future streaming channel ControlStream(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[ControlMessage, ControlMessage], error) + // Standalone disk inventory (managed volumes) + CreateDisk(ctx context.Context, in *CreateDiskRequest, opts ...grpc.CallOption) (*CreateDiskResponse, error) + ListDisks(ctx context.Context, in *ListDisksRequest, opts ...grpc.CallOption) (*ListDisksResponse, error) + GetDisk(ctx context.Context, in *GetDiskRequest, opts ...grpc.CallOption) (*GetDiskResponse, error) + DeleteDisk(ctx context.Context, in *DeleteDiskRequest, opts ...grpc.CallOption) (*DeleteDiskResponse, error) + // Object storage (Ceph RGW / S3-compatible buckets) + CreateBucket(ctx context.Context, in *CreateBucketRequest, opts ...grpc.CallOption) (*CreateBucketResponse, error) + ListBuckets(ctx context.Context, in *ListBucketsRequest, opts ...grpc.CallOption) (*ListBucketsResponse, error) + GetBucket(ctx context.Context, in *GetBucketRequest, opts ...grpc.CallOption) (*GetBucketResponse, error) + DeleteBucket(ctx context.Context, in *DeleteBucketRequest, opts ...grpc.CallOption) (*DeleteBucketResponse, error) + GetBucketAccess(ctx context.Context, in *GetBucketAccessRequest, opts ...grpc.CallOption) (*GetBucketAccessResponse, error) + ListBucketObjects(ctx context.Context, in *ListBucketObjectsRequest, opts ...grpc.CallOption) (*ListBucketObjectsResponse, error) } type agentControlClient struct { @@ -303,6 +326,106 @@ func (c *agentControlClient) ControlStream(ctx context.Context, opts ...grpc.Cal // This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name. type AgentControl_ControlStreamClient = grpc.BidiStreamingClient[ControlMessage, ControlMessage] +func (c *agentControlClient) CreateDisk(ctx context.Context, in *CreateDiskRequest, opts ...grpc.CallOption) (*CreateDiskResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(CreateDiskResponse) + err := c.cc.Invoke(ctx, AgentControl_CreateDisk_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *agentControlClient) ListDisks(ctx context.Context, in *ListDisksRequest, opts ...grpc.CallOption) (*ListDisksResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(ListDisksResponse) + err := c.cc.Invoke(ctx, AgentControl_ListDisks_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *agentControlClient) GetDisk(ctx context.Context, in *GetDiskRequest, opts ...grpc.CallOption) (*GetDiskResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(GetDiskResponse) + err := c.cc.Invoke(ctx, AgentControl_GetDisk_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *agentControlClient) DeleteDisk(ctx context.Context, in *DeleteDiskRequest, opts ...grpc.CallOption) (*DeleteDiskResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(DeleteDiskResponse) + err := c.cc.Invoke(ctx, AgentControl_DeleteDisk_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *agentControlClient) CreateBucket(ctx context.Context, in *CreateBucketRequest, opts ...grpc.CallOption) (*CreateBucketResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(CreateBucketResponse) + err := c.cc.Invoke(ctx, AgentControl_CreateBucket_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *agentControlClient) ListBuckets(ctx context.Context, in *ListBucketsRequest, opts ...grpc.CallOption) (*ListBucketsResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(ListBucketsResponse) + err := c.cc.Invoke(ctx, AgentControl_ListBuckets_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *agentControlClient) GetBucket(ctx context.Context, in *GetBucketRequest, opts ...grpc.CallOption) (*GetBucketResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(GetBucketResponse) + err := c.cc.Invoke(ctx, AgentControl_GetBucket_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *agentControlClient) DeleteBucket(ctx context.Context, in *DeleteBucketRequest, opts ...grpc.CallOption) (*DeleteBucketResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(DeleteBucketResponse) + err := c.cc.Invoke(ctx, AgentControl_DeleteBucket_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *agentControlClient) GetBucketAccess(ctx context.Context, in *GetBucketAccessRequest, opts ...grpc.CallOption) (*GetBucketAccessResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(GetBucketAccessResponse) + err := c.cc.Invoke(ctx, AgentControl_GetBucketAccess_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *agentControlClient) ListBucketObjects(ctx context.Context, in *ListBucketObjectsRequest, opts ...grpc.CallOption) (*ListBucketObjectsResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(ListBucketObjectsResponse) + err := c.cc.Invoke(ctx, AgentControl_ListBucketObjects_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + // AgentControlServer is the server API for AgentControl service. // All implementations must embed UnimplementedAgentControlServer // for forward compatibility. @@ -330,8 +453,9 @@ type AgentControlServer interface { ListWorkloads(context.Context, *ListWorkloadsRequest) (*ListWorkloadsResponse, error) GetWorkload(context.Context, *GetWorkloadRequest) (*GetWorkloadResponse, error) GetClusterSummary(context.Context, *GetClusterSummaryRequest) (*GetClusterSummaryResponse, error) - // Cluster-wide events: node lost, workload scheduled, drift detected, - // etc (see internal/scheduler/events.go for producers). ListEvents is a + // Cluster-wide events: node joined, node lost, node left, workload + // scheduled, drift detected, retries, reschedules, etc (see + // internal/scheduler/events.go for producers). ListEvents is a // plain unary call (auto-bridged to REST by persys-gateway's // reflection-based grpcbridge, no gateway changes needed). WatchEvents // is a server-streaming call — grpcbridge explicitly does not bridge @@ -343,6 +467,18 @@ type AgentControlServer interface { WatchEvents(*WatchEventsRequest, grpc.ServerStreamingServer[SchedulerEventView]) error // Optional future streaming channel ControlStream(grpc.BidiStreamingServer[ControlMessage, ControlMessage]) error + // Standalone disk inventory (managed volumes) + CreateDisk(context.Context, *CreateDiskRequest) (*CreateDiskResponse, error) + ListDisks(context.Context, *ListDisksRequest) (*ListDisksResponse, error) + GetDisk(context.Context, *GetDiskRequest) (*GetDiskResponse, error) + DeleteDisk(context.Context, *DeleteDiskRequest) (*DeleteDiskResponse, error) + // Object storage (Ceph RGW / S3-compatible buckets) + CreateBucket(context.Context, *CreateBucketRequest) (*CreateBucketResponse, error) + ListBuckets(context.Context, *ListBucketsRequest) (*ListBucketsResponse, error) + GetBucket(context.Context, *GetBucketRequest) (*GetBucketResponse, error) + DeleteBucket(context.Context, *DeleteBucketRequest) (*DeleteBucketResponse, error) + GetBucketAccess(context.Context, *GetBucketAccessRequest) (*GetBucketAccessResponse, error) + ListBucketObjects(context.Context, *ListBucketObjectsRequest) (*ListBucketObjectsResponse, error) mustEmbedUnimplementedAgentControlServer() } @@ -413,6 +549,36 @@ func (UnimplementedAgentControlServer) WatchEvents(*WatchEventsRequest, grpc.Ser func (UnimplementedAgentControlServer) ControlStream(grpc.BidiStreamingServer[ControlMessage, ControlMessage]) error { return status.Error(codes.Unimplemented, "method ControlStream not implemented") } +func (UnimplementedAgentControlServer) CreateDisk(context.Context, *CreateDiskRequest) (*CreateDiskResponse, error) { + return nil, status.Error(codes.Unimplemented, "method CreateDisk not implemented") +} +func (UnimplementedAgentControlServer) ListDisks(context.Context, *ListDisksRequest) (*ListDisksResponse, error) { + return nil, status.Error(codes.Unimplemented, "method ListDisks not implemented") +} +func (UnimplementedAgentControlServer) GetDisk(context.Context, *GetDiskRequest) (*GetDiskResponse, error) { + return nil, status.Error(codes.Unimplemented, "method GetDisk not implemented") +} +func (UnimplementedAgentControlServer) DeleteDisk(context.Context, *DeleteDiskRequest) (*DeleteDiskResponse, error) { + return nil, status.Error(codes.Unimplemented, "method DeleteDisk not implemented") +} +func (UnimplementedAgentControlServer) CreateBucket(context.Context, *CreateBucketRequest) (*CreateBucketResponse, error) { + return nil, status.Error(codes.Unimplemented, "method CreateBucket not implemented") +} +func (UnimplementedAgentControlServer) ListBuckets(context.Context, *ListBucketsRequest) (*ListBucketsResponse, error) { + return nil, status.Error(codes.Unimplemented, "method ListBuckets not implemented") +} +func (UnimplementedAgentControlServer) GetBucket(context.Context, *GetBucketRequest) (*GetBucketResponse, error) { + return nil, status.Error(codes.Unimplemented, "method GetBucket not implemented") +} +func (UnimplementedAgentControlServer) DeleteBucket(context.Context, *DeleteBucketRequest) (*DeleteBucketResponse, error) { + return nil, status.Error(codes.Unimplemented, "method DeleteBucket not implemented") +} +func (UnimplementedAgentControlServer) GetBucketAccess(context.Context, *GetBucketAccessRequest) (*GetBucketAccessResponse, error) { + return nil, status.Error(codes.Unimplemented, "method GetBucketAccess not implemented") +} +func (UnimplementedAgentControlServer) ListBucketObjects(context.Context, *ListBucketObjectsRequest) (*ListBucketObjectsResponse, error) { + return nil, status.Error(codes.Unimplemented, "method ListBucketObjects not implemented") +} func (UnimplementedAgentControlServer) mustEmbedUnimplementedAgentControlServer() {} func (UnimplementedAgentControlServer) testEmbeddedByValue() {} @@ -776,6 +942,186 @@ func _AgentControl_ControlStream_Handler(srv interface{}, stream grpc.ServerStre // This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name. type AgentControl_ControlStreamServer = grpc.BidiStreamingServer[ControlMessage, ControlMessage] +func _AgentControl_CreateDisk_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(CreateDiskRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AgentControlServer).CreateDisk(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: AgentControl_CreateDisk_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AgentControlServer).CreateDisk(ctx, req.(*CreateDiskRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _AgentControl_ListDisks_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(ListDisksRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AgentControlServer).ListDisks(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: AgentControl_ListDisks_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AgentControlServer).ListDisks(ctx, req.(*ListDisksRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _AgentControl_GetDisk_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(GetDiskRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AgentControlServer).GetDisk(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: AgentControl_GetDisk_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AgentControlServer).GetDisk(ctx, req.(*GetDiskRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _AgentControl_DeleteDisk_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(DeleteDiskRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AgentControlServer).DeleteDisk(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: AgentControl_DeleteDisk_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AgentControlServer).DeleteDisk(ctx, req.(*DeleteDiskRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _AgentControl_CreateBucket_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(CreateBucketRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AgentControlServer).CreateBucket(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: AgentControl_CreateBucket_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AgentControlServer).CreateBucket(ctx, req.(*CreateBucketRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _AgentControl_ListBuckets_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(ListBucketsRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AgentControlServer).ListBuckets(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: AgentControl_ListBuckets_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AgentControlServer).ListBuckets(ctx, req.(*ListBucketsRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _AgentControl_GetBucket_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(GetBucketRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AgentControlServer).GetBucket(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: AgentControl_GetBucket_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AgentControlServer).GetBucket(ctx, req.(*GetBucketRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _AgentControl_DeleteBucket_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(DeleteBucketRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AgentControlServer).DeleteBucket(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: AgentControl_DeleteBucket_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AgentControlServer).DeleteBucket(ctx, req.(*DeleteBucketRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _AgentControl_GetBucketAccess_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(GetBucketAccessRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AgentControlServer).GetBucketAccess(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: AgentControl_GetBucketAccess_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AgentControlServer).GetBucketAccess(ctx, req.(*GetBucketAccessRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _AgentControl_ListBucketObjects_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(ListBucketObjectsRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AgentControlServer).ListBucketObjects(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: AgentControl_ListBucketObjects_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AgentControlServer).ListBucketObjects(ctx, req.(*ListBucketObjectsRequest)) + } + return interceptor(ctx, in, info, handler) +} + // AgentControl_ServiceDesc is the grpc.ServiceDesc for AgentControl service. // It's only intended for direct use with grpc.RegisterService, // and not to be introspected or modified (even as a copy) @@ -855,6 +1201,46 @@ var AgentControl_ServiceDesc = grpc.ServiceDesc{ MethodName: "ListEvents", Handler: _AgentControl_ListEvents_Handler, }, + { + MethodName: "CreateDisk", + Handler: _AgentControl_CreateDisk_Handler, + }, + { + MethodName: "ListDisks", + Handler: _AgentControl_ListDisks_Handler, + }, + { + MethodName: "GetDisk", + Handler: _AgentControl_GetDisk_Handler, + }, + { + MethodName: "DeleteDisk", + Handler: _AgentControl_DeleteDisk_Handler, + }, + { + MethodName: "CreateBucket", + Handler: _AgentControl_CreateBucket_Handler, + }, + { + MethodName: "ListBuckets", + Handler: _AgentControl_ListBuckets_Handler, + }, + { + MethodName: "GetBucket", + Handler: _AgentControl_GetBucket_Handler, + }, + { + MethodName: "DeleteBucket", + Handler: _AgentControl_DeleteBucket_Handler, + }, + { + MethodName: "GetBucketAccess", + Handler: _AgentControl_GetBucketAccess_Handler, + }, + { + MethodName: "ListBucketObjects", + Handler: _AgentControl_ListBucketObjects_Handler, + }, }, Streams: []grpc.StreamDesc{ { diff --git a/persys-gateway/internal/middleware/access_log.go b/persys-gateway/internal/middleware/access_log.go new file mode 100644 index 0000000..d577946 --- /dev/null +++ b/persys-gateway/internal/middleware/access_log.go @@ -0,0 +1,27 @@ +package middleware + +import ( + + "github.com/gin-gonic/gin" +) + +// AccessLogSkipPaths are high-churn probes that should not fill access logs. +// Used with gin.LoggerWithConfig. +var AccessLogSkipPaths = []string{ + "/metrics", + "/health", + "/healthz", + "/ready", + "/readyz", + "/livez", + "/favicon.ico", +} + +// AccessLogger is gin.Logger that skips metrics/health scrape noise. +func AccessLogger() gin.HandlerFunc { + return gin.LoggerWithConfig(gin.LoggerConfig{ + SkipPaths: AccessLogSkipPaths, + // Also skip paths that only differ by trailing slash or cluster prefix noise. + + }) +} diff --git a/persys-gateway/internal/router/bindings.go b/persys-gateway/internal/router/bindings.go index 0bacedd..42c2c30 100644 --- a/persys-gateway/internal/router/bindings.go +++ b/persys-gateway/internal/router/bindings.go @@ -87,6 +87,18 @@ func ClusterControlBinding(clusterControl ClusterControlBindingInvoker, r *Route // DeleteNodeLabelRequest.NodeId (Key is a flat top-level field) {Method: "DeleteNodeLabel", Verb: "DELETE", Path: "/nodes/:id/labels", PathParams: map[string]string{"id": "node_id"}}, {Method: "GetClusterSummary", Verb: "GET", Path: "/cluster/metrics"}, + // Standalone disks (AgentControl mTLS gRPC) + {Method: "ListDisks", Verb: "GET", Path: "/disks"}, + {Method: "CreateDisk", Verb: "POST", Path: "/disks"}, + {Method: "GetDisk", Verb: "GET", Path: "/disks/:id", PathParams: map[string]string{"id": "disk_id"}}, + {Method: "DeleteDisk", Verb: "DELETE", Path: "/disks/:id", PathParams: map[string]string{"id": "disk_id"}}, + // Object storage (Ceph RGW / S3) + {Method: "ListBuckets", Verb: "GET", Path: "/buckets"}, + {Method: "CreateBucket", Verb: "POST", Path: "/buckets"}, + {Method: "GetBucket", Verb: "GET", Path: "/buckets/:id", PathParams: map[string]string{"id": "bucket_id"}}, + {Method: "DeleteBucket", Verb: "DELETE", Path: "/buckets/:id", PathParams: map[string]string{"id": "bucket_id"}}, + {Method: "GetBucketAccess", Verb: "GET", Path: "/buckets/:id/access", PathParams: map[string]string{"id": "bucket_id"}}, + {Method: "ListBucketObjects", Verb: "GET", Path: "/buckets/:id/objects", PathParams: map[string]string{"id": "bucket_id"}}, // RegisterNode, Heartbeat: intentionally NOT aliased — those // are compute-agent-to-scheduler internal calls, not // SDK/persysctl surface. Still technically reachable via diff --git a/persys-gateway/services/automation.service.go b/persys-gateway/services/automation.service.go index 807b3e6..a72c9b2 100644 --- a/persys-gateway/services/automation.service.go +++ b/persys-gateway/services/automation.service.go @@ -7,9 +7,8 @@ import ( "time" "github.com/persys-dev/persys-cloud/persys-gateway/config" + "github.com/persys-dev/persys-cloud/pkg/certmanager" automationv1 "github.com/persys-dev/persys-cloud/pkg/automation/automationv1" - "google.golang.org/grpc" - "google.golang.org/grpc/credentials" ) // AutomationService proxies gRPC calls from the gateway to persys-automation. @@ -17,6 +16,8 @@ type AutomationService struct { config *config.Config clientTLS *tls.Config timeout time.Duration + + certMgr *certmanager.Manager } // NewAutomationService builds an AutomationService, reusing the gateway's own @@ -29,6 +30,11 @@ func NewAutomationService(cfg *config.Config, clientTLS *tls.Config) (*Automatio return &AutomationService{config: cfg, clientTLS: clientTLS, timeout: timeout}, nil } +func (s *AutomationService) SetCertManager(m *certmanager.Manager) { + s.certMgr = m +} + + func (s *AutomationService) CreatePolicy(ctx context.Context, req *automationv1.CreatePolicyRequest) (*automationv1.CreatePolicyResponse, error) { resp, err := s.invoke(ctx, func(client automationv1.AutomationControlClient) (any, error) { return client.CreatePolicy(ctx, req) @@ -106,10 +112,7 @@ func (s *AutomationService) invoke(ctx context.Context, call func(automationv1.A automationTLS.ServerName = serverName } - conn, err := grpc.DialContext(callCtx, s.config.Automation.GRPCAddr, - grpc.WithTransportCredentials(credentials.NewTLS(automationTLS)), - grpc.WithBlock(), - ) + conn, err := dialGRPCTLS(callCtx, s.config.Automation.GRPCAddr, automationTLS, s.certMgr) if err != nil { return nil, fmt.Errorf("dial automation %s: %w", s.config.Automation.GRPCAddr, err) } diff --git a/persys-gateway/services/forgery.service.go b/persys-gateway/services/forgery.service.go index fb51538..e35aaa0 100644 --- a/persys-gateway/services/forgery.service.go +++ b/persys-gateway/services/forgery.service.go @@ -7,9 +7,9 @@ import ( "time" "github.com/persys-dev/persys-cloud/persys-gateway/config" + "github.com/persys-dev/persys-cloud/pkg/certmanager" forgeryv1 "github.com/persys-dev/persys-cloud/persys-gateway/internal/forgeryv1" "google.golang.org/grpc" - "google.golang.org/grpc/credentials" "google.golang.org/protobuf/proto" ) @@ -22,12 +22,18 @@ import ( type ForgeryService struct { cfg *config.Config clientTLS *tls.Config + + certMgr *certmanager.Manager } func NewForgeryService(cfg *config.Config, clientTLS *tls.Config) *ForgeryService { return &ForgeryService{cfg: cfg, clientTLS: clientTLS} } +func (s *ForgeryService) SetCertManager(m *certmanager.Manager) { + s.certMgr = m +} + func (s *ForgeryService) dial(ctx context.Context, timeout time.Duration) (*grpc.ClientConn, error) { var forgeryTLS *tls.Config if s.clientTLS != nil { @@ -41,10 +47,7 @@ func (s *ForgeryService) dial(ctx context.Context, timeout time.Duration) (*grpc callCtx, cancel := context.WithTimeout(ctx, timeout) defer cancel() - conn, err := grpc.DialContext(callCtx, s.cfg.Forgery.GRPCAddr, - grpc.WithTransportCredentials(credentials.NewTLS(forgeryTLS)), - grpc.WithBlock(), - ) + conn, err := dialGRPCTLS(callCtx, s.cfg.Forgery.GRPCAddr, forgeryTLS, s.certMgr) if err != nil { return nil, fmt.Errorf("dial forgery %s: %w", s.cfg.Forgery.GRPCAddr, err) } diff --git a/persys-gateway/services/github.impl.service.go b/persys-gateway/services/github.impl.service.go index feb141f..82156db 100644 --- a/persys-gateway/services/github.impl.service.go +++ b/persys-gateway/services/github.impl.service.go @@ -8,15 +8,17 @@ import ( "time" "github.com/persys-dev/persys-cloud/persys-gateway/config" + "github.com/persys-dev/persys-cloud/pkg/certmanager" forgeryv1 "github.com/persys-dev/persys-cloud/persys-gateway/internal/forgeryv1" "github.com/persys-dev/persys-cloud/persys-gateway/models" "google.golang.org/grpc" - "google.golang.org/grpc/credentials" ) type GithubServiceImpl struct { cfg *config.Config tlsClient *tls.Config + + certMgr *certmanager.Manager } // NewGithubService previously took an unused *mongo.Collection parameter @@ -27,6 +29,11 @@ func NewGithubService(cfg *config.Config, tlsClient *tls.Config) GithubService { return &GithubServiceImpl{cfg: cfg, tlsClient: tlsClient} } +func (g *GithubServiceImpl) SetCertManager(m *certmanager.Manager) { + g.certMgr = m +} + + func (g *GithubServiceImpl) SetAccessToken(user *models.DBResponse) error { ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() @@ -105,10 +112,7 @@ func (g *GithubServiceImpl) SetWebhook(user *models.DBResponse, repository strin } func (g *GithubServiceImpl) forgeryClient(ctx context.Context) (forgeryv1.ForgeryControlClient, *grpc.ClientConn, error) { - conn, err := grpc.DialContext(ctx, g.cfg.Forgery.GRPCAddr, - grpc.WithTransportCredentials(credentials.NewTLS(g.tlsClient)), - grpc.WithBlock(), - ) + conn, err := dialGRPCTLS(ctx, g.cfg.Forgery.GRPCAddr, g.tlsClient, g.certMgr) if err != nil { return nil, nil, err } diff --git a/persys-gateway/services/mtls_dial.go b/persys-gateway/services/mtls_dial.go new file mode 100644 index 0000000..2739216 --- /dev/null +++ b/persys-gateway/services/mtls_dial.go @@ -0,0 +1,74 @@ +package services + +import ( + "context" + "crypto/tls" + "crypto/x509" + "fmt" + "os" + + "github.com/persys-dev/persys-cloud/pkg/certmanager" + "google.golang.org/grpc" + "google.golang.org/grpc/credentials" +) + +// LiveClientTLSConfig returns a tls.Config that reloads the client keypair from +// disk on every handshake (GetClientCertificate) and uses the given CA pool. +// Call this once after certmanager has written certs; subsequent ForceRotate +// writes are picked up on the next handshake without rebuilding the config. +func LiveClientTLSConfig(certPath, keyPath, caPath string) (*tls.Config, error) { + caPEM, err := os.ReadFile(caPath) + if err != nil { + return nil, fmt.Errorf("read CA %s: %w", caPath, err) + } + pool := x509.NewCertPool() + if !pool.AppendCertsFromPEM(caPEM) { + return nil, fmt.Errorf("invalid CA bundle at %s", caPath) + } + // Validate keypair exists now. + if _, err := tls.LoadX509KeyPair(certPath, keyPath); err != nil { + return nil, fmt.Errorf("load client keypair: %w", err) + } + return &tls.Config{ + MinVersion: tls.VersionTLS12, + RootCAs: pool, + GetClientCertificate: func(*tls.CertificateRequestInfo) (*tls.Certificate, error) { + cert, err := tls.LoadX509KeyPair(certPath, keyPath) + if err != nil { + return nil, err + } + return &cert, nil + }, + }, nil +} + +// dialGRPCTLS dials addr with tlsCfg. When mgr is non-nil, cert-related TLS +// failures trigger ForceRotate and up to 3 dial attempts. +func dialGRPCTLS(ctx context.Context, addr string, tlsCfg *tls.Config, mgr *certmanager.Manager, extra ...grpc.DialOption) (*grpc.ClientConn, error) { + if tlsCfg == nil { + return nil, fmt.Errorf("dial %s: tls config is nil", addr) + } + doDial := func(ctx context.Context) (*grpc.ClientConn, error) { + opts := []grpc.DialOption{ + grpc.WithTransportCredentials(credentials.NewTLS(tlsCfg)), + grpc.WithBlock(), + } + opts = append(opts, extra...) + return grpc.DialContext(ctx, addr, opts...) + } + + if mgr == nil { + return doDial(ctx) + } + + var conn *grpc.ClientConn + err := certmanager.WithCertRetry(ctx, mgr, 3, func(ctx context.Context) error { + c, err := doDial(ctx) + if err != nil { + return err + } + conn = c + return nil + }) + return conn, err +} diff --git a/persys-gateway/services/scheduler.service.go b/persys-gateway/services/scheduler.service.go index 110f5d5..0e5278f 100644 --- a/persys-gateway/services/scheduler.service.go +++ b/persys-gateway/services/scheduler.service.go @@ -7,6 +7,7 @@ import ( "errors" "fmt" "github.com/persys-dev/persys-cloud/persys-gateway/config" + "github.com/persys-dev/persys-cloud/pkg/certmanager" "os" ) @@ -15,6 +16,7 @@ type ClusterControlService struct { clientTLS *tls.Config serverTLS *tls.Config schedulerPool *SchedulerPoolManager + certMgr *certmanager.Manager } func NewClusterControlService(cfg *config.Config) *ClusterControlService { @@ -38,26 +40,51 @@ func (s *ClusterControlService) Start(ctx context.Context) { } func (s *ClusterControlService) loadTLSConfigs() error { - cert, err := tls.LoadX509KeyPair(s.config.TLS.CertPath, s.config.TLS.KeyPath) + clientTLS, err := LiveClientTLSConfig(s.config.TLS.CertPath, s.config.TLS.KeyPath, s.config.TLS.CAPath) if err != nil { - return fmt.Errorf("failed to load client certificate: %w", err) + return fmt.Errorf("failed to load client TLS config: %w", err) } + s.clientTLS = clientTLS + // Server-side config still needs a concrete certificate for GetCertificate-style + // use; load once and also expose ClientCAs from the same CA path. + cert, err := tls.LoadX509KeyPair(s.config.TLS.CertPath, s.config.TLS.KeyPath) + if err != nil { + return fmt.Errorf("failed to load server certificate: %w", err) + } caCert, err := os.ReadFile(s.config.TLS.CAPath) if err != nil { return fmt.Errorf("failed to read CA certificate: %w", err) } - caCertPool := x509.NewCertPool() if !caCertPool.AppendCertsFromPEM(caCert) { return fmt.Errorf("failed to append CA certificate") } - - s.clientTLS = &tls.Config{Certificates: []tls.Certificate{cert}, RootCAs: caCertPool} - s.serverTLS = &tls.Config{Certificates: []tls.Certificate{cert}, ClientCAs: caCertPool, ClientAuth: tls.RequireAndVerifyClientCert} + s.serverTLS = &tls.Config{ + MinVersion: tls.VersionTLS12, + Certificates: []tls.Certificate{cert}, + ClientCAs: caCertPool, + ClientAuth: tls.RequireAndVerifyClientCert, + GetCertificate: func(*tls.ClientHelloInfo) (*tls.Certificate, error) { + c, err := tls.LoadX509KeyPair(s.config.TLS.CertPath, s.config.TLS.KeyPath) + if err != nil { + return nil, err + } + return &c, nil + }, + } return nil } +// SetCertManager wires certmanager so outbound scheduler dials can ForceRotate +// and retry on cert-related TLS failures. Propagates to the scheduler pool. +func (s *ClusterControlService) SetCertManager(m *certmanager.Manager) { + s.certMgr = m + if s.schedulerPool != nil { + s.schedulerPool.SetCertManager(m) + } +} + func (s *ClusterControlService) DiscoverAndPrintSchedulers() { s.schedulerPool.ForceDiscover(context.Background()) } diff --git a/persys-gateway/services/scheduler_events.service.go b/persys-gateway/services/scheduler_events.service.go index 327a24b..9d12a46 100644 --- a/persys-gateway/services/scheduler_events.service.go +++ b/persys-gateway/services/scheduler_events.service.go @@ -6,8 +6,6 @@ import ( "io" controlv1 "github.com/persys-dev/persys-cloud/persys-gateway/internal/controlv1" - "google.golang.org/grpc" - "google.golang.org/grpc/credentials" ) // WatchEventsForClient opens a server-streaming WatchEvents call against a @@ -32,10 +30,7 @@ func (s *ClusterControlService) WatchEventsForClient(ctx context.Context, cluste var lastErr error for _, target := range candidates { - conn, dialErr := grpc.DialContext(ctx, target.Address, - grpc.WithTransportCredentials(credentials.NewTLS(s.clientTLS)), - grpc.WithBlock(), - ) + conn, dialErr := dialGRPCTLS(ctx, target.Address, s.clientTLS, s.certMgr) if dialErr != nil { s.schedulerPool.MarkUnhealthy(clusterID, target.Address) lastErr = dialErr diff --git a/persys-gateway/services/scheduler_grpc.service.go b/persys-gateway/services/scheduler_grpc.service.go index 09690b3..aa4c3ac 100644 --- a/persys-gateway/services/scheduler_grpc.service.go +++ b/persys-gateway/services/scheduler_grpc.service.go @@ -6,10 +6,10 @@ import ( "strings" "time" + "github.com/persys-dev/persys-cloud/pkg/certmanager" controlv1 "github.com/persys-dev/persys-cloud/persys-gateway/internal/controlv1" "go.opentelemetry.io/otel" "google.golang.org/grpc" - "google.golang.org/grpc/credentials" "google.golang.org/grpc/metadata" "google.golang.org/protobuf/proto" ) @@ -167,10 +167,7 @@ func (s *ClusterControlService) invokeControlRPC(ctx context.Context, clusterID, var lastErr error for _, target := range candidates { callCtx, cancel := context.WithTimeout(ctx, 10*time.Second) - conn, dialErr := grpc.DialContext(callCtx, target.Address, - grpc.WithTransportCredentials(credentials.NewTLS(s.clientTLS)), - grpc.WithBlock(), - ) + conn, dialErr := dialGRPCTLS(callCtx, target.Address, s.clientTLS, s.certMgr) cancel() if dialErr != nil { s.schedulerPool.MarkUnhealthy(clusterID, target.Address) @@ -183,6 +180,9 @@ func (s *ClusterControlService) invokeControlRPC(ctx context.Context, clusterID, resp, rpcErr := call(clientFromContext(client, callWithTrace)) _ = conn.Close() if rpcErr != nil { + if certmanager.IsCertRelatedTLSError(rpcErr) && s.certMgr != nil { + _ = s.certMgr.ForceRotate(ctx) + } s.schedulerPool.MarkUnhealthy(clusterID, target.Address) lastErr = rpcErr continue @@ -216,10 +216,7 @@ func (s *ClusterControlService) InvokeDynamic(ctx context.Context, clusterID, se var lastErr error for _, target := range candidates { callCtx, cancel := context.WithTimeout(ctx, 10*time.Second) - conn, dialErr := grpc.DialContext(callCtx, target.Address, - grpc.WithTransportCredentials(credentials.NewTLS(s.clientTLS)), - grpc.WithBlock(), - ) + conn, dialErr := dialGRPCTLS(callCtx, target.Address, s.clientTLS, s.certMgr) cancel() if dialErr != nil { s.schedulerPool.MarkUnhealthy(clusterID, target.Address) @@ -231,6 +228,9 @@ func (s *ClusterControlService) InvokeDynamic(ctx context.Context, clusterID, se rpcErr := conn.Invoke(callWithTrace, fullMethod, in, out) _ = conn.Close() if rpcErr != nil { + if certmanager.IsCertRelatedTLSError(rpcErr) && s.certMgr != nil { + _ = s.certMgr.ForceRotate(ctx) + } s.schedulerPool.MarkUnhealthy(clusterID, target.Address) lastErr = rpcErr continue @@ -258,10 +258,7 @@ func (s *ClusterControlService) DialForReflection(ctx context.Context, clusterID if err != nil || len(candidates) == 0 { return nil, fmt.Errorf("no scheduler candidates for cluster %q: %w", clusterID, err) } - return grpc.DialContext(ctx, candidates[0].Address, - grpc.WithTransportCredentials(credentials.NewTLS(s.clientTLS)), - grpc.WithBlock(), - ) + return dialGRPCTLS(ctx, candidates[0].Address, s.clientTLS, s.certMgr) } // Forgery methods (TriggerBuild, UpsertProject, ForwardWebhookTest, diff --git a/persys-gateway/services/scheduler_pool.go b/persys-gateway/services/scheduler_pool.go index c062041..4f28995 100644 --- a/persys-gateway/services/scheduler_pool.go +++ b/persys-gateway/services/scheduler_pool.go @@ -14,9 +14,8 @@ import ( "time" "github.com/persys-dev/persys-cloud/persys-gateway/config" + "github.com/persys-dev/persys-cloud/pkg/certmanager" "github.com/sirupsen/logrus" - "google.golang.org/grpc" - "google.golang.org/grpc/credentials" ) var ( @@ -52,6 +51,7 @@ type Cluster struct { type SchedulerPoolManager struct { cfg *config.Config tlsClient *tls.Config + certMgr *certmanager.Manager healthPath string healthInterval time.Duration discoveryInterval time.Duration @@ -100,6 +100,12 @@ func NewSchedulerPoolManager(cfg *config.Config, tlsClient *tls.Config) (*Schedu return m, nil } +// SetCertManager enables ForceRotate + retry on cert-related TLS failures +// during health probes and (via ClusterControlService) control-plane dials. +func (m *SchedulerPoolManager) SetCertManager(mgr *certmanager.Manager) { + m.certMgr = mgr +} + func (m *SchedulerPoolManager) Start(ctx context.Context) { m.logger.WithFields(logrus.Fields{ "default_cluster_id": m.cfg.Scheduler.DefaultClusterID, @@ -381,12 +387,7 @@ func (m *SchedulerPoolManager) checkInstanceHealth(ctx context.Context, address healthCtx, cancel := context.WithTimeout(ctx, 4*time.Second) defer cancel() - conn, err := grpc.DialContext( - healthCtx, - address, - grpc.WithTransportCredentials(credentials.NewTLS(m.tlsClient)), - grpc.WithBlock(), - ) + conn, err := dialGRPCTLS(healthCtx, address, m.tlsClient, m.certMgr) if err != nil { m.logger.WithFields(logrus.Fields{ "scheduler": address, diff --git a/persys-gateway/services/webhook.service.go b/persys-gateway/services/webhook.service.go index eeaf95e..76a06d6 100644 --- a/persys-gateway/services/webhook.service.go +++ b/persys-gateway/services/webhook.service.go @@ -15,13 +15,12 @@ import ( "time" "github.com/persys-dev/persys-cloud/persys-gateway/config" + "github.com/persys-dev/persys-cloud/pkg/certmanager" forgeryv1 "github.com/persys-dev/persys-cloud/persys-gateway/internal/forgeryv1" "github.com/persys-dev/persys-cloud/persys-gateway/internal/store" "github.com/persys-dev/persys-cloud/persys-gateway/models" "go.opentelemetry.io/otel" "go.opentelemetry.io/otel/propagation" - "google.golang.org/grpc" - "google.golang.org/grpc/credentials" ) type WebhookService interface { @@ -39,6 +38,8 @@ type webhookService struct { cacheMu sync.Mutex deliverySeenAt map[string]time.Time jobs chan forwardJob + + certMgr *certmanager.Manager } type forwardJob struct { @@ -92,6 +93,11 @@ func NewWebhookService(cfg *config.Config, tlsClient *tls.Config, st *store.Stor }, nil } +func (w *webhookService) SetCertManager(m *certmanager.Manager) { + w.certMgr = m +} + + func (w *webhookService) Start(ctx context.Context) { for i := 0; i < 2; i++ { go w.runWorker(ctx) @@ -285,10 +291,7 @@ func (w *webhookService) forwardGRPC(ctx context.Context, eventName, repo, clust dialCtx, cancel := context.WithTimeout(ctx, 10*time.Second) defer cancel() - conn, err := grpc.DialContext(dialCtx, w.cfg.Forgery.GRPCAddr, - grpc.WithTransportCredentials(credentials.NewTLS(w.tlsConfig)), - grpc.WithBlock(), - ) + conn, err := dialGRPCTLS(dialCtx, w.cfg.Forgery.GRPCAddr, w.tlsConfig, w.certMgr) if err != nil { return err }