Skip to content

Commit 6dfc35c

Browse files
authored
Merge pull request #36 from persys-dev/persysctl/pkg-minor-updates
Persysctl/pkg minor updates
2 parents f63beca + 35202d8 commit 6dfc35c

8 files changed

Lines changed: 982 additions & 130 deletions

File tree

.gitignore

Lines changed: 2 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -4,15 +4,12 @@
44
.idea
55
*.log
66
__pycache__
7-
.old/azure-pipeline/azure-pipelines.yml
87
.old/**
98
*.bak
109
terraform.tfvars
1110
**/bin/
12-
certs/
13-
vault-mtls-mock/certs-runtime/
14-
vault-mtls-mock/
1511
.env
1612
persys-automation/DESIGN_SPEC.md
1713
persys-intelligence/DESIGN_SPEC.md
18-
third_party/
14+
.zip
15+
.tar.gz

persysctl/cmd/metrics.go

Lines changed: 88 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,10 +3,24 @@ package cmd
33
import (
44
"encoding/json"
55
"fmt"
6+
"time"
67

78
"github.com/spf13/cobra"
89
)
910

11+
var (
12+
metricsFrom string
13+
metricsTo string
14+
meterWorkloadID string
15+
meterNodeID string
16+
meterType string
17+
meterLimit int
18+
eventsType string
19+
eventsWorkloadID string
20+
eventsNodeID string
21+
eventsLimit int
22+
)
23+
1024
var metricsCmd = &cobra.Command{
1125
Use: "metrics",
1226
Short: "View Persys Compute metrics",
@@ -23,6 +37,80 @@ var metricsCmd = &cobra.Command{
2337
},
2438
}
2539

40+
var meterCmd = &cobra.Command{
41+
Use: "meter",
42+
Short: "Inspect workload usage via the gateway meter API",
43+
}
44+
45+
var meterSummaryCmd = &cobra.Command{
46+
Use: "summary [workload-id]",
47+
Short: "Get usage summary for a workload",
48+
Args: cobra.ExactArgs(1),
49+
Run: func(cmd *cobra.Command, args []string) {
50+
c, _, err := newClientWithTrace()
51+
cobra.CheckErr(err)
52+
defer c.Close()
53+
from, err := parseTimeFlag(metricsFrom)
54+
cobra.CheckErr(err)
55+
to, err := parseTimeFlag(metricsTo)
56+
cobra.CheckErr(err)
57+
resp, err := c.MeterSummary(args[0], from, to)
58+
cobra.CheckErr(err)
59+
printJSON(resp)
60+
},
61+
}
62+
63+
var meterListCmd = &cobra.Command{
64+
Use: "list",
65+
Short: "List current workload usage samples",
66+
Run: func(cmd *cobra.Command, args []string) {
67+
c, _, err := newClientWithTrace()
68+
cobra.CheckErr(err)
69+
defer c.Close()
70+
resp, err := c.MeterListWorkloads(meterType)
71+
cobra.CheckErr(err)
72+
printJSON(resp)
73+
},
74+
}
75+
76+
var eventsCmd = &cobra.Command{
77+
Use: "events",
78+
Short: "Watch scheduler events through the gateway SSE endpoint",
79+
}
80+
81+
var eventsWatchCmd = &cobra.Command{
82+
Use: "watch",
83+
Short: "Stream scheduler events matching optional filters",
84+
Run: func(cmd *cobra.Command, args []string) {
85+
c, _, err := newClientWithTrace()
86+
cobra.CheckErr(err)
87+
defer c.Close()
88+
resp, err := c.WatchGatewayEvents(eventsType, eventsWorkloadID, eventsNodeID, eventsLimit)
89+
cobra.CheckErr(err)
90+
printJSON(resp)
91+
},
92+
}
93+
94+
func parseTimeFlag(raw string) (time.Time, error) {
95+
if raw == "" {
96+
return time.Time{}, nil
97+
}
98+
return time.Parse(time.RFC3339, raw)
99+
}
100+
26101
func init() {
27102
rootCmd.AddCommand(metricsCmd)
103+
rootCmd.AddCommand(meterCmd)
104+
rootCmd.AddCommand(eventsCmd)
105+
meterCmd.AddCommand(meterSummaryCmd)
106+
meterCmd.AddCommand(meterListCmd)
107+
eventsCmd.AddCommand(eventsWatchCmd)
108+
109+
meterSummaryCmd.Flags().StringVar(&metricsFrom, "from", "", "Start time in RFC3339 format")
110+
meterSummaryCmd.Flags().StringVar(&metricsTo, "to", "", "End time in RFC3339 format")
111+
meterListCmd.Flags().StringVar(&meterType, "type", "", "Workload type filter")
112+
eventsWatchCmd.Flags().StringVar(&eventsType, "type", "", "Event type filter")
113+
eventsWatchCmd.Flags().StringVar(&eventsWorkloadID, "workload-id", "", "Filter by workload ID")
114+
eventsWatchCmd.Flags().StringVar(&eventsNodeID, "node-id", "", "Filter by node ID")
115+
eventsWatchCmd.Flags().IntVar(&eventsLimit, "limit", 20, "Maximum number of events to print before exiting")
28116
}

