From c757ef9a913a3e6ed9f79eee33b14592fca048c0 Mon Sep 17 00:00:00 2001 From: milx Date: Fri, 7 Aug 2026 19:27:42 +0330 Subject: [PATCH 1/6] Feat: Add Meter and Events Commands inder metrics --- persysctl/cmd/metrics.go | 88 ++++++++++++++++++++++++++++++++++++++++ 1 file changed, 88 insertions(+) diff --git a/persysctl/cmd/metrics.go b/persysctl/cmd/metrics.go index b71a8e2..f328bcb 100644 --- a/persysctl/cmd/metrics.go +++ b/persysctl/cmd/metrics.go @@ -3,10 +3,24 @@ package cmd import ( "encoding/json" "fmt" + "time" "github.com/spf13/cobra" ) +var ( + metricsFrom string + metricsTo string + meterWorkloadID string + meterNodeID string + meterType string + meterLimit int + eventsType string + eventsWorkloadID string + eventsNodeID string + eventsLimit int +) + var metricsCmd = &cobra.Command{ Use: "metrics", Short: "View Persys Compute metrics", @@ -23,6 +37,80 @@ var metricsCmd = &cobra.Command{ }, } +var meterCmd = &cobra.Command{ + Use: "meter", + Short: "Inspect workload usage via the gateway meter API", +} + +var meterSummaryCmd = &cobra.Command{ + Use: "summary [workload-id]", + Short: "Get usage summary for a workload", + Args: cobra.ExactArgs(1), + Run: func(cmd *cobra.Command, args []string) { + c, _, err := newClientWithTrace() + cobra.CheckErr(err) + defer c.Close() + from, err := parseTimeFlag(metricsFrom) + cobra.CheckErr(err) + to, err := parseTimeFlag(metricsTo) + cobra.CheckErr(err) + resp, err := c.MeterSummary(args[0], from, to) + cobra.CheckErr(err) + printJSON(resp) + }, +} + +var meterListCmd = &cobra.Command{ + Use: "list", + Short: "List current workload usage samples", + Run: func(cmd *cobra.Command, args []string) { + c, _, err := newClientWithTrace() + cobra.CheckErr(err) + defer c.Close() + resp, err := c.MeterListWorkloads(meterType) + cobra.CheckErr(err) + printJSON(resp) + }, +} + +var eventsCmd = &cobra.Command{ + Use: "events", + Short: "Watch scheduler events through the gateway SSE endpoint", +} + +var eventsWatchCmd = &cobra.Command{ + Use: "watch", + Short: "Stream scheduler events matching optional filters", + Run: func(cmd *cobra.Command, args []string) { + c, _, err := newClientWithTrace() + cobra.CheckErr(err) + defer c.Close() + resp, err := c.WatchGatewayEvents(eventsType, eventsWorkloadID, eventsNodeID, eventsLimit) + cobra.CheckErr(err) + printJSON(resp) + }, +} + +func parseTimeFlag(raw string) (time.Time, error) { + if raw == "" { + return time.Time{}, nil + } + return time.Parse(time.RFC3339, raw) +} + func init() { rootCmd.AddCommand(metricsCmd) + rootCmd.AddCommand(meterCmd) + rootCmd.AddCommand(eventsCmd) + meterCmd.AddCommand(meterSummaryCmd) + meterCmd.AddCommand(meterListCmd) + eventsCmd.AddCommand(eventsWatchCmd) + + meterSummaryCmd.Flags().StringVar(&metricsFrom, "from", "", "Start time in RFC3339 format") + meterSummaryCmd.Flags().StringVar(&metricsTo, "to", "", "End time in RFC3339 format") + meterListCmd.Flags().StringVar(&meterType, "type", "", "Workload type filter") + eventsWatchCmd.Flags().StringVar(&eventsType, "type", "", "Event type filter") + eventsWatchCmd.Flags().StringVar(&eventsWorkloadID, "workload-id", "", "Filter by workload ID") + eventsWatchCmd.Flags().StringVar(&eventsNodeID, "node-id", "", "Filter by node ID") + eventsWatchCmd.Flags().IntVar(&eventsLimit, "limit", 20, "Maximum number of events to print before exiting") } From 4203ba4c67f00f54e1823424922b7b334e342f1e Mon Sep 17 00:00:00 2001 From: milx Date: Fri, 7 Aug 2026 19:28:07 +0330 Subject: [PATCH 2/6] Feat: Add Client Unit Test --- persysctl/internal/client/client_test.go | 86 ++++++++++++++++++++++++ 1 file changed, 86 insertions(+) create mode 100644 persysctl/internal/client/client_test.go diff --git a/persysctl/internal/client/client_test.go b/persysctl/internal/client/client_test.go new file mode 100644 index 0000000..98d832a --- /dev/null +++ b/persysctl/internal/client/client_test.go @@ -0,0 +1,86 @@ +package client + +import ( + "net/http" + "net/http/httptest" + "testing" + "time" + + "github.com/persys-dev/persysctl/internal/config" +) + +func TestMeterSummaryGatewayPath(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.Method != http.MethodGet { + t.Fatalf("unexpected method: %s", r.Method) + } + if r.URL.Path != "/meter/v1/workloads/wl-123/summary" { + t.Fatalf("unexpected path: %s", r.URL.Path) + } + if got := r.URL.Query().Get("from"); got == "" { + t.Fatal("missing from query") + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"workload_id":"wl-123","sample_count":7}`)) + })) + defer server.Close() + + cli := &Client{cfg: config.Config{Transport: "http", APIEndpoint: server.URL}, httpClient: server.Client()} + from := time.Date(2026, 7, 1, 0, 0, 0, 0, time.UTC) + to := time.Date(2026, 7, 2, 0, 0, 0, 0, time.UTC) + + resp, err := cli.MeterSummary("wl-123", from, to) + if err != nil { + t.Fatalf("MeterSummary() error = %v", err) + } + if resp["workload_id"] != "wl-123" { + t.Fatalf("MeterSummary() workload_id = %v", resp["workload_id"]) + } +} + +func TestWatchGatewayEventsParsesSSE(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/events/watch" { + t.Fatalf("unexpected path: %s", r.URL.Path) + } + w.Header().Set("Content-Type", "text/event-stream") + _, _ = w.Write([]byte("event: event\ndata: {\"type\":\"workload\",\"workload_id\":\"wl-9\"}\n\n")) + })) + defer server.Close() + + cli := &Client{cfg: config.Config{Transport: "http", APIEndpoint: server.URL}, httpClient: server.Client()} + events, err := cli.WatchGatewayEvents("workload", "wl-9", "", 10) + if err != nil { + t.Fatalf("WatchGatewayEvents() error = %v", err) + } + if len(events) != 1 { + t.Fatalf("WatchGatewayEvents() len = %d", len(events)) + } + if got := events[0]["workload_id"]; got != "wl-9" { + t.Fatalf("WatchGatewayEvents() workload_id = %v", got) + } +} + +func TestWatchGatewayEventsTimesOutOnIdleStream(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "text/event-stream") + if f, ok := w.(http.Flusher); ok { + f.Flush() + } + <-r.Context().Done() + })) + defer server.Close() + + cli := &Client{cfg: config.Config{Transport: "http", APIEndpoint: server.URL, RPCTimeoutSeconds: 1}, httpClient: server.Client()} + start := time.Now() + events, err := cli.WatchGatewayEvents("", "", "", 20) + if err != nil { + t.Fatalf("WatchGatewayEvents() error = %v", err) + } + if len(events) != 0 { + t.Fatalf("WatchGatewayEvents() idle stream returned %d events", len(events)) + } + if elapsed := time.Since(start); elapsed > 3*time.Second { + t.Fatalf("WatchGatewayEvents() took too long: %s", elapsed) + } +} From 0204371cca3a92a7c85e7adf5f36dc6b37d1b3fb Mon Sep 17 00:00:00 2001 From: milx Date: Fri, 7 Aug 2026 19:28:36 +0330 Subject: [PATCH 3/6] Feat: Add Meter And Event Calls to api-gateway --- persysctl/internal/client/client.go | 165 ++++++++++++++++++++++++++++ 1 file changed, 165 insertions(+) diff --git a/persysctl/internal/client/client.go b/persysctl/internal/client/client.go index 0447cdd..16163aa 100644 --- a/persysctl/internal/client/client.go +++ b/persysctl/internal/client/client.go @@ -1,6 +1,7 @@ package client import ( + "bufio" "bytes" "context" "crypto/tls" @@ -1011,6 +1012,170 @@ func (c *Client) GatewayClusters() (*GatewayClustersResponse, error) { return &out, nil } +func (c *Client) MeterListWorkloads(workloadType string) (map[string]any, error) { + if c.cfg.Transport != "http" { + return nil, fmt.Errorf("meter API is available only with http transport") + } + path := "/meter/v1/workloads" + q := url.Values{} + if trimmed := strings.TrimSpace(workloadType); trimmed != "" { + q.Set("workload_type", trimmed) + } + if encoded := q.Encode(); encoded != "" { + path += "?" + encoded + } + var out map[string]any + if err := c.httpJSONRequest("GET", path, nil, &out); err != nil { + return nil, err + } + return out, nil +} + +func (c *Client) MeterGetWorkload(workloadID string) (map[string]any, error) { + if c.cfg.Transport != "http" { + return nil, fmt.Errorf("meter API is available only with http transport") + } + path := "/meter/v1/workloads/" + url.PathEscape(strings.TrimSpace(workloadID)) + var out map[string]any + if err := c.httpJSONRequest("GET", path, nil, &out); err != nil { + return nil, err + } + return out, nil +} + +func (c *Client) MeterHistory(workloadID string, from, to time.Time, limit int) (map[string]any, error) { + if c.cfg.Transport != "http" { + return nil, fmt.Errorf("meter API is available only with http transport") + } + path := "/meter/v1/workloads/" + url.PathEscape(strings.TrimSpace(workloadID)) + "/history" + q := url.Values{} + if !from.IsZero() { + q.Set("from", from.UTC().Format(time.RFC3339)) + } + if !to.IsZero() { + q.Set("to", to.UTC().Format(time.RFC3339)) + } + if limit > 0 { + q.Set("limit", strconv.Itoa(limit)) + } + if encoded := q.Encode(); encoded != "" { + path += "?" + encoded + } + var out map[string]any + if err := c.httpJSONRequest("GET", path, nil, &out); err != nil { + return nil, err + } + return out, nil +} + +func (c *Client) MeterSummary(workloadID string, from, to time.Time) (map[string]any, error) { + if c.cfg.Transport != "http" { + return nil, fmt.Errorf("meter API is available only with http transport") + } + path := "/meter/v1/workloads/" + url.PathEscape(strings.TrimSpace(workloadID)) + "/summary" + q := url.Values{} + if !from.IsZero() { + q.Set("from", from.UTC().Format(time.RFC3339)) + } + if !to.IsZero() { + q.Set("to", to.UTC().Format(time.RFC3339)) + } + if encoded := q.Encode(); encoded != "" { + path += "?" + encoded + } + var out map[string]any + if err := c.httpJSONRequest("GET", path, nil, &out); err != nil { + return nil, err + } + return out, nil +} + +func (c *Client) WatchGatewayEvents(eventType, workloadID, nodeID string, limit int) ([]map[string]any, error) { + if c.cfg.Transport != "http" { + return nil, fmt.Errorf("gateway event stream is available only with http transport") + } + if limit <= 0 { + limit = 20 + } + + path := "/events/watch" + q := url.Values{} + if trimmed := strings.TrimSpace(eventType); trimmed != "" { + q.Set("type", trimmed) + } + if trimmed := strings.TrimSpace(workloadID); trimmed != "" { + q.Set("workload_id", trimmed) + } + if trimmed := strings.TrimSpace(nodeID); trimmed != "" { + q.Set("node_id", trimmed) + } + if encoded := q.Encode(); encoded != "" { + path += "?" + encoded + } + + timeout := time.Duration(c.cfg.RPCTimeoutSeconds) * time.Second + if timeout <= 0 { + timeout = 15 * time.Second + } + ctx, cancel := context.WithTimeout(context.Background(), timeout) + defer cancel() + + req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.cfg.APIEndpoint+path, nil) + if err != nil { + return nil, fmt.Errorf("failed to create request: %w", err) + } + req.Header.Set("Accept", "text/event-stream") + resp, err := c.httpClient.Do(req) + if err != nil { + if ctx.Err() == context.DeadlineExceeded { + return nil, nil + } + return nil, fmt.Errorf("failed to send request: %w", err) + } + defer resp.Body.Close() + if resp.StatusCode < 200 || resp.StatusCode >= 300 { + body, _ := io.ReadAll(resp.Body) + return nil, fmt.Errorf("API returned status %d: %s", resp.StatusCode, string(body)) + } + + scanner := bufio.NewScanner(resp.Body) + scanner.Buffer(make([]byte, 0, 64*1024), 1024*1024) + var ( + events []map[string]any + dataLine string + ) + for scanner.Scan() { + line := strings.TrimSpace(scanner.Text()) + if strings.HasPrefix(line, "event:") { + continue + } + if strings.HasPrefix(line, "data:") { + dataLine = strings.TrimSpace(strings.TrimPrefix(line, "data:")) + continue + } + if line == "" && dataLine != "" { + var event map[string]any + if err := json.Unmarshal([]byte(dataLine), &event); err == nil { + events = append(events, event) + if len(events) >= limit { + return events, nil + } + } + dataLine = "" + } + if ctx.Err() != nil { + return events, nil + } + } + if err := scanner.Err(); err != nil { + if ctx.Err() == context.DeadlineExceeded { + return events, nil + } + return nil, fmt.Errorf("failed to read event stream: %w", err) + } + return events, nil +} + func (c *Client) TriggerForgeryBuild(req ForgeryBuildTriggerRequest) (map[string]interface{}, error) { if c.cfg.Transport != "http" { return nil, fmt.Errorf("forgery trigger-build is available only with http transport") From 340d8772f132c3ebc415999512e93ff708b83204 Mon Sep 17 00:00:00 2001 From: milx Date: Fri, 7 Aug 2026 19:29:10 +0330 Subject: [PATCH 4/6] Update: Vault Cert TTL To 1 Hour --- persysctl/internal/config/config.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/persysctl/internal/config/config.go b/persysctl/internal/config/config.go index 9a62c27..45c7995 100644 --- a/persysctl/internal/config/config.go +++ b/persysctl/internal/config/config.go @@ -138,7 +138,7 @@ func GetConfig() Config { if cfg.VaultPKIRole == "" { cfg.VaultPKIRole = "persysctl" } - cfg.VaultCertTTL = durationOr(stringWithEnv("vault_cert_ttl", "PERSYS_VAULT_CERT_TTL"), 24*time.Hour) + cfg.VaultCertTTL = durationOr(stringWithEnv("vault_cert_ttl", "PERSYS_VAULT_CERT_TTL"), 1*time.Hour) cfg.VaultServiceName = strings.TrimSpace(stringWithEnv("vault_service_name", "PERSYS_VAULT_SERVICE_NAME")) if cfg.VaultServiceName == "" { cfg.VaultServiceName = "persysctl" From 07bc3673dec4885e6cbcdc65976e69388feb9ed6 Mon Sep 17 00:00:00 2001 From: milx Date: Fri, 7 Aug 2026 19:30:16 +0330 Subject: [PATCH 5/6] Chore: Update Protobuf generated files for agent/scheduler in global pkg --- pkg/agent/api/v1/agent.pb.go | 132 ++++- pkg/scheduler/controlv1/control.pb.go | 531 ++++++++++++++++----- pkg/scheduler/controlv1/control_grpc.pb.go | 101 +++- 3 files changed, 640 insertions(+), 124 deletions(-) diff --git a/pkg/agent/api/v1/agent.pb.go b/pkg/agent/api/v1/agent.pb.go index c8f6996..eb9ee71 100644 --- a/pkg/agent/api/v1/agent.pb.go +++ b/pkg/agent/api/v1/agent.pb.go @@ -1221,8 +1221,15 @@ type VMSpec struct { Metadata map[string]string `protobuf:"bytes,7,rep,name=metadata,proto3" json:"metadata,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"` CloudInitConfig *CloudInitConfig `protobuf:"bytes,8,opt,name=cloud_init_config,json=cloudInitConfig,proto3" json:"cloud_init_config,omitempty"` // advanced cloud-init settings ManagedVolumes []*ManagedVolumeSpec `protobuf:"bytes,9,rep,name=managed_volumes,json=managedVolumes,proto3" json:"managed_volumes,omitempty"` - unknownFields protoimpl.UnknownFields - sizeCache protoimpl.SizeCache + // Happy-path OS image: catalog name or absolute path to a read-only base + // image. Agent creates a writable qcow2 overlay; base is never mutated. + OsImage string `protobuf:"bytes,10,opt,name=os_image,json=osImage,proto3" json:"os_image,omitempty"` + // Root disk size in GB when synthesizing from os_image (default 10). + DiskGb int64 `protobuf:"varint,11,opt,name=disk_gb,json=diskGb,proto3" json:"disk_gb,omitempty"` + // Optional runtime selector: "libvirt" (default) or "firecracker". + Runtime string `protobuf:"bytes,12,opt,name=runtime,proto3" json:"runtime,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache } func (x *VMSpec) Reset() { @@ -1318,12 +1325,36 @@ func (x *VMSpec) GetManagedVolumes() []*ManagedVolumeSpec { return nil } +func (x *VMSpec) GetOsImage() string { + if x != nil { + return x.OsImage + } + return "" +} + +func (x *VMSpec) GetDiskGb() int64 { + if x != nil { + return x.DiskGb + } + return 0 +} + +func (x *VMSpec) GetRuntime() string { + if x != nil { + return x.Runtime + } + return "" +} + type CloudInitConfig struct { state protoimpl.MessageState `protogen:"open.v1"` UserData string `protobuf:"bytes,1,opt,name=user_data,json=userData,proto3" json:"user_data,omitempty"` // cloud-init user-data script MetaData string `protobuf:"bytes,2,opt,name=meta_data,json=metaData,proto3" json:"meta_data,omitempty"` // cloud-init meta-data (JSON) NetworkConfig string `protobuf:"bytes,3,opt,name=network_config,json=networkConfig,proto3" json:"network_config,omitempty"` // cloud-init network config (YAML) VendorData string `protobuf:"bytes,4,opt,name=vendor_data,json=vendorData,proto3" json:"vendor_data,omitempty"` // cloud-init vendor-data + Username string `protobuf:"bytes,5,opt,name=username,proto3" json:"username,omitempty"` // default login user when generating user-data + SshPublicKey string `protobuf:"bytes,6,opt,name=ssh_public_key,json=sshPublicKey,proto3" json:"ssh_public_key,omitempty"` // inject authorized key instead of password + Password string `protobuf:"bytes,7,opt,name=password,proto3" json:"password,omitempty"` // fixed password (otherwise random) unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -1386,6 +1417,27 @@ func (x *CloudInitConfig) GetVendorData() string { return "" } +func (x *CloudInitConfig) GetUsername() string { + if x != nil { + return x.Username + } + return "" +} + +func (x *CloudInitConfig) GetSshPublicKey() string { + if x != nil { + return x.SshPublicKey + } + return "" +} + +func (x *CloudInitConfig) GetPassword() string { + if x != nil { + return x.Password + } + return "" +} + type ManagedVolumeSpec struct { state protoimpl.MessageState `protogen:"open.v1"` Name string `protobuf:"bytes,1,opt,name=name,proto3" json:"name,omitempty"` @@ -1720,12 +1772,14 @@ func (x *RestartPolicy) GetMaxRetryCount() int32 { type DiskConfig struct { state protoimpl.MessageState `protogen:"open.v1"` - Path string `protobuf:"bytes,1,opt,name=path,proto3" json:"path,omitempty"` // path to disk image or ISO + Path string `protobuf:"bytes,1,opt,name=path,proto3" json:"path,omitempty"` // path to disk image or ISO (leave empty with os_image set) Device string `protobuf:"bytes,2,opt,name=device,proto3" json:"device,omitempty"` // vda, vdb, etc. Format string `protobuf:"bytes,3,opt,name=format,proto3" json:"format,omitempty"` // qcow2, raw, iso SizeGb int64 `protobuf:"varint,4,opt,name=size_gb,json=sizeGb,proto3" json:"size_gb,omitempty"` - Type string `protobuf:"bytes,5,opt,name=type,proto3" json:"type,omitempty"` // disk or cdrom (for ISO) - Boot bool `protobuf:"varint,6,opt,name=boot,proto3" json:"boot,omitempty"` // true if this is the boot disk/ISO + Type string `protobuf:"bytes,5,opt,name=type,proto3" json:"type,omitempty"` // disk or cdrom (for ISO) + Boot bool `protobuf:"varint,6,opt,name=boot,proto3" json:"boot,omitempty"` // true if this is the boot disk/ISO + BackingFile string `protobuf:"bytes,7,opt,name=backing_file,json=backingFile,proto3" json:"backing_file,omitempty"` // optional explicit backing image for overlay + Storage string `protobuf:"bytes,8,opt,name=storage,proto3" json:"storage,omitempty"` // local|nfs|ceph-rbd hint unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -1802,11 +1856,28 @@ func (x *DiskConfig) GetBoot() bool { return false } +func (x *DiskConfig) GetBackingFile() string { + if x != nil { + return x.BackingFile + } + return "" +} + +func (x *DiskConfig) GetStorage() string { + if x != nil { + return x.Storage + } + return "" +} + type NetworkConfig struct { state protoimpl.MessageState `protogen:"open.v1"` - Network string `protobuf:"bytes,1,opt,name=network,proto3" json:"network,omitempty"` // network name or bridge + Network string `protobuf:"bytes,1,opt,name=network,proto3" json:"network,omitempty"` // libvirt network name or bridge (default: "default") MacAddress string `protobuf:"bytes,2,opt,name=mac_address,json=macAddress,proto3" json:"mac_address,omitempty"` - IpAddress string `protobuf:"bytes,3,opt,name=ip_address,json=ipAddress,proto3" json:"ip_address,omitempty"` // optional static IP + IpAddress string `protobuf:"bytes,3,opt,name=ip_address,json=ipAddress,proto3" json:"ip_address,omitempty"` // optional static guest IP + HostDevName string `protobuf:"bytes,4,opt,name=host_dev_name,json=hostDevName,proto3" json:"host_dev_name,omitempty"` // Firecracker host TAP device + Model string `protobuf:"bytes,5,opt,name=model,proto3" json:"model,omitempty"` // virtio (default) + Bridge string `protobuf:"bytes,6,opt,name=bridge,proto3" json:"bridge,omitempty"` // optional explicit bridge unknownFields protoimpl.UnknownFields sizeCache protoimpl.SizeCache } @@ -1862,6 +1933,27 @@ func (x *NetworkConfig) GetIpAddress() string { return "" } +func (x *NetworkConfig) GetHostDevName() string { + if x != nil { + return x.HostDevName + } + return "" +} + +func (x *NetworkConfig) GetModel() string { + if x != nil { + return x.Model + } + return "" +} + +func (x *NetworkConfig) GetBridge() string { + if x != nil { + return x.Bridge + } + return "" +} + type WorkloadStatus struct { state protoimpl.MessageState `protogen:"open.v1"` Id string `protobuf:"bytes,1,opt,name=id,proto3" json:"id,omitempty"` @@ -2187,7 +2279,7 @@ const file_agent_proto_rawDesc = "" + "\x03env\x18\x03 \x03(\v2%.persys.agent.v1.ComposeSpec.EnvEntryR\x03env\x1a6\n" + "\bEnvEntry\x12\x10\n" + "\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" + - "\x05value\x18\x02 \x01(\tR\x05value:\x028\x01\"\xf8\x03\n" + + "\x05value\x18\x02 \x01(\tR\x05value:\x028\x01\"\xc6\x04\n" + "\x06VMSpec\x12\x12\n" + "\x04name\x18\x01 \x01(\tR\x04name\x12\x14\n" + "\x05vcpus\x18\x02 \x01(\x05R\x05vcpus\x12\x1b\n" + @@ -2198,16 +2290,23 @@ const file_agent_proto_rawDesc = "" + "cloud_init\x18\x06 \x01(\tR\tcloudInit\x12A\n" + "\bmetadata\x18\a \x03(\v2%.persys.agent.v1.VMSpec.MetadataEntryR\bmetadata\x12L\n" + "\x11cloud_init_config\x18\b \x01(\v2 .persys.agent.v1.CloudInitConfigR\x0fcloudInitConfig\x12K\n" + - "\x0fmanaged_volumes\x18\t \x03(\v2\".persys.agent.v1.ManagedVolumeSpecR\x0emanagedVolumes\x1a;\n" + + "\x0fmanaged_volumes\x18\t \x03(\v2\".persys.agent.v1.ManagedVolumeSpecR\x0emanagedVolumes\x12\x19\n" + + "\bos_image\x18\n" + + " \x01(\tR\aosImage\x12\x17\n" + + "\adisk_gb\x18\v \x01(\x03R\x06diskGb\x12\x18\n" + + "\aruntime\x18\f \x01(\tR\aruntime\x1a;\n" + "\rMetadataEntry\x12\x10\n" + "\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" + - "\x05value\x18\x02 \x01(\tR\x05value:\x028\x01\"\x93\x01\n" + + "\x05value\x18\x02 \x01(\tR\x05value:\x028\x01\"\xf1\x01\n" + "\x0fCloudInitConfig\x12\x1b\n" + "\tuser_data\x18\x01 \x01(\tR\buserData\x12\x1b\n" + "\tmeta_data\x18\x02 \x01(\tR\bmetaData\x12%\n" + "\x0enetwork_config\x18\x03 \x01(\tR\rnetworkConfig\x12\x1f\n" + "\vvendor_data\x18\x04 \x01(\tR\n" + - "vendorData\"\xf3\x01\n" + + "vendorData\x12\x1a\n" + + "\busername\x18\x05 \x01(\tR\busername\x12$\n" + + "\x0essh_public_key\x18\x06 \x01(\tR\fsshPublicKey\x12\x1a\n" + + "\bpassword\x18\a \x01(\tR\bpassword\"\xf3\x01\n" + "\x11ManagedVolumeSpec\x12\x12\n" + "\x04name\x18\x01 \x01(\tR\x04name\x12\x16\n" + "\x06driver\x18\x02 \x01(\tR\x06driver\x12\x17\n" + @@ -2234,7 +2333,7 @@ const file_agent_proto_rawDesc = "" + "\x11memory_swap_bytes\x18\x03 \x01(\x03R\x0fmemorySwapBytes\"O\n" + "\rRestartPolicy\x12\x16\n" + "\x06policy\x18\x01 \x01(\tR\x06policy\x12&\n" + - "\x0fmax_retry_count\x18\x02 \x01(\x05R\rmaxRetryCount\"\x91\x01\n" + + "\x0fmax_retry_count\x18\x02 \x01(\x05R\rmaxRetryCount\"\xce\x01\n" + "\n" + "DiskConfig\x12\x12\n" + "\x04path\x18\x01 \x01(\tR\x04path\x12\x16\n" + @@ -2242,13 +2341,18 @@ const file_agent_proto_rawDesc = "" + "\x06format\x18\x03 \x01(\tR\x06format\x12\x17\n" + "\asize_gb\x18\x04 \x01(\x03R\x06sizeGb\x12\x12\n" + "\x04type\x18\x05 \x01(\tR\x04type\x12\x12\n" + - "\x04boot\x18\x06 \x01(\bR\x04boot\"i\n" + + "\x04boot\x18\x06 \x01(\bR\x04boot\x12!\n" + + "\fbacking_file\x18\a \x01(\tR\vbackingFile\x12\x18\n" + + "\astorage\x18\b \x01(\tR\astorage\"\xbb\x01\n" + "\rNetworkConfig\x12\x18\n" + "\anetwork\x18\x01 \x01(\tR\anetwork\x12\x1f\n" + "\vmac_address\x18\x02 \x01(\tR\n" + "macAddress\x12\x1d\n" + "\n" + - "ip_address\x18\x03 \x01(\tR\tipAddress\"\x97\x04\n" + + "ip_address\x18\x03 \x01(\tR\tipAddress\x12\"\n" + + "\rhost_dev_name\x18\x04 \x01(\tR\vhostDevName\x12\x14\n" + + "\x05model\x18\x05 \x01(\tR\x05model\x12\x16\n" + + "\x06bridge\x18\x06 \x01(\tR\x06bridge\"\x97\x04\n" + "\x0eWorkloadStatus\x12\x0e\n" + "\x02id\x18\x01 \x01(\tR\x02id\x121\n" + "\x04type\x18\x02 \x01(\x0e2\x1d.persys.agent.v1.WorkloadTypeR\x04type\x12\x1f\n" + diff --git a/pkg/scheduler/controlv1/control.pb.go b/pkg/scheduler/controlv1/control.pb.go index a2d0fc4..8d4807d 100644 --- a/pkg/scheduler/controlv1/control.pb.go +++ b/pkg/scheduler/controlv1/control.pb.go @@ -3483,6 +3483,279 @@ func (x *ListWorkloadsRequest) GetStatus() string { return "" } +type SchedulerEventView struct { + state protoimpl.MessageState `protogen:"open.v1"` + Id string `protobuf:"bytes,1,opt,name=id,proto3" json:"id,omitempty"` + Type string `protobuf:"bytes,2,opt,name=type,proto3" json:"type,omitempty"` // e.g. "NodeLost", "WorkloadScheduled", "DriftDetected" + WorkloadId string `protobuf:"bytes,3,opt,name=workload_id,json=workloadId,proto3" json:"workload_id,omitempty"` // optional, empty if not workload-scoped + NodeId string `protobuf:"bytes,4,opt,name=node_id,json=nodeId,proto3" json:"node_id,omitempty"` // optional, empty if not node-scoped + Reason string `protobuf:"bytes,5,opt,name=reason,proto3" json:"reason,omitempty"` + Timestamp *timestamppb.Timestamp `protobuf:"bytes,6,opt,name=timestamp,proto3" json:"timestamp,omitempty"` + // Free-form auxiliary data. Values are stringified on the way out + // (models.SchedulerEvent.Details is map[string]interface{} on the Go + // side) — this is a deliberate simplification over a + // google.protobuf.Struct, since event details are informational/ + // display-oriented, not structured data a client needs to + // round-trip losslessly. + Details map[string]string `protobuf:"bytes,7,rep,name=details,proto3" json:"details,omitempty" protobuf_key:"bytes,1,opt,name=key" protobuf_val:"bytes,2,opt,name=value"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *SchedulerEventView) Reset() { + *x = SchedulerEventView{} + mi := &file_control_proto_msgTypes[49] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *SchedulerEventView) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*SchedulerEventView) ProtoMessage() {} + +func (x *SchedulerEventView) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[49] + 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 SchedulerEventView.ProtoReflect.Descriptor instead. +func (*SchedulerEventView) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{49} +} + +func (x *SchedulerEventView) GetId() string { + if x != nil { + return x.Id + } + return "" +} + +func (x *SchedulerEventView) GetType() string { + if x != nil { + return x.Type + } + return "" +} + +func (x *SchedulerEventView) GetWorkloadId() string { + if x != nil { + return x.WorkloadId + } + return "" +} + +func (x *SchedulerEventView) GetNodeId() string { + if x != nil { + return x.NodeId + } + return "" +} + +func (x *SchedulerEventView) GetReason() string { + if x != nil { + return x.Reason + } + return "" +} + +func (x *SchedulerEventView) GetTimestamp() *timestamppb.Timestamp { + if x != nil { + return x.Timestamp + } + return nil +} + +func (x *SchedulerEventView) GetDetails() map[string]string { + if x != nil { + return x.Details + } + return nil +} + +type ListEventsRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + Limit int64 `protobuf:"varint,1,opt,name=limit,proto3" json:"limit,omitempty"` // 0 means server default + Type string `protobuf:"bytes,2,opt,name=type,proto3" json:"type,omitempty"` // optional filter + WorkloadId string `protobuf:"bytes,3,opt,name=workload_id,json=workloadId,proto3" json:"workload_id,omitempty"` // optional filter + NodeId string `protobuf:"bytes,4,opt,name=node_id,json=nodeId,proto3" json:"node_id,omitempty"` // optional filter + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ListEventsRequest) Reset() { + *x = ListEventsRequest{} + mi := &file_control_proto_msgTypes[50] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ListEventsRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ListEventsRequest) ProtoMessage() {} + +func (x *ListEventsRequest) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[50] + 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 ListEventsRequest.ProtoReflect.Descriptor instead. +func (*ListEventsRequest) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{50} +} + +func (x *ListEventsRequest) GetLimit() int64 { + if x != nil { + return x.Limit + } + return 0 +} + +func (x *ListEventsRequest) GetType() string { + if x != nil { + return x.Type + } + return "" +} + +func (x *ListEventsRequest) GetWorkloadId() string { + if x != nil { + return x.WorkloadId + } + return "" +} + +func (x *ListEventsRequest) GetNodeId() string { + if x != nil { + return x.NodeId + } + return "" +} + +type ListEventsResponse struct { + state protoimpl.MessageState `protogen:"open.v1"` + Events []*SchedulerEventView `protobuf:"bytes,1,rep,name=events,proto3" json:"events,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *ListEventsResponse) Reset() { + *x = ListEventsResponse{} + mi := &file_control_proto_msgTypes[51] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *ListEventsResponse) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*ListEventsResponse) ProtoMessage() {} + +func (x *ListEventsResponse) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[51] + 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 ListEventsResponse.ProtoReflect.Descriptor instead. +func (*ListEventsResponse) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{51} +} + +func (x *ListEventsResponse) GetEvents() []*SchedulerEventView { + if x != nil { + return x.Events + } + return nil +} + +type WatchEventsRequest struct { + state protoimpl.MessageState `protogen:"open.v1"` + // Same optional filters as ListEventsRequest. The stream first replays + // recent matching events (server-side default limit), then continues + // with new matching events as they're emitted. + Type string `protobuf:"bytes,1,opt,name=type,proto3" json:"type,omitempty"` + WorkloadId string `protobuf:"bytes,2,opt,name=workload_id,json=workloadId,proto3" json:"workload_id,omitempty"` + NodeId string `protobuf:"bytes,3,opt,name=node_id,json=nodeId,proto3" json:"node_id,omitempty"` + unknownFields protoimpl.UnknownFields + sizeCache protoimpl.SizeCache +} + +func (x *WatchEventsRequest) Reset() { + *x = WatchEventsRequest{} + mi := &file_control_proto_msgTypes[52] + ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) + ms.StoreMessageInfo(mi) +} + +func (x *WatchEventsRequest) String() string { + return protoimpl.X.MessageStringOf(x) +} + +func (*WatchEventsRequest) ProtoMessage() {} + +func (x *WatchEventsRequest) ProtoReflect() protoreflect.Message { + mi := &file_control_proto_msgTypes[52] + 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 WatchEventsRequest.ProtoReflect.Descriptor instead. +func (*WatchEventsRequest) Descriptor() ([]byte, []int) { + return file_control_proto_rawDescGZIP(), []int{52} +} + +func (x *WatchEventsRequest) GetType() string { + if x != nil { + return x.Type + } + return "" +} + +func (x *WatchEventsRequest) GetWorkloadId() string { + if x != nil { + return x.WorkloadId + } + return "" +} + +func (x *WatchEventsRequest) GetNodeId() string { + if x != nil { + return x.NodeId + } + return "" +} + type GetWorkloadRequest struct { state protoimpl.MessageState `protogen:"open.v1"` WorkloadId string `protobuf:"bytes,1,opt,name=workload_id,json=workloadId,proto3" json:"workload_id,omitempty"` @@ -3492,7 +3765,7 @@ type GetWorkloadRequest struct { func (x *GetWorkloadRequest) Reset() { *x = GetWorkloadRequest{} - mi := &file_control_proto_msgTypes[49] + mi := &file_control_proto_msgTypes[53] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3504,7 +3777,7 @@ func (x *GetWorkloadRequest) String() string { func (*GetWorkloadRequest) ProtoMessage() {} func (x *GetWorkloadRequest) ProtoReflect() protoreflect.Message { - mi := &file_control_proto_msgTypes[49] + mi := &file_control_proto_msgTypes[53] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -3517,7 +3790,7 @@ func (x *GetWorkloadRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use GetWorkloadRequest.ProtoReflect.Descriptor instead. func (*GetWorkloadRequest) Descriptor() ([]byte, []int) { - return file_control_proto_rawDescGZIP(), []int{49} + return file_control_proto_rawDescGZIP(), []int{53} } func (x *GetWorkloadRequest) GetWorkloadId() string { @@ -3536,7 +3809,7 @@ type ListWorkloadsResponse struct { func (x *ListWorkloadsResponse) Reset() { *x = ListWorkloadsResponse{} - mi := &file_control_proto_msgTypes[50] + mi := &file_control_proto_msgTypes[54] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3548,7 +3821,7 @@ func (x *ListWorkloadsResponse) String() string { func (*ListWorkloadsResponse) ProtoMessage() {} func (x *ListWorkloadsResponse) ProtoReflect() protoreflect.Message { - mi := &file_control_proto_msgTypes[50] + mi := &file_control_proto_msgTypes[54] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -3561,7 +3834,7 @@ func (x *ListWorkloadsResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use ListWorkloadsResponse.ProtoReflect.Descriptor instead. func (*ListWorkloadsResponse) Descriptor() ([]byte, []int) { - return file_control_proto_rawDescGZIP(), []int{50} + return file_control_proto_rawDescGZIP(), []int{54} } func (x *ListWorkloadsResponse) GetWorkloads() []*WorkloadView { @@ -3580,7 +3853,7 @@ type GetWorkloadResponse struct { func (x *GetWorkloadResponse) Reset() { *x = GetWorkloadResponse{} - mi := &file_control_proto_msgTypes[51] + mi := &file_control_proto_msgTypes[55] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3592,7 +3865,7 @@ func (x *GetWorkloadResponse) String() string { func (*GetWorkloadResponse) ProtoMessage() {} func (x *GetWorkloadResponse) ProtoReflect() protoreflect.Message { - mi := &file_control_proto_msgTypes[51] + mi := &file_control_proto_msgTypes[55] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -3605,7 +3878,7 @@ func (x *GetWorkloadResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use GetWorkloadResponse.ProtoReflect.Descriptor instead. func (*GetWorkloadResponse) Descriptor() ([]byte, []int) { - return file_control_proto_rawDescGZIP(), []int{51} + return file_control_proto_rawDescGZIP(), []int{55} } func (x *GetWorkloadResponse) GetWorkload() *WorkloadView { @@ -3637,7 +3910,7 @@ type WorkloadView struct { func (x *WorkloadView) Reset() { *x = WorkloadView{} - mi := &file_control_proto_msgTypes[52] + mi := &file_control_proto_msgTypes[56] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3649,7 +3922,7 @@ func (x *WorkloadView) String() string { func (*WorkloadView) ProtoMessage() {} func (x *WorkloadView) ProtoReflect() protoreflect.Message { - mi := &file_control_proto_msgTypes[52] + mi := &file_control_proto_msgTypes[56] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -3662,7 +3935,7 @@ func (x *WorkloadView) ProtoReflect() protoreflect.Message { // Deprecated: Use WorkloadView.ProtoReflect.Descriptor instead. func (*WorkloadView) Descriptor() ([]byte, []int) { - return file_control_proto_rawDescGZIP(), []int{52} + return file_control_proto_rawDescGZIP(), []int{56} } func (x *WorkloadView) GetWorkloadId() string { @@ -3771,7 +4044,7 @@ type GetClusterSummaryRequest struct { func (x *GetClusterSummaryRequest) Reset() { *x = GetClusterSummaryRequest{} - mi := &file_control_proto_msgTypes[53] + mi := &file_control_proto_msgTypes[57] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3783,7 +4056,7 @@ func (x *GetClusterSummaryRequest) String() string { func (*GetClusterSummaryRequest) ProtoMessage() {} func (x *GetClusterSummaryRequest) ProtoReflect() protoreflect.Message { - mi := &file_control_proto_msgTypes[53] + mi := &file_control_proto_msgTypes[57] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -3796,7 +4069,7 @@ func (x *GetClusterSummaryRequest) ProtoReflect() protoreflect.Message { // Deprecated: Use GetClusterSummaryRequest.ProtoReflect.Descriptor instead. func (*GetClusterSummaryRequest) Descriptor() ([]byte, []int) { - return file_control_proto_rawDescGZIP(), []int{53} + return file_control_proto_rawDescGZIP(), []int{57} } type GetClusterSummaryResponse struct { @@ -3816,7 +4089,7 @@ type GetClusterSummaryResponse struct { func (x *GetClusterSummaryResponse) Reset() { *x = GetClusterSummaryResponse{} - mi := &file_control_proto_msgTypes[54] + mi := &file_control_proto_msgTypes[58] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3828,7 +4101,7 @@ func (x *GetClusterSummaryResponse) String() string { func (*GetClusterSummaryResponse) ProtoMessage() {} func (x *GetClusterSummaryResponse) ProtoReflect() protoreflect.Message { - mi := &file_control_proto_msgTypes[54] + mi := &file_control_proto_msgTypes[58] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -3841,7 +4114,7 @@ func (x *GetClusterSummaryResponse) ProtoReflect() protoreflect.Message { // Deprecated: Use GetClusterSummaryResponse.ProtoReflect.Descriptor instead. func (*GetClusterSummaryResponse) Descriptor() ([]byte, []int) { - return file_control_proto_rawDescGZIP(), []int{54} + return file_control_proto_rawDescGZIP(), []int{58} } func (x *GetClusterSummaryResponse) GetTotalNodes() int32 { @@ -3922,7 +4195,7 @@ type ControlMessage struct { func (x *ControlMessage) Reset() { *x = ControlMessage{} - mi := &file_control_proto_msgTypes[55] + mi := &file_control_proto_msgTypes[59] ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) ms.StoreMessageInfo(mi) } @@ -3934,7 +4207,7 @@ func (x *ControlMessage) String() string { func (*ControlMessage) ProtoMessage() {} func (x *ControlMessage) ProtoReflect() protoreflect.Message { - mi := &file_control_proto_msgTypes[55] + mi := &file_control_proto_msgTypes[59] if x != nil { ms := protoimpl.X.MessageStateOf(protoimpl.Pointer(x)) if ms.LoadMessageInfo() == nil { @@ -3947,7 +4220,7 @@ func (x *ControlMessage) ProtoReflect() protoreflect.Message { // Deprecated: Use ControlMessage.ProtoReflect.Descriptor instead. func (*ControlMessage) Descriptor() ([]byte, []int) { - return file_control_proto_rawDescGZIP(), []int{55} + return file_control_proto_rawDescGZIP(), []int{59} } func (x *ControlMessage) GetMessage() isControlMessage_Message { @@ -4315,7 +4588,32 @@ const file_control_proto_rawDesc = "" + "\x05value\x18\x02 \x01(\tR\x05value:\x028\x01\"G\n" + "\x14ListWorkloadsRequest\x12\x17\n" + "\anode_id\x18\x01 \x01(\tR\x06nodeId\x12\x16\n" + - "\x06status\x18\x02 \x01(\tR\x06status\"5\n" + + "\x06status\x18\x02 \x01(\tR\x06status\"\xce\x02\n" + + "\x12SchedulerEventView\x12\x0e\n" + + "\x02id\x18\x01 \x01(\tR\x02id\x12\x12\n" + + "\x04type\x18\x02 \x01(\tR\x04type\x12\x1f\n" + + "\vworkload_id\x18\x03 \x01(\tR\n" + + "workloadId\x12\x17\n" + + "\anode_id\x18\x04 \x01(\tR\x06nodeId\x12\x16\n" + + "\x06reason\x18\x05 \x01(\tR\x06reason\x128\n" + + "\ttimestamp\x18\x06 \x01(\v2\x1a.google.protobuf.TimestampR\ttimestamp\x12L\n" + + "\adetails\x18\a \x03(\v22.persys.control.v1.SchedulerEventView.DetailsEntryR\adetails\x1a:\n" + + "\fDetailsEntry\x12\x10\n" + + "\x03key\x18\x01 \x01(\tR\x03key\x12\x14\n" + + "\x05value\x18\x02 \x01(\tR\x05value:\x028\x01\"w\n" + + "\x11ListEventsRequest\x12\x14\n" + + "\x05limit\x18\x01 \x01(\x03R\x05limit\x12\x12\n" + + "\x04type\x18\x02 \x01(\tR\x04type\x12\x1f\n" + + "\vworkload_id\x18\x03 \x01(\tR\n" + + "workloadId\x12\x17\n" + + "\anode_id\x18\x04 \x01(\tR\x06nodeId\"S\n" + + "\x12ListEventsResponse\x12=\n" + + "\x06events\x18\x01 \x03(\v2%.persys.control.v1.SchedulerEventViewR\x06events\"b\n" + + "\x12WatchEventsRequest\x12\x12\n" + + "\x04type\x18\x01 \x01(\tR\x04type\x12\x1f\n" + + "\vworkload_id\x18\x02 \x01(\tR\n" + + "workloadId\x12\x17\n" + + "\anode_id\x18\x03 \x01(\tR\x06nodeId\"5\n" + "\x12GetWorkloadRequest\x12\x1f\n" + "\vworkload_id\x18\x01 \x01(\tR\n" + "workloadId\"V\n" + @@ -4376,7 +4674,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\xf0\r\n" + + "\x0eVM_BOOT_FAILED\x10\b2\xaa\x0f\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" + @@ -4395,6 +4693,9 @@ const file_control_proto_rawDesc = "" + "\rListWorkloads\x12'.persys.control.v1.ListWorkloadsRequest\x1a(.persys.control.v1.ListWorkloadsResponse\x12\\\n" + "\vGetWorkload\x12%.persys.control.v1.GetWorkloadRequest\x1a&.persys.control.v1.GetWorkloadResponse\x12n\n" + "\x11GetClusterSummary\x12+.persys.control.v1.GetClusterSummaryRequest\x1a,.persys.control.v1.GetClusterSummaryResponse\x12Y\n" + + "\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" var ( @@ -4410,7 +4711,7 @@ func file_control_proto_rawDescGZIP() []byte { } var file_control_proto_enumTypes = make([]protoimpl.EnumInfo, 2) -var file_control_proto_msgTypes = make([]protoimpl.MessageInfo, 61) +var file_control_proto_msgTypes = make([]protoimpl.MessageInfo, 66) var file_control_proto_goTypes = []any{ (AutomationActionType)(0), // 0: persys.control.v1.AutomationActionType (FailureReason)(0), // 1: persys.control.v1.FailureReason @@ -4463,56 +4764,61 @@ var file_control_proto_goTypes = []any{ (*GetNodeResponse)(nil), // 48: persys.control.v1.GetNodeResponse (*NodeView)(nil), // 49: persys.control.v1.NodeView (*ListWorkloadsRequest)(nil), // 50: persys.control.v1.ListWorkloadsRequest - (*GetWorkloadRequest)(nil), // 51: persys.control.v1.GetWorkloadRequest - (*ListWorkloadsResponse)(nil), // 52: persys.control.v1.ListWorkloadsResponse - (*GetWorkloadResponse)(nil), // 53: persys.control.v1.GetWorkloadResponse - (*WorkloadView)(nil), // 54: persys.control.v1.WorkloadView - (*GetClusterSummaryRequest)(nil), // 55: persys.control.v1.GetClusterSummaryRequest - (*GetClusterSummaryResponse)(nil), // 56: persys.control.v1.GetClusterSummaryResponse - (*ControlMessage)(nil), // 57: persys.control.v1.ControlMessage - nil, // 58: persys.control.v1.RegisterNodeRequest.LabelsEntry - nil, // 59: persys.control.v1.WorkloadSpec.MetadataEntry - nil, // 60: persys.control.v1.ContainerSpec.EnvEntry - nil, // 61: persys.control.v1.ComposeSpec.EnvEntry - nil, // 62: persys.control.v1.NodeView.LabelsEntry - (*timestamppb.Timestamp)(nil), // 63: google.protobuf.Timestamp + (*SchedulerEventView)(nil), // 51: persys.control.v1.SchedulerEventView + (*ListEventsRequest)(nil), // 52: persys.control.v1.ListEventsRequest + (*ListEventsResponse)(nil), // 53: persys.control.v1.ListEventsResponse + (*WatchEventsRequest)(nil), // 54: persys.control.v1.WatchEventsRequest + (*GetWorkloadRequest)(nil), // 55: persys.control.v1.GetWorkloadRequest + (*ListWorkloadsResponse)(nil), // 56: persys.control.v1.ListWorkloadsResponse + (*GetWorkloadResponse)(nil), // 57: persys.control.v1.GetWorkloadResponse + (*WorkloadView)(nil), // 58: persys.control.v1.WorkloadView + (*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 } var file_control_proto_depIdxs = []int32{ 0, // 0: persys.control.v1.AutomationSuggestion.action_type:type_name -> persys.control.v1.AutomationActionType - 63, // 1: persys.control.v1.AutomationSuggestion.suggested_at:type_name -> google.protobuf.Timestamp + 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 - 63, // 3: persys.control.v1.SubmitAutomationSuggestionResponse.decided_at:type_name -> google.protobuf.Timestamp + 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 - 58, // 5: persys.control.v1.RegisterNodeRequest.labels:type_name -> persys.control.v1.RegisterNodeRequest.LabelsEntry - 63, // 6: persys.control.v1.RegisterNodeRequest.timestamp:type_name -> google.protobuf.Timestamp + 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 - 63, // 8: persys.control.v1.RegisterNodeResponse.lease_expires_at:type_name -> google.protobuf.Timestamp + 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 - 63, // 11: persys.control.v1.HeartbeatRequest.timestamp:type_name -> google.protobuf.Timestamp + 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 - 63, // 13: persys.control.v1.HeartbeatResponse.lease_expires_at:type_name -> google.protobuf.Timestamp + 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 - 59, // 20: persys.control.v1.WorkloadSpec.metadata:type_name -> persys.control.v1.WorkloadSpec.MetadataEntry - 60, // 21: persys.control.v1.ContainerSpec.env:type_name -> persys.control.v1.ContainerSpec.EnvEntry + 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 - 61, // 25: persys.control.v1.ComposeSpec.env:type_name -> persys.control.v1.ComposeSpec.EnvEntry + 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 - 63, // 30: persys.control.v1.WorkloadUsageSnapshot.collected_at:type_name -> google.protobuf.Timestamp - 63, // 31: persys.control.v1.ReasonDetail.last_transition:type_name -> google.protobuf.Timestamp - 63, // 32: persys.control.v1.ReasonDetail.next_retry_at:type_name -> google.protobuf.Timestamp + 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 - 63, // 34: persys.control.v1.WorkloadStatus.last_transition:type_name -> google.protobuf.Timestamp + 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 @@ -4524,63 +4830,70 @@ var file_control_proto_depIdxs = []int32{ 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 - 63, // 46: persys.control.v1.NodeView.status_updated_at:type_name -> google.protobuf.Timestamp - 63, // 47: persys.control.v1.NodeView.last_heartbeat:type_name -> google.protobuf.Timestamp - 62, // 48: persys.control.v1.NodeView.labels:type_name -> persys.control.v1.NodeView.LabelsEntry + 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 - 54, // 50: persys.control.v1.ListWorkloadsResponse.workloads:type_name -> persys.control.v1.WorkloadView - 54, // 51: persys.control.v1.GetWorkloadResponse.workload:type_name -> persys.control.v1.WorkloadView - 63, // 52: persys.control.v1.WorkloadView.retry_next_at:type_name -> google.protobuf.Timestamp - 63, // 53: persys.control.v1.WorkloadView.last_updated:type_name -> google.protobuf.Timestamp - 28, // 54: persys.control.v1.WorkloadView.reason:type_name -> persys.control.v1.ReasonDetail - 27, // 55: persys.control.v1.WorkloadView.usage:type_name -> persys.control.v1.WorkloadUsageSnapshot - 63, // 56: persys.control.v1.WorkloadView.created_at:type_name -> google.protobuf.Timestamp - 63, // 57: persys.control.v1.GetClusterSummaryResponse.generated_at:type_name -> google.protobuf.Timestamp - 5, // 58: persys.control.v1.ControlMessage.register:type_name -> persys.control.v1.RegisterNodeRequest - 9, // 59: persys.control.v1.ControlMessage.heartbeat:type_name -> persys.control.v1.HeartbeatRequest - 12, // 60: persys.control.v1.ControlMessage.apply:type_name -> persys.control.v1.ApplyWorkloadRequest - 14, // 61: persys.control.v1.ControlMessage.delete:type_name -> persys.control.v1.DeleteWorkloadRequest - 5, // 62: persys.control.v1.AgentControl.RegisterNode:input_type -> persys.control.v1.RegisterNodeRequest - 9, // 63: persys.control.v1.AgentControl.Heartbeat:input_type -> persys.control.v1.HeartbeatRequest - 12, // 64: persys.control.v1.AgentControl.ApplyWorkload:input_type -> persys.control.v1.ApplyWorkloadRequest - 14, // 65: persys.control.v1.AgentControl.DeleteWorkload:input_type -> persys.control.v1.DeleteWorkloadRequest - 30, // 66: persys.control.v1.AgentControl.RetryWorkload:input_type -> persys.control.v1.RetryWorkloadRequest - 32, // 67: persys.control.v1.AgentControl.DrainNode:input_type -> persys.control.v1.DrainNodeRequest - 34, // 68: persys.control.v1.AgentControl.UndrainNode:input_type -> persys.control.v1.UndrainNodeRequest - 37, // 69: persys.control.v1.AgentControl.TaintNode:input_type -> persys.control.v1.TaintNodeRequest - 39, // 70: persys.control.v1.AgentControl.UntaintNode:input_type -> persys.control.v1.UntaintNodeRequest - 41, // 71: persys.control.v1.AgentControl.SetNodeLabel:input_type -> persys.control.v1.SetNodeLabelRequest - 43, // 72: persys.control.v1.AgentControl.DeleteNodeLabel:input_type -> persys.control.v1.DeleteNodeLabelRequest - 3, // 73: persys.control.v1.AgentControl.SubmitAutomationSuggestion:input_type -> persys.control.v1.SubmitAutomationSuggestionRequest - 45, // 74: persys.control.v1.AgentControl.ListNodes:input_type -> persys.control.v1.ListNodesRequest - 46, // 75: persys.control.v1.AgentControl.GetNode:input_type -> persys.control.v1.GetNodeRequest - 50, // 76: persys.control.v1.AgentControl.ListWorkloads:input_type -> persys.control.v1.ListWorkloadsRequest - 51, // 77: persys.control.v1.AgentControl.GetWorkload:input_type -> persys.control.v1.GetWorkloadRequest - 55, // 78: persys.control.v1.AgentControl.GetClusterSummary:input_type -> persys.control.v1.GetClusterSummaryRequest - 57, // 79: persys.control.v1.AgentControl.ControlStream:input_type -> persys.control.v1.ControlMessage - 8, // 80: persys.control.v1.AgentControl.RegisterNode:output_type -> persys.control.v1.RegisterNodeResponse - 11, // 81: persys.control.v1.AgentControl.Heartbeat:output_type -> persys.control.v1.HeartbeatResponse - 13, // 82: persys.control.v1.AgentControl.ApplyWorkload:output_type -> persys.control.v1.ApplyWorkloadResponse - 15, // 83: persys.control.v1.AgentControl.DeleteWorkload:output_type -> persys.control.v1.DeleteWorkloadResponse - 31, // 84: persys.control.v1.AgentControl.RetryWorkload:output_type -> persys.control.v1.RetryWorkloadResponse - 33, // 85: persys.control.v1.AgentControl.DrainNode:output_type -> persys.control.v1.DrainNodeResponse - 35, // 86: persys.control.v1.AgentControl.UndrainNode:output_type -> persys.control.v1.UndrainNodeResponse - 38, // 87: persys.control.v1.AgentControl.TaintNode:output_type -> persys.control.v1.TaintNodeResponse - 40, // 88: persys.control.v1.AgentControl.UntaintNode:output_type -> persys.control.v1.UntaintNodeResponse - 42, // 89: persys.control.v1.AgentControl.SetNodeLabel:output_type -> persys.control.v1.SetNodeLabelResponse - 44, // 90: persys.control.v1.AgentControl.DeleteNodeLabel:output_type -> persys.control.v1.DeleteNodeLabelResponse - 4, // 91: persys.control.v1.AgentControl.SubmitAutomationSuggestion:output_type -> persys.control.v1.SubmitAutomationSuggestionResponse - 47, // 92: persys.control.v1.AgentControl.ListNodes:output_type -> persys.control.v1.ListNodesResponse - 48, // 93: persys.control.v1.AgentControl.GetNode:output_type -> persys.control.v1.GetNodeResponse - 52, // 94: persys.control.v1.AgentControl.ListWorkloads:output_type -> persys.control.v1.ListWorkloadsResponse - 53, // 95: persys.control.v1.AgentControl.GetWorkload:output_type -> persys.control.v1.GetWorkloadResponse - 56, // 96: persys.control.v1.AgentControl.GetClusterSummary:output_type -> persys.control.v1.GetClusterSummaryResponse - 57, // 97: persys.control.v1.AgentControl.ControlStream:output_type -> persys.control.v1.ControlMessage - 80, // [80:98] is the sub-list for method output_type - 62, // [62:80] is the sub-list for method input_type - 62, // [62:62] is the sub-list for extension type_name - 62, // [62:62] is the sub-list for extension extendee - 0, // [0:62] is the sub-list for field type_name + 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 } func init() { file_control_proto_init() } @@ -4593,7 +4906,7 @@ func file_control_proto_init() { (*WorkloadSpec_Compose)(nil), (*WorkloadSpec_Vm)(nil), } - file_control_proto_msgTypes[55].OneofWrappers = []any{ + file_control_proto_msgTypes[59].OneofWrappers = []any{ (*ControlMessage_Register)(nil), (*ControlMessage_Heartbeat)(nil), (*ControlMessage_Apply)(nil), @@ -4605,7 +4918,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: 61, + NumMessages: 66, NumExtensions: 0, NumServices: 1, }, diff --git a/pkg/scheduler/controlv1/control_grpc.pb.go b/pkg/scheduler/controlv1/control_grpc.pb.go index 4a70e36..4094860 100644 --- a/pkg/scheduler/controlv1/control_grpc.pb.go +++ b/pkg/scheduler/controlv1/control_grpc.pb.go @@ -36,6 +36,8 @@ const ( AgentControl_ListWorkloads_FullMethodName = "/persys.control.v1.AgentControl/ListWorkloads" AgentControl_GetWorkload_FullMethodName = "/persys.control.v1.AgentControl/GetWorkload" AgentControl_GetClusterSummary_FullMethodName = "/persys.control.v1.AgentControl/GetClusterSummary" + AgentControl_ListEvents_FullMethodName = "/persys.control.v1.AgentControl/ListEvents" + AgentControl_WatchEvents_FullMethodName = "/persys.control.v1.AgentControl/WatchEvents" AgentControl_ControlStream_FullMethodName = "/persys.control.v1.AgentControl/ControlStream" ) @@ -66,6 +68,18 @@ 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 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 + // streaming RPCs, so consumers that need HTTP (e.g. a browser + // dashboard) go through a hand-written SSE endpoint on the gateway + // instead of the generic bridge; a gRPC client (e.g. persysctl) can + // call it directly. + ListEvents(ctx context.Context, in *ListEventsRequest, opts ...grpc.CallOption) (*ListEventsResponse, error) + 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) } @@ -248,9 +262,38 @@ func (c *agentControlClient) GetClusterSummary(ctx context.Context, in *GetClust return out, nil } +func (c *agentControlClient) ListEvents(ctx context.Context, in *ListEventsRequest, opts ...grpc.CallOption) (*ListEventsResponse, error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + out := new(ListEventsResponse) + err := c.cc.Invoke(ctx, AgentControl_ListEvents_FullMethodName, in, out, cOpts...) + if err != nil { + return nil, err + } + return out, nil +} + +func (c *agentControlClient) WatchEvents(ctx context.Context, in *WatchEventsRequest, opts ...grpc.CallOption) (grpc.ServerStreamingClient[SchedulerEventView], error) { + cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) + stream, err := c.cc.NewStream(ctx, &AgentControl_ServiceDesc.Streams[0], AgentControl_WatchEvents_FullMethodName, cOpts...) + if err != nil { + return nil, err + } + x := &grpc.GenericClientStream[WatchEventsRequest, SchedulerEventView]{ClientStream: stream} + if err := x.ClientStream.SendMsg(in); err != nil { + return nil, err + } + if err := x.ClientStream.CloseSend(); err != nil { + return nil, err + } + return x, nil +} + +// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name. +type AgentControl_WatchEventsClient = grpc.ServerStreamingClient[SchedulerEventView] + func (c *agentControlClient) ControlStream(ctx context.Context, opts ...grpc.CallOption) (grpc.BidiStreamingClient[ControlMessage, ControlMessage], error) { cOpts := append([]grpc.CallOption{grpc.StaticMethod()}, opts...) - stream, err := c.cc.NewStream(ctx, &AgentControl_ServiceDesc.Streams[0], AgentControl_ControlStream_FullMethodName, cOpts...) + stream, err := c.cc.NewStream(ctx, &AgentControl_ServiceDesc.Streams[1], AgentControl_ControlStream_FullMethodName, cOpts...) if err != nil { return nil, err } @@ -288,6 +331,18 @@ 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 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 + // streaming RPCs, so consumers that need HTTP (e.g. a browser + // dashboard) go through a hand-written SSE endpoint on the gateway + // instead of the generic bridge; a gRPC client (e.g. persysctl) can + // call it directly. + ListEvents(context.Context, *ListEventsRequest) (*ListEventsResponse, error) + WatchEvents(*WatchEventsRequest, grpc.ServerStreamingServer[SchedulerEventView]) error // Optional future streaming channel ControlStream(grpc.BidiStreamingServer[ControlMessage, ControlMessage]) error mustEmbedUnimplementedAgentControlServer() @@ -351,6 +406,12 @@ func (UnimplementedAgentControlServer) GetWorkload(context.Context, *GetWorkload func (UnimplementedAgentControlServer) GetClusterSummary(context.Context, *GetClusterSummaryRequest) (*GetClusterSummaryResponse, error) { return nil, status.Error(codes.Unimplemented, "method GetClusterSummary not implemented") } +func (UnimplementedAgentControlServer) ListEvents(context.Context, *ListEventsRequest) (*ListEventsResponse, error) { + return nil, status.Error(codes.Unimplemented, "method ListEvents not implemented") +} +func (UnimplementedAgentControlServer) WatchEvents(*WatchEventsRequest, grpc.ServerStreamingServer[SchedulerEventView]) error { + return status.Error(codes.Unimplemented, "method WatchEvents not implemented") +} func (UnimplementedAgentControlServer) ControlStream(grpc.BidiStreamingServer[ControlMessage, ControlMessage]) error { return status.Error(codes.Unimplemented, "method ControlStream not implemented") } @@ -681,6 +742,35 @@ func _AgentControl_GetClusterSummary_Handler(srv interface{}, ctx context.Contex return interceptor(ctx, in, info, handler) } +func _AgentControl_ListEvents_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { + in := new(ListEventsRequest) + if err := dec(in); err != nil { + return nil, err + } + if interceptor == nil { + return srv.(AgentControlServer).ListEvents(ctx, in) + } + info := &grpc.UnaryServerInfo{ + Server: srv, + FullMethod: AgentControl_ListEvents_FullMethodName, + } + handler := func(ctx context.Context, req interface{}) (interface{}, error) { + return srv.(AgentControlServer).ListEvents(ctx, req.(*ListEventsRequest)) + } + return interceptor(ctx, in, info, handler) +} + +func _AgentControl_WatchEvents_Handler(srv interface{}, stream grpc.ServerStream) error { + m := new(WatchEventsRequest) + if err := stream.RecvMsg(m); err != nil { + return err + } + return srv.(AgentControlServer).WatchEvents(m, &grpc.GenericServerStream[WatchEventsRequest, SchedulerEventView]{ServerStream: stream}) +} + +// This type alias is provided for backwards compatibility with existing code that references the prior non-generic stream type by name. +type AgentControl_WatchEventsServer = grpc.ServerStreamingServer[SchedulerEventView] + func _AgentControl_ControlStream_Handler(srv interface{}, stream grpc.ServerStream) error { return srv.(AgentControlServer).ControlStream(&grpc.GenericServerStream[ControlMessage, ControlMessage]{ServerStream: stream}) } @@ -763,8 +853,17 @@ var AgentControl_ServiceDesc = grpc.ServiceDesc{ MethodName: "GetClusterSummary", Handler: _AgentControl_GetClusterSummary_Handler, }, + { + MethodName: "ListEvents", + Handler: _AgentControl_ListEvents_Handler, + }, }, Streams: []grpc.StreamDesc{ + { + StreamName: "WatchEvents", + Handler: _AgentControl_WatchEvents_Handler, + ServerStreams: true, + }, { StreamName: "ControlStream", Handler: _AgentControl_ControlStream_Handler, From 35202d83381ec76f38140274feb61db146819d6b Mon Sep 17 00:00:00 2001 From: milx Date: Fri, 7 Aug 2026 19:30:49 +0330 Subject: [PATCH 6/6] Update: gitignore to ignore .zip + .tar.gz --- .gitignore | 7 ++----- 1 file changed, 2 insertions(+), 5 deletions(-) diff --git a/.gitignore b/.gitignore index 4eb9d7c..2884ee9 100644 --- a/.gitignore +++ b/.gitignore @@ -4,15 +4,12 @@ .idea *.log __pycache__ -.old/azure-pipeline/azure-pipelines.yml .old/** *.bak terraform.tfvars **/bin/ -certs/ -vault-mtls-mock/certs-runtime/ -vault-mtls-mock/ .env persys-automation/DESIGN_SPEC.md persys-intelligence/DESIGN_SPEC.md -third_party/ \ No newline at end of file +.zip +.tar.gz \ No newline at end of file