From e2c0398ec9262f4677f37aa8ba52015294bc17de Mon Sep 17 00:00:00 2001 From: Goon Date: Fri, 12 Jun 2026 09:48:01 +0700 Subject: [PATCH] feat(cli): add trace operator commands --- README.md | 10 ++ cmd/gateway_http_client.go | 44 ++----- cmd/gateway_http_client_resolution.go | 88 +++++++++++++ cmd/gateway_http_client_test.go | 131 +++++++++++++++++++ cmd/pairing.go | 18 +-- cmd/root.go | 7 + cmd/traces_cmd.go | 122 ++++++++++++++++++ cmd/traces_cmd_test.go | 95 ++++++++++++++ cmd/traces_output.go | 177 ++++++++++++++++++++++++++ cmd/traces_paths.go | 169 ++++++++++++++++++++++++ cmd/traces_run.go | 93 ++++++++++++++ cmd/traces_run_test.go | 160 +++++++++++++++++++++++ docs/10-tracing-observability.md | 32 +++++ docs/18-http-api.md | 20 +++ docs/project-changelog.md | 18 +++ 15 files changed, 1139 insertions(+), 45 deletions(-) create mode 100644 cmd/gateway_http_client_resolution.go create mode 100644 cmd/gateway_http_client_test.go create mode 100644 cmd/traces_cmd.go create mode 100644 cmd/traces_cmd_test.go create mode 100644 cmd/traces_output.go create mode 100644 cmd/traces_paths.go create mode 100644 cmd/traces_run.go create mode 100644 cmd/traces_run_test.go diff --git a/README.md b/README.md index b1bc9ce6..42c9ec9d 100644 --- a/README.md +++ b/README.md @@ -184,6 +184,16 @@ make logs # Tail logs (goclaw service) make reset # Wipe volumes and rebuild from scratch ``` +**Operator CLI:** + +The main `goclaw` binary can also inspect local or remote gateways: + +```bash +goclaw traces list --status error +goclaw traces get -o json +goclaw --server https://goclaw.example.com --token "$GOCLAW_GATEWAY_TOKEN" traces follow --session +``` + **Optional services** — enable with `WITH_*` flags: | Flag | Service | What it does | diff --git a/cmd/gateway_http_client.go b/cmd/gateway_http_client.go index b1788a67..ec6fe4d5 100644 --- a/cmd/gateway_http_client.go +++ b/cmd/gateway_http_client.go @@ -8,8 +8,6 @@ import ( "net/http" "os" "time" - - "github.com/nextlevelbuilder/goclaw/internal/config" ) // gatewayHTTPError represents a structured error from the gateway HTTP API. @@ -27,35 +25,7 @@ var httpClient = &http.Client{Timeout: 10 * time.Second} // healthClient has a shorter timeout for quick health checks. var healthClient = &http.Client{Timeout: 3 * time.Second} -// resolveGatewayBaseURL reads host/port from config and returns http://host:port. -func resolveGatewayBaseURL() string { - cfg, err := config.Load(resolveConfigPath()) - if err != nil { - return "http://127.0.0.1:18790" - } - host := cfg.Gateway.Host - if host == "" || host == "0.0.0.0" { - host = "127.0.0.1" - } - port := cfg.Gateway.Port - if port == 0 { - port = 18790 - } - return fmt.Sprintf("http://%s:%d", host, port) -} - -// resolveGatewayToken returns the gateway auth token. -// Priority: GOCLAW_GATEWAY_TOKEN env → config file token. -func resolveGatewayToken() string { - if t := os.Getenv("GOCLAW_GATEWAY_TOKEN"); t != "" { - return t - } - cfg, _ := config.Load(resolveConfigPath()) - if cfg != nil { - return cfg.Gateway.Token - } - return "" -} +const gatewayHTTPResponseLimit = 1 << 20 // gatewayHTTPDo sends an HTTP request to the gateway with auth and returns the parsed JSON response. func gatewayHTTPDo(method, path string, body any) (map[string]any, error) { @@ -103,6 +73,10 @@ func gatewayHTTPDelete(path string) error { // gatewayHTTPDoRaw executes an HTTP request and returns the raw response bytes. // Shared by both map-based and typed response functions. func gatewayHTTPDoRaw(method, path string, body any) ([]byte, int, error) { + return gatewayHTTPDoRawWithLimit(method, path, body, gatewayHTTPResponseLimit) +} + +func gatewayHTTPDoRawWithLimit(method, path string, body any, limit int64) ([]byte, int, error) { base := resolveGatewayBaseURL() var bodyReader io.Reader @@ -130,7 +104,13 @@ func gatewayHTTPDoRaw(method, path string, body any) ([]byte, int, error) { } defer resp.Body.Close() - raw, _ := io.ReadAll(io.LimitReader(resp.Body, 1<<20)) + raw, err := io.ReadAll(io.LimitReader(resp.Body, limit+1)) + if err != nil { + return nil, resp.StatusCode, fmt.Errorf("read gateway response: %w", err) + } + if int64(len(raw)) > limit { + return nil, resp.StatusCode, fmt.Errorf("gateway response exceeds %d bytes", limit) + } return raw, resp.StatusCode, nil } diff --git a/cmd/gateway_http_client_resolution.go b/cmd/gateway_http_client_resolution.go new file mode 100644 index 00000000..37204dbf --- /dev/null +++ b/cmd/gateway_http_client_resolution.go @@ -0,0 +1,88 @@ +package cmd + +import ( + "fmt" + "net/url" + "os" + "strings" + + "github.com/nextlevelbuilder/goclaw/internal/config" +) + +// resolveGatewayBaseURL reads host/port from config and returns http://host:port. +func resolveGatewayBaseURL() string { + if base := firstNonEmpty(gatewayServerOverride, os.Getenv("GOCLAW_SERVER"), os.Getenv("GOCLAW_GATEWAY_URL")); base != "" { + return normalizeGatewayBaseURL(base) + } + + cfg, err := config.Load(resolveConfigPath()) + if err != nil { + return "http://127.0.0.1:18790" + } + host := cfg.Gateway.Host + if host == "" || host == "0.0.0.0" { + host = "127.0.0.1" + } + port := cfg.Gateway.Port + if port == 0 { + port = 18790 + } + return fmt.Sprintf("http://%s:%d", host, port) +} + +// resolveGatewayToken returns the gateway auth token. +// Priority: --token flag -> GOCLAW_GATEWAY_TOKEN env -> config file token. +func resolveGatewayToken() string { + if t := strings.TrimSpace(gatewayTokenOverride); t != "" { + return t + } + if t := os.Getenv("GOCLAW_GATEWAY_TOKEN"); t != "" { + return t + } + cfg, _ := config.Load(resolveConfigPath()) + if cfg != nil { + return cfg.Gateway.Token + } + return "" +} + +func resolveGatewayWebSocketURL() (string, error) { + baseURL := resolveGatewayBaseURL() + parsed, err := url.Parse(baseURL) + if err != nil { + return "", fmt.Errorf("parse gateway URL %q: %w", baseURL, err) + } + switch parsed.Scheme { + case "https": + parsed.Scheme = "wss" + case "http": + parsed.Scheme = "ws" + case "ws", "wss": + default: + return "", fmt.Errorf("unsupported gateway URL scheme %q", parsed.Scheme) + } + parsed.Path = strings.TrimRight(parsed.Path, "/") + "/ws" + parsed.RawQuery = "" + parsed.Fragment = "" + return parsed.String(), nil +} + +func normalizeGatewayBaseURL(raw string) string { + base := strings.TrimSpace(raw) + if base == "" { + return "" + } + if !strings.HasPrefix(base, "http://") && !strings.HasPrefix(base, "https://") { + base = "http://" + base + } + return strings.TrimRight(base, "/") +} + +func firstNonEmpty(values ...string) string { + for _, v := range values { + if trimmed := strings.TrimSpace(v); trimmed != "" { + return trimmed + } + } + return "" +} diff --git a/cmd/gateway_http_client_test.go b/cmd/gateway_http_client_test.go new file mode 100644 index 00000000..af5d25d1 --- /dev/null +++ b/cmd/gateway_http_client_test.go @@ -0,0 +1,131 @@ +package cmd + +import ( + "net/http" + "net/http/httptest" + "testing" +) + +func TestResolveGatewayClientOverrides(t *testing.T) { + t.Setenv("GOCLAW_GATEWAY_TOKEN", "env-token") + t.Setenv("GOCLAW_SERVER", "") + t.Setenv("GOCLAW_GATEWAY_URL", "") + + oldServer := gatewayServerOverride + oldToken := gatewayTokenOverride + oldCfg := cfgFile + t.Cleanup(func() { + gatewayServerOverride = oldServer + gatewayTokenOverride = oldToken + cfgFile = oldCfg + }) + + cfgFile = "/path/that/does/not/exist" + gatewayServerOverride = "https://goclaw.example.com/" + gatewayTokenOverride = "flag-token" + + if got := resolveGatewayBaseURL(); got != "https://goclaw.example.com" { + t.Fatalf("resolveGatewayBaseURL() = %q, want trimmed override", got) + } + if got := resolveGatewayToken(); got != "flag-token" { + t.Fatalf("resolveGatewayToken() = %q, want flag token", got) + } + + gatewayTokenOverride = "" + if got := resolveGatewayToken(); got != "env-token" { + t.Fatalf("resolveGatewayToken() = %q, want env token", got) + } + + gatewayServerOverride = "" + t.Setenv("GOCLAW_SERVER", "remote.example.com:18790/") + if got := resolveGatewayBaseURL(); got != "http://remote.example.com:18790" { + t.Fatalf("resolveGatewayBaseURL() = %q, want normalized env URL", got) + } +} + +func TestGatewayHTTPDoRawUsesServerAndTokenOverride(t *testing.T) { + oldServer := gatewayServerOverride + oldToken := gatewayTokenOverride + oldClient := httpClient + t.Cleanup(func() { + gatewayServerOverride = oldServer + gatewayTokenOverride = oldToken + httpClient = oldClient + }) + + var sawRequest bool + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + sawRequest = true + if r.URL.Path != "/v1/traces" { + t.Errorf("path = %q, want /v1/traces", r.URL.Path) + } + if got := r.Header.Get("Authorization"); got != "Bearer flag-token" { + t.Errorf("Authorization = %q, want Bearer flag-token", got) + } + if got := r.Header.Get("X-GoClaw-User-Id"); got != "system" { + t.Errorf("X-GoClaw-User-Id = %q, want system", got) + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{"ok":true}`)) + })) + defer srv.Close() + + gatewayServerOverride = srv.URL + gatewayTokenOverride = "flag-token" + httpClient = srv.Client() + + raw, status, err := gatewayHTTPDoRaw(http.MethodGet, "/v1/traces", nil) + if err != nil { + t.Fatalf("gatewayHTTPDoRaw: %v", err) + } + if status != http.StatusOK { + t.Fatalf("status = %d, want 200", status) + } + if string(raw) != `{"ok":true}` { + t.Fatalf("raw = %s", raw) + } + if !sawRequest { + t.Fatal("test server did not receive request") + } +} + +func TestResolveGatewayWebSocketURLUsesServerOverride(t *testing.T) { + oldServer := gatewayServerOverride + oldCfg := cfgFile + t.Cleanup(func() { + gatewayServerOverride = oldServer + cfgFile = oldCfg + }) + + cfgFile = "/path/that/does/not/exist" + gatewayServerOverride = "https://goclaw.example.com/base/" + + got, err := resolveGatewayWebSocketURL() + if err != nil { + t.Fatalf("resolveGatewayWebSocketURL: %v", err) + } + if got != "wss://goclaw.example.com/base/ws" { + t.Fatalf("resolveGatewayWebSocketURL() = %q", got) + } +} + +func TestGatewayHTTPDoRawWithLimitRejectsOversizedResponse(t *testing.T) { + oldServer := gatewayServerOverride + oldClient := httpClient + t.Cleanup(func() { + gatewayServerOverride = oldServer + httpClient = oldClient + }) + + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + _, _ = w.Write([]byte("12345")) + })) + defer srv.Close() + + gatewayServerOverride = srv.URL + httpClient = srv.Client() + + if _, _, err := gatewayHTTPDoRawWithLimit(http.MethodGet, "/big", nil, 4); err == nil { + t.Fatal("expected oversized response error") + } +} diff --git a/cmd/pairing.go b/cmd/pairing.go index aed8b9a0..d2d652e4 100644 --- a/cmd/pairing.go +++ b/cmd/pairing.go @@ -3,14 +3,12 @@ package cmd import ( "encoding/json" "fmt" - "net/url" "os" "time" "github.com/gorilla/websocket" "github.com/spf13/cobra" - "github.com/nextlevelbuilder/goclaw/internal/config" "github.com/nextlevelbuilder/goclaw/pkg/protocol" ) @@ -187,26 +185,20 @@ func pairingRevokeCmd() *cobra.Command { // gatewayRPC connects to the running gateway, authenticates, sends an RPC call, and returns the response. func gatewayRPC(method string, params json.RawMessage) (*protocol.ResponseFrame, error) { - cfg, err := config.Load(resolveConfigPath()) + wsURL, err := resolveGatewayWebSocketURL() if err != nil { - return nil, fmt.Errorf("load config: %w", err) + return nil, err } - host := cfg.Gateway.Host - if host == "0.0.0.0" { - host = "127.0.0.1" - } - - u := url.URL{Scheme: "ws", Host: fmt.Sprintf("%s:%d", host, cfg.Gateway.Port), Path: "/ws"} - conn, _, err := websocket.DefaultDialer.Dial(u.String(), nil) + conn, _, err := websocket.DefaultDialer.Dial(wsURL, nil) if err != nil { - return nil, fmt.Errorf("connect to gateway at %s: %w", u.String(), err) + return nil, fmt.Errorf("connect to gateway at %s: %w", wsURL, err) } defer conn.Close() // Step 1: Send connect handshake connectParams, _ := json.Marshal(map[string]any{ - "token": cfg.Gateway.Token, + "token": resolveGatewayToken(), "protocol": protocol.ProtocolVersion, }) connectReq := protocol.RequestFrame{ diff --git a/cmd/root.go b/cmd/root.go index 8f4fd9a4..4e812115 100644 --- a/cmd/root.go +++ b/cmd/root.go @@ -15,6 +15,10 @@ var Version = "dev" var ( cfgFile string verbose bool + + gatewayServerOverride string + gatewayTokenOverride string + gatewayOutputFormat string ) var rootCmd = &cobra.Command{ @@ -29,6 +33,8 @@ var rootCmd = &cobra.Command{ func init() { rootCmd.PersistentFlags().StringVar(&cfgFile, "config", "", "config file (default: config.json or $GOCLAW_CONFIG)") rootCmd.PersistentFlags().BoolVarP(&verbose, "verbose", "v", false, "enable debug logging") + rootCmd.PersistentFlags().StringVar(&gatewayServerOverride, "server", "", "gateway server URL override") + rootCmd.PersistentFlags().StringVar(&gatewayTokenOverride, "token", "", "gateway bearer token override") rootCmd.AddCommand(onboardCmd()) rootCmd.AddCommand(versionCmd()) @@ -42,6 +48,7 @@ func init() { rootCmd.AddCommand(cronCmd()) rootCmd.AddCommand(skillsCmd()) rootCmd.AddCommand(sessionsCmd()) + rootCmd.AddCommand(tracesCmd()) rootCmd.AddCommand(migrateCmd()) rootCmd.AddCommand(upgradeCmd()) rootCmd.AddCommand(backupCmd()) diff --git a/cmd/traces_cmd.go b/cmd/traces_cmd.go new file mode 100644 index 00000000..456c023e --- /dev/null +++ b/cmd/traces_cmd.go @@ -0,0 +1,122 @@ +package cmd + +import "github.com/spf13/cobra" + +const traceExportResponseLimit = 64 << 20 + +func tracesCmd() *cobra.Command { + cmd := &cobra.Command{ + Use: "traces", + Short: "Inspect gateway traces", + } + cmd.PersistentFlags().StringVarP(&gatewayOutputFormat, "output", "o", "table", "output format (table|json)") + cmd.AddCommand(tracesListCmd()) + cmd.AddCommand(tracesGetCmd()) + cmd.AddCommand(tracesExportCmd()) + cmd.AddCommand(tracesFollowCmd()) + cmd.AddCommand(tracesTimelineCmd()) + return cmd +} + +func tracesListCmd() *cobra.Command { + var opts traceListOptions + cmd := &cobra.Command{ + Use: "list", + Short: "List traces", + RunE: func(cmd *cobra.Command, args []string) error { + requireRunningGatewayHTTP() + return runTracesList(opts) + }, + } + addTraceListFlags(cmd, &opts) + return cmd +} + +func tracesGetCmd() *cobra.Command { + return &cobra.Command{ + Use: "get ", + Short: "Get trace details with spans", + Args: cobra.ExactArgs(1), + RunE: func(cmd *cobra.Command, args []string) error { + requireRunningGatewayHTTP() + return runTracesGet(args[0]) + }, + } +} + +func tracesExportCmd() *cobra.Command { + var filePath string + cmd := &cobra.Command{ + Use: "export ", + Short: "Export a gzipped trace tree", + Args: cobra.ExactArgs(1), + RunE: func(cmd *cobra.Command, args []string) error { + requireRunningGatewayHTTP() + return runTracesExport(args[0], filePath) + }, + } + cmd.Flags().StringVar(&filePath, "file", "", "write gzip export to file (use - for stdout)") + return cmd +} + +func tracesFollowCmd() *cobra.Command { + var opts traceFollowOptions + cmd := &cobra.Command{ + Use: "follow", + Short: "Poll trace changes for a session or agent", + RunE: func(cmd *cobra.Command, args []string) error { + requireRunningGatewayHTTP() + return runTracesFollow(opts) + }, + } + cmd.Flags().StringVar(&opts.SessionKey, "session", "", "filter by session key") + cmd.Flags().StringVar(&opts.AgentID, "agent-id", "", "filter by agent UUID") + cmd.Flags().StringVar(&opts.UserID, "user", "", "filter by user ID for admin callers") + cmd.Flags().StringVar(&opts.Status, "status", "", "filter by trace status") + cmd.Flags().StringVar(&opts.Channel, "channel", "", "filter by raw channel") + cmd.Flags().StringVar(&opts.Since, "since", "", "RFC3339 lower bound for changed traces") + cmd.Flags().IntVar(&opts.Limit, "limit", 0, "page size, max 200") + cmd.Flags().BoolVar(&opts.IncludeSpans, "include-spans", false, "include spans grouped by trace ID") + return cmd +} + +func tracesTimelineCmd() *cobra.Command { + var opts traceTimelineOptions + cmd := &cobra.Command{ + Use: "timeline ", + Short: "Show the run timeline linked to a trace", + Args: cobra.ExactArgs(1), + RunE: func(cmd *cobra.Command, args []string) error { + requireRunningGatewayHTTP() + return runTracesTimeline(args[0], opts) + }, + } + cmd.Flags().IntVar(&opts.Limit, "limit", 0, "page size, max 500") + cmd.Flags().IntVar(&opts.Offset, "offset", 0, "pagination offset") + return cmd +} + +func addTraceListFlags(cmd *cobra.Command, opts *traceListOptions) { + cmd.Flags().StringVarP(&opts.Query, "query", "q", "", "search trace text, IDs, labels, and span previews") + cmd.Flags().StringVar(&opts.AgentID, "agent-id", "", "filter by agent UUID") + cmd.Flags().StringVar(&opts.UserID, "user", "", "filter by user ID for admin callers") + cmd.Flags().StringVar(&opts.SessionKey, "session", "", "filter by session key") + cmd.Flags().StringVar(&opts.Status, "status", "", "filter by trace status") + cmd.Flags().StringVar(&opts.Channel, "channel", "", "filter by raw channel") + cmd.Flags().StringVar(&opts.AgentQuery, "agent", "", "search agent display name or key") + cmd.Flags().StringVar(&opts.ChannelQuery, "channel-query", "", "search channel instance labels") + cmd.Flags().StringVar(&opts.ToolName, "tool", "", "search span tool names") + cmd.Flags().StringVar(&opts.From, "from", "", "start time lower bound, RFC3339") + cmd.Flags().StringVar(&opts.To, "to", "", "start time upper bound, RFC3339") + cmd.Flags().StringVar(&opts.Since, "since", "", "alias for --from") + cmd.Flags().StringVar(&opts.Until, "until", "", "alias for --to") + cmd.Flags().StringVar(&opts.HasToolCalls, "has-tool-calls", "", "filter true or false") + cmd.Flags().IntVar(&opts.MinInputTokens, "min-input-tokens", 0, "minimum input tokens") + cmd.Flags().IntVar(&opts.MaxInputTokens, "max-input-tokens", 0, "maximum input tokens") + cmd.Flags().IntVar(&opts.MinOutputTokens, "min-output-tokens", 0, "minimum output tokens") + cmd.Flags().IntVar(&opts.MaxOutputTokens, "max-output-tokens", 0, "maximum output tokens") + cmd.Flags().IntVar(&opts.MinToolCalls, "min-tool-calls", 0, "minimum tool calls") + cmd.Flags().IntVar(&opts.MaxToolCalls, "max-tool-calls", 0, "maximum tool calls") + cmd.Flags().IntVar(&opts.Limit, "limit", 0, "page size, max 200") + cmd.Flags().IntVar(&opts.Offset, "offset", 0, "pagination offset") +} diff --git a/cmd/traces_cmd_test.go b/cmd/traces_cmd_test.go new file mode 100644 index 00000000..9cba5db8 --- /dev/null +++ b/cmd/traces_cmd_test.go @@ -0,0 +1,95 @@ +package cmd + +import ( + "net/url" + "testing" +) + +func TestBuildTraceListPathIncludesFilters(t *testing.T) { + path := buildTraceListPath(traceListOptions{ + Query: "provider fail", + SessionKey: "session A", + Status: "running", + AgentQuery: "coder", + HasToolCalls: "true", + Limit: 25, + Offset: 10, + }) + + u, err := url.Parse(path) + if err != nil { + t.Fatalf("parse path: %v", err) + } + if u.Path != "/v1/traces" { + t.Fatalf("path = %q, want /v1/traces", u.Path) + } + q := u.Query() + assertQueryValue(t, q, "q", "provider fail") + assertQueryValue(t, q, "session_key", "session A") + assertQueryValue(t, q, "status", "running") + assertQueryValue(t, q, "agent", "coder") + assertQueryValue(t, q, "has_tool_calls", "true") + assertQueryValue(t, q, "limit", "25") + assertQueryValue(t, q, "offset", "10") +} + +func TestBuildTraceFollowPathRequiresScope(t *testing.T) { + if _, err := buildTraceFollowPath(traceFollowOptions{}); err == nil { + t.Fatal("expected missing scope error") + } + + path, err := buildTraceFollowPath(traceFollowOptions{ + SessionKey: "session A", + Since: "2026-06-12T01:00:00Z", + Limit: 20, + IncludeSpans: true, + }) + if err != nil { + t.Fatalf("buildTraceFollowPath: %v", err) + } + u, err := url.Parse(path) + if err != nil { + t.Fatalf("parse path: %v", err) + } + if u.Path != "/v1/traces/follow" { + t.Fatalf("path = %q, want /v1/traces/follow", u.Path) + } + q := u.Query() + assertQueryValue(t, q, "session_key", "session A") + assertQueryValue(t, q, "since", "2026-06-12T01:00:00Z") + assertQueryValue(t, q, "limit", "20") + assertQueryValue(t, q, "include_spans", "true") +} + +func TestTraceTimelinePathUsesRunIDFromDetail(t *testing.T) { + runID, err := traceRunIDFromDetail(traceDetailResponse{ + Trace: traceDataForCLI{RunID: "run-123"}, + }) + if err != nil { + t.Fatalf("traceRunIDFromDetail: %v", err) + } + if runID != "run-123" { + t.Fatalf("runID = %q, want run-123", runID) + } + + if _, err := traceRunIDFromDetail(traceDetailResponse{}); err == nil { + t.Fatal("expected missing run_id error") + } +} + +func TestValidateTraceOutputFormatRejectsUnsupportedValues(t *testing.T) { + oldOutput := gatewayOutputFormat + t.Cleanup(func() { gatewayOutputFormat = oldOutput }) + + gatewayOutputFormat = "yaml" + if err := validateTraceOutputFormat(); err == nil { + t.Fatal("expected unsupported output format error") + } +} + +func assertQueryValue(t *testing.T, q url.Values, key, want string) { + t.Helper() + if got := q.Get(key); got != want { + t.Fatalf("query %s = %q, want %q", key, got, want) + } +} diff --git a/cmd/traces_output.go b/cmd/traces_output.go new file mode 100644 index 00000000..b83c9c11 --- /dev/null +++ b/cmd/traces_output.go @@ -0,0 +1,177 @@ +package cmd + +import ( + "bytes" + "compress/gzip" + "encoding/json" + "fmt" + "io" + "os" + "strings" + "text/tabwriter" + "time" + + "github.com/google/uuid" +) + +func outputFormatIsJSON() bool { + return strings.EqualFold(strings.TrimSpace(gatewayOutputFormat), "json") +} + +func validateTraceOutputFormat() error { + switch strings.ToLower(strings.TrimSpace(gatewayOutputFormat)) { + case "", "table", "json": + return nil + default: + return fmt.Errorf("unsupported output format %q; use table or json", gatewayOutputFormat) + } +} + +func printJSON(value any) error { + data, err := json.MarshalIndent(value, "", " ") + if err != nil { + return err + } + fmt.Println(string(data)) + return nil +} + +func printTraceList(resp traceListResponse) error { + if outputFormatIsJSON() { + return printJSON(resp) + } + if len(resp.Traces) == 0 { + fmt.Println("No traces found.") + return nil + } + tw := tabwriter.NewWriter(os.Stdout, 0, 0, 2, ' ', 0) + fmt.Fprintln(tw, "TRACE\tSTATUS\tAGENT\tSESSION\tRUN\tTOKENS\tSTARTED") + for _, tr := range resp.Traces { + fmt.Fprintf(tw, "%s\t%s\t%s\t%s\t%s\t%d/%d\t%s\n", + shortTraceID(tr.ID.String()), + tr.Status, + shortOptionalUUID(tr.AgentID), + truncateStr(tr.SessionKey, 36), + truncateStr(tr.RunID, 24), + tr.TotalInputTokens, + tr.TotalOutputTokens, + formatTraceTime(tr.StartTime), + ) + } + return tw.Flush() +} + +func printTraceDetail(resp traceDetailResponse) error { + if outputFormatIsJSON() { + return printJSON(resp) + } + tr := resp.Trace + fmt.Printf("Trace: %s\n", tr.ID) + fmt.Printf("Status: %s\n", tr.Status) + fmt.Printf("Agent: %s\n", shortOptionalUUID(tr.AgentID)) + fmt.Printf("Session: %s\n", tr.SessionKey) + fmt.Printf("Run: %s\n", tr.RunID) + fmt.Printf("Tokens: %d input / %d output\n", tr.TotalInputTokens, tr.TotalOutputTokens) + if tr.Error != "" { + fmt.Printf("Error: %s\n", tr.Error) + } + if len(resp.Spans) == 0 { + fmt.Println("\nNo spans found.") + return nil + } + fmt.Println() + tw := tabwriter.NewWriter(os.Stdout, 0, 0, 2, ' ', 0) + fmt.Fprintln(tw, "SPAN\tTYPE\tSTATUS\tPROVIDER\tMODEL\tDURATION\tNAME") + for _, sp := range resp.Spans { + fmt.Fprintf(tw, "%s\t%s\t%s\t%s\t%s\t%dms\t%s\n", + shortTraceID(sp.ID.String()), + sp.SpanType, + sp.Status, + sp.Provider, + sp.Model, + sp.DurationMS, + truncateStr(sp.Name, 48), + ) + } + return tw.Flush() +} + +func printTraceFollow(resp traceFollowResponse) error { + if outputFormatIsJSON() { + return printJSON(resp) + } + if resp.NextSince != "" { + fmt.Printf("Next since: %s\n\n", resp.NextSince) + } + return printTraceList(traceListResponse{ + Traces: resp.Traces, + Total: len(resp.Traces), + Limit: resp.Limit, + }) +} + +func printTraceTimeline(resp traceTimelineResponse) error { + if outputFormatIsJSON() { + return printJSON(resp) + } + if len(resp.Items) == 0 { + fmt.Println("No timeline items found.") + return nil + } + tw := tabwriter.NewWriter(os.Stdout, 0, 0, 2, ' ', 0) + fmt.Fprintln(tw, "SEQ\tTYPE\tSTATUS\tTITLE\tPREVIEW\tCREATED") + for _, item := range resp.Items { + fmt.Fprintf(tw, "%d\t%s\t%s\t%s\t%s\t%s\n", + item.Seq, + item.ItemType, + item.Status, + truncateStr(item.Title, 32), + truncateStr(firstNonEmpty(item.Preview, item.Content), 60), + formatTraceTime(item.CreatedAt), + ) + } + return tw.Flush() +} + +func printGzipJSON(raw []byte) error { + gr, err := gzip.NewReader(bytes.NewReader(raw)) + if err != nil { + return err + } + defer gr.Close() + data, err := io.ReadAll(gr) + if err != nil { + return err + } + var value any + if err := json.Unmarshal(data, &value); err != nil { + return err + } + return printJSON(value) +} + +func defaultTraceExportPath(traceID string) string { + short := shortTraceID(traceID) + return fmt.Sprintf("trace-%s-%s.json.gz", short, time.Now().UTC().Format("20060102")) +} + +func shortTraceID(id string) string { + if len(id) <= 8 { + return id + } + return id[:8] +} + +func shortOptionalUUID(id *uuid.UUID) string { + if id == nil { + return "-" + } + return shortTraceID(id.String()) +} + +func formatTraceTime(t time.Time) string { + if t.IsZero() { + return "-" + } + return t.Local().Format(time.DateTime) +} diff --git a/cmd/traces_paths.go b/cmd/traces_paths.go new file mode 100644 index 00000000..edc7765a --- /dev/null +++ b/cmd/traces_paths.go @@ -0,0 +1,169 @@ +package cmd + +import ( + "fmt" + "net/url" + "strconv" + "strings" + + "github.com/nextlevelbuilder/goclaw/internal/store" +) + +type traceDataForCLI = store.TraceData +type spanDataForCLI = store.SpanData +type timelineItemForCLI = store.RunTimelineItem + +type traceListOptions struct { + Query string + AgentID string + UserID string + SessionKey string + Status string + Channel string + AgentQuery string + ChannelQuery string + ToolName string + From string + To string + Since string + Until string + HasToolCalls string + MinInputTokens int + MaxInputTokens int + MinOutputTokens int + MaxOutputTokens int + MinToolCalls int + MaxToolCalls int + Limit int + Offset int +} + +type traceFollowOptions struct { + AgentID string + SessionKey string + Status string + Channel string + UserID string + Since string + Limit int + IncludeSpans bool +} + +type traceTimelineOptions struct { + Limit int + Offset int +} + +type traceListResponse struct { + Traces []traceDataForCLI `json:"traces"` + Total int `json:"total"` + Limit int `json:"limit"` + Offset int `json:"offset"` +} + +type traceDetailResponse struct { + Trace traceDataForCLI `json:"trace"` + Spans []spanDataForCLI `json:"spans"` +} + +type traceFollowResponse struct { + Traces []traceDataForCLI `json:"traces"` + SpansByTraceID map[string][]spanDataForCLI `json:"spans_by_trace_id"` + ServerTime string `json:"server_time"` + NextSince string `json:"next_since"` + Limit int `json:"limit"` +} + +type traceTimelineResponse struct { + RunID string `json:"run_id"` + SessionKey string `json:"session_key"` + Items []timelineItemForCLI `json:"items"` + Limit int `json:"limit"` + Offset int `json:"offset"` +} + +func buildTraceListPath(opts traceListOptions) string { + values := url.Values{} + addQuery(values, "q", opts.Query) + addQuery(values, "agent_id", opts.AgentID) + addQuery(values, "user_id", opts.UserID) + addQuery(values, "session_key", opts.SessionKey) + addQuery(values, "status", opts.Status) + addQuery(values, "channel", opts.Channel) + addQuery(values, "agent", opts.AgentQuery) + addQuery(values, "channel_query", opts.ChannelQuery) + addQuery(values, "tool_name", opts.ToolName) + addQuery(values, "from", firstNonEmpty(opts.From, opts.Since)) + addQuery(values, "to", firstNonEmpty(opts.To, opts.Until)) + addQuery(values, "has_tool_calls", opts.HasToolCalls) + addIntQuery(values, "min_input_tokens", opts.MinInputTokens) + addIntQuery(values, "max_input_tokens", opts.MaxInputTokens) + addIntQuery(values, "min_output_tokens", opts.MinOutputTokens) + addIntQuery(values, "max_output_tokens", opts.MaxOutputTokens) + addIntQuery(values, "min_tool_calls", opts.MinToolCalls) + addIntQuery(values, "max_tool_calls", opts.MaxToolCalls) + addIntQuery(values, "limit", opts.Limit) + addIntQuery(values, "offset", opts.Offset) + return pathWithQuery("/v1/traces", values) +} + +func buildTraceFollowPath(opts traceFollowOptions) (string, error) { + if strings.TrimSpace(opts.SessionKey) == "" && strings.TrimSpace(opts.AgentID) == "" { + return "", fmt.Errorf("traces follow requires --session or --agent-id") + } + values := url.Values{} + addQuery(values, "agent_id", opts.AgentID) + addQuery(values, "session_key", opts.SessionKey) + addQuery(values, "status", opts.Status) + addQuery(values, "channel", opts.Channel) + addQuery(values, "user_id", opts.UserID) + addQuery(values, "since", opts.Since) + addIntQuery(values, "limit", opts.Limit) + if opts.IncludeSpans { + values.Set("include_spans", "true") + } + return pathWithQuery("/v1/traces/follow", values), nil +} + +func buildTraceTimelinePath(runID, sessionKey string, opts traceTimelineOptions) string { + values := url.Values{} + addQuery(values, "session_key", sessionKey) + addIntQuery(values, "limit", opts.Limit) + addIntQuery(values, "offset", opts.Offset) + return pathWithQuery("/v1/runs/"+url.PathEscape(runID)+"/timeline", values) +} + +func traceDetailPath(traceID string) string { + return "/v1/traces/" + url.PathEscape(traceID) +} + +func traceExportPath(traceID string) string { + return "/v1/traces/" + url.PathEscape(traceID) + "/export" +} + +func traceRunIDFromDetail(detail traceDetailResponse) (string, error) { + runID := strings.TrimSpace(detail.Trace.RunID) + if runID == "" { + return "", fmt.Errorf("trace has no run_id; timeline is unavailable") + } + return runID, nil +} + +func addQuery(values url.Values, key, value string) { + if trimmed := strings.TrimSpace(value); trimmed != "" { + values.Set(key, trimmed) + } +} + +func addIntQuery(values url.Values, key string, value int) { + if value > 0 { + values.Set(key, strconv.Itoa(value)) + } +} + +func pathWithQuery(path string, values url.Values) string { + if len(values) == 0 { + return path + } + return path + "?" + values.Encode() +} diff --git a/cmd/traces_run.go b/cmd/traces_run.go new file mode 100644 index 00000000..ed93ce85 --- /dev/null +++ b/cmd/traces_run.go @@ -0,0 +1,93 @@ +package cmd + +import ( + "fmt" + "net/http" + "os" +) + +func runTracesList(opts traceListOptions) error { + if err := validateTraceOutputFormat(); err != nil { + return err + } + resp, err := gatewayHTTPGetTyped[traceListResponse](buildTraceListPath(opts)) + if err != nil { + return err + } + return printTraceList(resp) +} + +func runTracesGet(traceID string) error { + if err := validateTraceOutputFormat(); err != nil { + return err + } + resp, err := gatewayHTTPGetTyped[traceDetailResponse](traceDetailPath(traceID)) + if err != nil { + return err + } + return printTraceDetail(resp) +} + +func runTracesExport(traceID, filePath string) error { + if err := validateTraceOutputFormat(); err != nil { + return err + } + raw, status, err := gatewayHTTPDoRawWithLimit(http.MethodGet, traceExportPath(traceID), nil, traceExportResponseLimit) + if err != nil { + return err + } + if status >= 400 { + return parseHTTPError(raw, status) + } + if outputFormatIsJSON() { + return printGzipJSON(raw) + } + if filePath == "-" { + _, err = os.Stdout.Write(raw) + return err + } + if filePath == "" { + filePath = defaultTraceExportPath(traceID) + } + if err := os.WriteFile(filePath, raw, 0o600); err != nil { + return err + } + fmt.Printf("Exported trace to %s\n", filePath) + return nil +} + +func runTracesFollow(opts traceFollowOptions) error { + if err := validateTraceOutputFormat(); err != nil { + return err + } + path, err := buildTraceFollowPath(opts) + if err != nil { + return err + } + resp, err := gatewayHTTPGetTyped[traceFollowResponse](path) + if err != nil { + return err + } + return printTraceFollow(resp) +} + +func runTracesTimeline(traceID string, opts traceTimelineOptions) error { + if err := validateTraceOutputFormat(); err != nil { + return err + } + detail, err := gatewayHTTPGetTyped[traceDetailResponse](traceDetailPath(traceID)) + if err != nil { + return err + } + runID, err := traceRunIDFromDetail(detail) + if err != nil { + return err + } + resp, err := gatewayHTTPGetTyped[traceTimelineResponse]( + buildTraceTimelinePath(runID, detail.Trace.SessionKey, opts), + ) + if err != nil { + return err + } + return printTraceTimeline(resp) +} diff --git a/cmd/traces_run_test.go b/cmd/traces_run_test.go new file mode 100644 index 00000000..50b430ba --- /dev/null +++ b/cmd/traces_run_test.go @@ -0,0 +1,160 @@ +package cmd + +import ( + "bytes" + "compress/gzip" + "io" + "net/http" + "net/http/httptest" + "os" + "path/filepath" + "strings" + "testing" +) + +func TestRunTracesGetJSONOutput(t *testing.T) { + traceID := "11111111-1111-1111-1111-111111111111" + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/v1/traces/"+traceID { + t.Fatalf("path = %q", r.URL.Path) + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{ + "trace": { + "id": "` + traceID + `", + "run_id": "run-123", + "start_time": "2026-06-12T01:00:00Z", + "created_at": "2026-06-12T01:00:00Z", + "status": "completed", + "total_input_tokens": 10, + "total_output_tokens": 20 + }, + "spans": [] + }`)) + })) + defer srv.Close() + withTraceTestGateway(t, srv) + gatewayOutputFormat = "json" + + out, err := captureStdout(t, func() error { + return runTracesGet(traceID) + }) + if err != nil { + t.Fatalf("runTracesGet: %v", err) + } + if !strings.Contains(out, `"run_id": "run-123"`) { + t.Fatalf("output missing run_id: %s", out) + } +} + +func TestRunTracesExportWritesFileAndPrintsJSON(t *testing.T) { + var gzipPayload bytes.Buffer + gz := gzip.NewWriter(&gzipPayload) + _, _ = gz.Write([]byte(`{"trace":{"id":"trace-1"},"spans":[],"sub_traces":[]}`)) + if err := gz.Close(); err != nil { + t.Fatalf("close gzip: %v", err) + } + + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/v1/traces/trace-1/export" { + t.Fatalf("path = %q", r.URL.Path) + } + w.Header().Set("Content-Type", "application/gzip") + _, _ = w.Write(gzipPayload.Bytes()) + })) + defer srv.Close() + withTraceTestGateway(t, srv) + + gatewayOutputFormat = "table" + outFile := filepath.Join(t.TempDir(), "trace.json.gz") + if _, err := captureStdout(t, func() error { + return runTracesExport("trace-1", outFile) + }); err != nil { + t.Fatalf("runTracesExport file: %v", err) + } + written, err := os.ReadFile(outFile) + if err != nil { + t.Fatalf("read export: %v", err) + } + if !bytes.Equal(written, gzipPayload.Bytes()) { + t.Fatal("written gzip payload mismatch") + } + + gatewayOutputFormat = "json" + out, err := captureStdout(t, func() error { + return runTracesExport("trace-1", "") + }) + if err != nil { + t.Fatalf("runTracesExport json: %v", err) + } + if !strings.Contains(out, `"trace"`) || !strings.Contains(out, `"sub_traces"`) { + t.Fatalf("json output missing trace tree fields: %s", out) + } +} + +func TestRunTracesFollowJSONOutput(t *testing.T) { + srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/v1/traces/follow" { + t.Fatalf("path = %q", r.URL.Path) + } + if got := r.URL.Query().Get("session_key"); got != "session-1" { + t.Fatalf("session_key = %q", got) + } + w.Header().Set("Content-Type", "application/json") + _, _ = w.Write([]byte(`{ + "traces": [], + "spans_by_trace_id": {}, + "server_time": "2026-06-12T01:00:00Z", + "next_since": "2026-06-12T01:00:00Z", + "limit": 50 + }`)) + })) + defer srv.Close() + withTraceTestGateway(t, srv) + gatewayOutputFormat = "json" + + out, err := captureStdout(t, func() error { + return runTracesFollow(traceFollowOptions{SessionKey: "session-1"}) + }) + if err != nil { + t.Fatalf("runTracesFollow: %v", err) + } + if !strings.Contains(out, `"next_since": "2026-06-12T01:00:00Z"`) { + t.Fatalf("output missing next_since: %s", out) + } +} + +func withTraceTestGateway(t *testing.T, srv *httptest.Server) { + t.Helper() + oldServer := gatewayServerOverride + oldToken := gatewayTokenOverride + oldOutput := gatewayOutputFormat + oldClient := httpClient + t.Cleanup(func() { + gatewayServerOverride = oldServer + gatewayTokenOverride = oldToken + gatewayOutputFormat = oldOutput + httpClient = oldClient + }) + gatewayServerOverride = srv.URL + gatewayTokenOverride = "test-token" + httpClient = srv.Client() +} + +func captureStdout(t *testing.T, fn func() error) (string, error) { + t.Helper() + oldStdout := os.Stdout + r, w, err := os.Pipe() + if err != nil { + t.Fatalf("pipe: %v", err) + } + os.Stdout = w + runErr := fn() + _ = w.Close() + os.Stdout = oldStdout + out, readErr := io.ReadAll(r) + if readErr != nil { + t.Fatalf("read stdout: %v", readErr) + } + return string(out), runErr +} diff --git a/docs/10-tracing-observability.md b/docs/10-tracing-observability.md index f69f998f..6ed99ed6 100644 --- a/docs/10-tracing-observability.md +++ b/docs/10-tracing-observability.md @@ -187,7 +187,10 @@ worker.Stop() | Method | Path | Description | |--------|------|-------------| | GET | `/v1/traces` | List traces with pagination and filters | +| GET | `/v1/traces/follow` | Poll trace changes for one session or agent | | GET | `/v1/traces/{id}` | Get trace details with all spans | +| GET | `/v1/traces/{id}/export` | Export a gzipped trace tree with spans and sub-traces | +| GET | `/v1/runs/{runID}/timeline` | Get persisted run archive timeline items | ### Query Filters @@ -210,6 +213,35 @@ worker.Stop() | `limit` | int | Page size (default 50) | | `offset` | int | Pagination offset | +### Operator CLI + +The main `goclaw` binary can also act as a thin operator client for trace +inspection: + +```bash +goclaw traces list --status error --limit 20 +goclaw traces get -o json +goclaw traces export --file trace.json.gz +goclaw traces follow --session --since 2026-06-12T01:00:00Z +goclaw traces timeline +``` + +By default, commands use the same local gateway config and +`GOCLAW_GATEWAY_TOKEN` behavior as existing admin commands. For remote +operations, use explicit overrides: + +```bash +goclaw --server https://goclaw.example.com --token "$GOCLAW_GATEWAY_TOKEN" traces get -o json +``` + +`--server` also applies to existing WebSocket/RPC-backed admin commands such as +`sessions`, `cron`, and `pairing`. The URL can also come from `GOCLAW_SERVER` +or `GOCLAW_GATEWAY_URL`; `--server` wins when both are set. + +The standalone `nextlevelbuilder/goclaw-cli` can remain a compatibility tool, +but first-party trace operator workflows are now available from the main +server/runtime binary. + --- ## 8. Delegation History diff --git a/docs/18-http-api.md b/docs/18-http-api.md index f8c24449..745b22d9 100644 --- a/docs/18-http-api.md +++ b/docs/18-http-api.md @@ -1393,6 +1393,26 @@ Follow response: } ``` +Main binary operator commands wrap the same endpoints: + +```bash +goclaw traces list --query "provider fail" --status error +goclaw traces get -o json +goclaw traces export --file trace.json.gz +goclaw traces follow --session +goclaw traces timeline +``` + +Remote gateways use the shared client overrides: + +```bash +goclaw --server https://goclaw.example.com --token "$GOCLAW_GATEWAY_TOKEN" traces list -o json +``` + +The same `--server` / `--token` resolver is shared with WebSocket/RPC-backed +admin commands. `GOCLAW_SERVER` or `GOCLAW_GATEWAY_URL` can provide the base URL +when the flag is omitted. + ### Run Timeline `GET /v1/runs/{runID}/timeline` returns display-safe archive entries for one diff --git a/docs/project-changelog.md b/docs/project-changelog.md index 3b820c20..28a69dc4 100644 --- a/docs/project-changelog.md +++ b/docs/project-changelog.md @@ -6,6 +6,24 @@ Significant changes, features, and fixes in reverse chronological order. ## 2026-06-12 +### Operator trace CLI (issue #158) + +**Changes** + +- Added first-class `goclaw traces` operator commands to the main binary: + `list`, `get`, `export`, `follow`, and `timeline`. +- Added remote client overrides with `--server` and `--token`; trace commands + support trace-scoped output selection via `--output` / `-o`. +- Shared the URL/token resolver with WebSocket/RPC-backed admin commands so + `sessions`, `cron`, and `pairing` can use the same remote gateway overrides. +- Kept trace commands as thin wrappers over existing HTTP endpoints so the + server/runtime and operator CLI share one binary without new API contracts. + +**Tests** + +- Added command/client regressions for gateway URL and token overrides, trace + query serialization, follow scope validation, and timeline run ID handling. + ### Mid-flight request preservation (issue #137) **Fixes**