persysctl/internal/client/client.go

Lines changed: 165 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
package client
22

33
import (
4+
"bufio"
45
"bytes"
56
"context"
67
"crypto/tls"
@@ -1011,6 +1012,170 @@ func (c *Client) GatewayClusters() (*GatewayClustersResponse, error) {
10111012
return &out, nil
10121013
}
10131014

1015+
func (c *Client) MeterListWorkloads(workloadType string) (map[string]any, error) {
1016+
if c.cfg.Transport != "http" {
1017+
return nil, fmt.Errorf("meter API is available only with http transport")
1018+
}
1019+
path := "/meter/v1/workloads"
1020+
q := url.Values{}
1021+
if trimmed := strings.TrimSpace(workloadType); trimmed != "" {
1022+
q.Set("workload_type", trimmed)
1023+
}
1024+
if encoded := q.Encode(); encoded != "" {
1025+
path += "?" + encoded
1026+
}
1027+
var out map[string]any
1028+
if err := c.httpJSONRequest("GET", path, nil, &out); err != nil {
1029+
return nil, err
1030+
}
1031+
return out, nil
1032+
}
1033+
1034+
func (c *Client) MeterGetWorkload(workloadID string) (map[string]any, error) {
1035+
if c.cfg.Transport != "http" {
1036+
return nil, fmt.Errorf("meter API is available only with http transport")
1037+
}
1038+
path := "/meter/v1/workloads/" + url.PathEscape(strings.TrimSpace(workloadID))
1039+
var out map[string]any
1040+
if err := c.httpJSONRequest("GET", path, nil, &out); err != nil {
1041+
return nil, err
1042+
}
1043+
return out, nil
1044+
}
1045+
1046+
func (c *Client) MeterHistory(workloadID string, from, to time.Time, limit int) (map[string]any, error) {
1047+
if c.cfg.Transport != "http" {
1048+
return nil, fmt.Errorf("meter API is available only with http transport")
1049+
}
1050+
path := "/meter/v1/workloads/" + url.PathEscape(strings.TrimSpace(workloadID)) + "/history"
1051+
q := url.Values{}
1052+
if !from.IsZero() {
1053+
q.Set("from", from.UTC().Format(time.RFC3339))
1054+
}
1055+
if !to.IsZero() {
1056+
q.Set("to", to.UTC().Format(time.RFC3339))
1057+
}
1058+
if limit > 0 {
1059+
q.Set("limit", strconv.Itoa(limit))
1060+
}
1061+
if encoded := q.Encode(); encoded != "" {
1062+
path += "?" + encoded
1063+
}
1064+
var out map[string]any
1065+
if err := c.httpJSONRequest("GET", path, nil, &out); err != nil {
1066+
return nil, err
1067+
}
1068+
return out, nil
1069+
}
1070+
1071+
func (c *Client) MeterSummary(workloadID string, from, to time.Time) (map[string]any, error) {
1072+
if c.cfg.Transport != "http" {
1073+
return nil, fmt.Errorf("meter API is available only with http transport")
1074+
}
1075+
path := "/meter/v1/workloads/" + url.PathEscape(strings.TrimSpace(workloadID)) + "/summary"
1076+
q := url.Values{}
1077+
if !from.IsZero() {
1078+
q.Set("from", from.UTC().Format(time.RFC3339))
1079+
}
1080+
if !to.IsZero() {
1081+
q.Set("to", to.UTC().Format(time.RFC3339))
1082+
}
1083+
if encoded := q.Encode(); encoded != "" {
1084+
path += "?" + encoded
1085+
}
1086+
var out map[string]any
1087+
if err := c.httpJSONRequest("GET", path, nil, &out); err != nil {
1088+
return nil, err
1089+
}
1090+
return out, nil
1091+
}
1092+
1093+
func (c *Client) WatchGatewayEvents(eventType, workloadID, nodeID string, limit int) ([]map[string]any, error) {
1094+
if c.cfg.Transport != "http" {
1095+
return nil, fmt.Errorf("gateway event stream is available only with http transport")
1096+
}
1097+
if limit <= 0 {
1098+
limit = 20
1099+
}
1100+
1101+
path := "/events/watch"
1102+
q := url.Values{}
1103+
if trimmed := strings.TrimSpace(eventType); trimmed != "" {
1104+
q.Set("type", trimmed)
1105+
}
1106+
if trimmed := strings.TrimSpace(workloadID); trimmed != "" {
1107+
q.Set("workload_id", trimmed)
1108+
}
1109+
if trimmed := strings.TrimSpace(nodeID); trimmed != "" {
1110+
q.Set("node_id", trimmed)
1111+
}
1112+
if encoded := q.Encode(); encoded != "" {
1113+
path += "?" + encoded
1114+
}
1115+
1116+
timeout := time.Duration(c.cfg.RPCTimeoutSeconds) * time.Second
1117+
if timeout <= 0 {
1118+
timeout = 15 * time.Second
1119+
}
1120+
ctx, cancel := context.WithTimeout(context.Background(), timeout)
1121+
defer cancel()
1122+
1123+
req, err := http.NewRequestWithContext(ctx, http.MethodGet, c.cfg.APIEndpoint+path, nil)
1124+
if err != nil {
1125+
return nil, fmt.Errorf("failed to create request: %w", err)
1126+
}
1127+
req.Header.Set("Accept", "text/event-stream")
1128+
resp, err := c.httpClient.Do(req)
1129+
if err != nil {
1130+
if ctx.Err() == context.DeadlineExceeded {
1131+
return nil, nil
1132+
}
1133+
return nil, fmt.Errorf("failed to send request: %w", err)
1134+
}
1135+
defer resp.Body.Close()
1136+
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
1137+
body, _ := io.ReadAll(resp.Body)
1138+
return nil, fmt.Errorf("API returned status %d: %s", resp.StatusCode, string(body))
1139+
}
1140+
1141+
scanner := bufio.NewScanner(resp.Body)
1142+
scanner.Buffer(make([]byte, 0, 64*1024), 1024*1024)
1143+
var (
1144+
events []map[string]any
1145+
dataLine string
1146+
)
1147+
for scanner.Scan() {
1148+
line := strings.TrimSpace(scanner.Text())
1149+
if strings.HasPrefix(line, "event:") {
1150+
continue
1151+
}
1152+
if strings.HasPrefix(line, "data:") {
1153+
dataLine = strings.TrimSpace(strings.TrimPrefix(line, "data:"))
1154+
continue
1155+
}
1156+
if line == "" && dataLine != "" {
1157+
var event map[string]any
1158+
if err := json.Unmarshal([]byte(dataLine), &event); err == nil {
1159+
events = append(events, event)
1160+
if len(events) >= limit {
1161+
return events, nil
1162+
}
1163+
}
1164+
dataLine = ""
1165+
}
1166+
if ctx.Err() != nil {
1167+
return events, nil
1168+
}
1169+
}
1170+
if err := scanner.Err(); err != nil {
1171+
if ctx.Err() == context.DeadlineExceeded {
1172+
return events, nil
1173+
}
1174+
return nil, fmt.Errorf("failed to read event stream: %w", err)
1175+
}
1176+
return events, nil
1177+
}
1178+
10141179
func (c *Client) TriggerForgeryBuild(req ForgeryBuildTriggerRequest) (map[string]interface{}, error) {
10151180
if c.cfg.Transport != "http" {
10161181
return nil, fmt.Errorf("forgery trigger-build is available only with http transport")
Lines changed: 86 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,86 @@
1+
package client
2+
3+
import (
4+
"net/http"
5+
"net/http/httptest"
6+
"testing"
7+
"time"
8+
9+
"github.com/persys-dev/persysctl/internal/config"
10+
)
11+
12+
func TestMeterSummaryGatewayPath(t *testing.T) {
13+
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
14+
if r.Method != http.MethodGet {
15+
t.Fatalf("unexpected method: %s", r.Method)
16+
}
17+
if r.URL.Path != "/meter/v1/workloads/wl-123/summary" {
18+
t.Fatalf("unexpected path: %s", r.URL.Path)
19+
}
20+
if got := r.URL.Query().Get("from"); got == "" {
21+
t.Fatal("missing from query")
22+
}
23+
w.Header().Set("Content-Type", "application/json")
24+
_, _ = w.Write([]byte(`{"workload_id":"wl-123","sample_count":7}`))
25+
}))
26+
defer server.Close()
27+
28+
cli := &Client{cfg: config.Config{Transport: "http", APIEndpoint: server.URL}, httpClient: server.Client()}
29+
from := time.Date(2026, 7, 1, 0, 0, 0, 0, time.UTC)
30+
to := time.Date(2026, 7, 2, 0, 0, 0, 0, time.UTC)
31+
32+
resp, err := cli.MeterSummary("wl-123", from, to)
33+
if err != nil {
34+
t.Fatalf("MeterSummary() error = %v", err)
35+
}
36+
if resp["workload_id"] != "wl-123" {
37+
t.Fatalf("MeterSummary() workload_id = %v", resp["workload_id"])
38+
}
39+
}
40+
41+
func TestWatchGatewayEventsParsesSSE(t *testing.T) {
42+
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
43+
if r.URL.Path != "/events/watch" {
44+
t.Fatalf("unexpected path: %s", r.URL.Path)
45+
}
46+
w.Header().Set("Content-Type", "text/event-stream")
47+
_, _ = w.Write([]byte("event: event\ndata: {\"type\":\"workload\",\"workload_id\":\"wl-9\"}\n\n"))
48+
}))
49+
defer server.Close()
50+
51+
cli := &Client{cfg: config.Config{Transport: "http", APIEndpoint: server.URL}, httpClient: server.Client()}
52+
events, err := cli.WatchGatewayEvents("workload", "wl-9", "", 10)
53+
if err != nil {
54+
t.Fatalf("WatchGatewayEvents() error = %v", err)
55+
}
56+
if len(events) != 1 {
57+
t.Fatalf("WatchGatewayEvents() len = %d", len(events))
58+
}
59+
if got := events[0]["workload_id"]; got != "wl-9" {
60+
t.Fatalf("WatchGatewayEvents() workload_id = %v", got)
61+
}
62+
}
63+
64+
func TestWatchGatewayEventsTimesOutOnIdleStream(t *testing.T) {
65+
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
66+
w.Header().Set("Content-Type", "text/event-stream")
67+
if f, ok := w.(http.Flusher); ok {
68+
f.Flush()
69+
}
70+
<-r.Context().Done()
71+
}))
72+
defer server.Close()
73+
74+
cli := &Client{cfg: config.Config{Transport: "http", APIEndpoint: server.URL, RPCTimeoutSeconds: 1}, httpClient: server.Client()}
75+
start := time.Now()
76+
events, err := cli.WatchGatewayEvents("", "", "", 20)
77+
if err != nil {
78+
t.Fatalf("WatchGatewayEvents() error = %v", err)
79+
}
80+
if len(events) != 0 {
81+
t.Fatalf("WatchGatewayEvents() idle stream returned %d events", len(events))
82+
}
83+
if elapsed := time.Since(start); elapsed > 3*time.Second {
84+
t.Fatalf("WatchGatewayEvents() took too long: %s", elapsed)
85+
}
86+
}

0 commit comments

Comments
 (0)