Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
7 changes: 2 additions & 5 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -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/
.zip
.tar.gz
88 changes: 88 additions & 0 deletions persysctl/cmd/metrics.go
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand All @@ -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")
}
165 changes: 165 additions & 0 deletions persysctl/internal/client/client.go
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
package client

import (
"bufio"
"bytes"
"context"
"crypto/tls"
Expand Down Expand Up @@ -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")
Expand Down
86 changes: 86 additions & 0 deletions persysctl/internal/client/client_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
}
Loading
Loading