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
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -277,6 +277,7 @@ pb promql active-queries
# Datasets
pb dataset list
pb dataset info <dataset>
pb dataset schema <dataset>
pb dataset add <dataset>
pb dataset remove <dataset>

Expand Down Expand Up @@ -353,6 +354,7 @@ Commands that support `-o json` return structured output:
pb status -o json
pb profile list -o json
pb dataset list -o json
pb dataset schema <dataset> -o json
pb sql list -o json
pb sql run "SELECT count(*) FROM backend" --from=1h --output json
pb promql run "up" --dataset otel_metrics --instant --output json
Expand Down
1 change: 1 addition & 0 deletions cmd/agent.go
Original file line number Diff line number Diff line change
Expand Up @@ -107,6 +107,7 @@ func newAgentManifest() agentManifest {
{Command: "pb profile list -o json", Description: "List locally configured profiles without credentials", Scope: "local", Mutates: false, RequiresProfile: false, Constraints: []string{"Requires the local pb config file"}},
{Command: "pb dataset list -o json", Description: "List datasets", Scope: "server", Mutates: false, RequiresProfile: true},
{Command: "pb dataset info <dataset> -o json", Description: "Read dataset statistics and configuration", Scope: "server", Mutates: false, RequiresProfile: true},
{Command: "pb dataset schema <dataset> -o json", Description: "Read dataset field names and types", Scope: "server", Mutates: false, RequiresProfile: true},
{Command: "pb user list -o json", Description: "List users and their roles", Scope: "server", Mutates: false, RequiresProfile: true},
{Command: "pb role list -o json", Description: "List roles and privileges", Scope: "server", Mutates: false, RequiresProfile: true},
{Command: "pb sql run \"<SELECT query>\" --from <time> --to <time> -o json", Description: "Run a read-only SQL query", Scope: "server", Mutates: false, RequiresProfile: true, Constraints: []string{"Use SELECT-only SQL", "Never use --save-as because it creates a saved query"}},
Expand Down
2 changes: 1 addition & 1 deletion cmd/agent_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,7 @@ func TestAgentManifestContainsOnlyReadOnlyCommands(t *testing.T) {
}
seen[item.Command] = struct{}{}
}
for _, command := range []string{"pb -o json", "pb help <command...> -o json", "pb dataset list -o json", "pb sql list -o json", "pb promql active-queries -o json"} {
for _, command := range []string{"pb -o json", "pb help <command...> -o json", "pb dataset list -o json", "pb dataset schema <dataset> -o json", "pb sql list -o json", "pb promql active-queries -o json"} {
if _, exists := seen[command]; !exists {
t.Fatalf("read-only command missing from agent catalog: %s", command)
}
Expand Down
97 changes: 97 additions & 0 deletions cmd/dataset_schema.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,97 @@
// Copyright (c) 2024 Parseable, Inc
//
// This program is free software: you can redistribute it and/or modify
// it under the terms of the GNU Affero General Public License as published by
// the Free Software Foundation, either version 3 of the License, or
// (at your option) any later version.
//
// This program is distributed in the hope that it will be useful
// but WITHOUT ANY WARRANTY; without even the implied warranty of
// MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the
// GNU Affero General Public License for more details.
//
// You should have received a copy of the GNU Affero General Public License
// along with this program. If not, see <http://www.gnu.org/licenses/>.

package cmd

import (
"bytes"
"encoding/json"
"fmt"
"io"
"net/http"

internalHTTP "github.com/parseablehq/pb/pkg/http"
"github.com/spf13/cobra"
)

// SchemaDatasetCmd prints the fields and types in a dataset's current schema.
var SchemaDatasetCmd = &cobra.Command{
Use: "schema <dataset>",
Short: "Show dataset field names and types",
Example: " pb dataset schema backend_logs\n pb dataset schema backend_logs -o json",
Args: cobra.ExactArgs(1),
RunE: func(cmd *cobra.Command, args []string) error {
format, err := commandOutputFormat(cmd)
if err != nil {
return err
}

client := internalHTTP.DefaultClient(&DefaultProfile)
req, err := client.NewRequest(http.MethodGet, "logstream/"+args[0]+"/schema", nil)
if err != nil {
return err
}
resp, err := client.Client.Do(req)
if err != nil {
return err
}
defer resp.Body.Close()

body, err := io.ReadAll(resp.Body)
if err != nil {
return err
}
if resp.StatusCode != http.StatusOK {
return responseStatusError("get dataset schema", resp.StatusCode, resp.Status, body)
}

var schema struct {
Fields []struct {
Name string `json:"name"`
DataType json.RawMessage `json:"data_type"`
} `json:"fields"`
}
if err := json.Unmarshal(body, &schema); err != nil {
return newInvalidResponseError(fmt.Errorf("invalid dataset schema: %w", err))
}
if schema.Fields == nil {
return newInvalidResponseError(fmt.Errorf("dataset schema has no fields array"))
}

if format == outputJSON {
return writeRawJSON(cmd.OutOrStdout(), body)
}

out := cmd.OutOrStdout()
for _, field := range schema.Fields {
var typeName string
if err := json.Unmarshal(field.DataType, &typeName); err != nil {
var compact bytes.Buffer
if err := json.Compact(&compact, field.DataType); err != nil {
return newInvalidResponseError(fmt.Errorf("invalid type for field %q: %w", field.Name, err))
}
typeName = compact.String()
}
if _, err := fmt.Fprintln(out, (&DatasetListItem{Name: field.Name, Type: typeName}).Render()); err != nil {
return err
}
}
return nil
},
}

func init() {
SchemaDatasetCmd.Flags().StringP("output", "o", "text", "Output format (text|json)")
}
75 changes: 75 additions & 0 deletions cmd/dataset_schema_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,75 @@
package cmd

import (
"bytes"
"encoding/json"
"net/http"
"testing"
)

const testDatasetSchema = `{"fields":[{"name":"service_name","data_type":"Utf8","nullable":true},{"name":"timestamp","data_type":{"Timestamp":["Nanosecond",null]}}],"metadata":{"source":"test"}}`

func TestDatasetSchemaOutput(t *testing.T) {
useCommandTransport(t, func(req *http.Request) (*http.Response, error) {
if req.Method != http.MethodGet || req.URL.Path != "/api/v1/logstream/events/schema" {
t.Fatalf("unexpected request: %s %s", req.Method, req.URL.Path)
}
return commandHTTPResponse(http.StatusOK, "200 OK", testDatasetSchema), nil
})

for _, test := range []struct {
format string
want string
}{
{format: "text", want: (&DatasetListItem{Name: "service_name", Type: "Utf8"}).Render() + "\n" +
(&DatasetListItem{Name: "timestamp", Type: `{"Timestamp":["Nanosecond",null]}`}).Render() + "\n"},
{format: "json"},
} {
t.Run(test.format, func(t *testing.T) {
if err := SchemaDatasetCmd.Flags().Set("output", test.format); err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = SchemaDatasetCmd.Flags().Set("output", "text") })
var output bytes.Buffer
SchemaDatasetCmd.SetOut(&output)
t.Cleanup(func() { SchemaDatasetCmd.SetOut(nil) })

if err := SchemaDatasetCmd.RunE(SchemaDatasetCmd, []string{"events"}); err != nil {
t.Fatal(err)
}
if test.format == "text" {
if got := output.String(); got != test.want {
t.Fatalf("unexpected schema text: %q", got)
}
return
}
var got, want any
if err := json.Unmarshal(output.Bytes(), &got); err != nil {
t.Fatalf("invalid JSON output: %v", err)
}
if err := json.Unmarshal([]byte(testDatasetSchema), &want); err != nil {
t.Fatal(err)
}
gotJSON, _ := json.Marshal(got)
wantJSON, _ := json.Marshal(want)
if !bytes.Equal(gotJSON, wantJSON) {
t.Fatalf("schema JSON changed: got %s, want %s", gotJSON, wantJSON)
}
})
}
}

func TestDatasetSchemaNotFound(t *testing.T) {
useCommandTransport(t, func(*http.Request) (*http.Response, error) {
return commandHTTPResponse(http.StatusNotFound, "404 Not Found", "missing dataset"), nil
})
var output bytes.Buffer
SchemaDatasetCmd.SetOut(&output)
t.Cleanup(func() { SchemaDatasetCmd.SetOut(nil) })
if err := SchemaDatasetCmd.RunE(SchemaDatasetCmd, []string{"missing"}); err == nil || errorDetails(err).Code != ErrorNotFound {
t.Fatalf("expected not-found error, got %v", err)
}
if output.Len() != 0 {
t.Fatalf("unexpected output on error: %q", output.String())
}
}
1 change: 1 addition & 0 deletions main.go
Original file line number Diff line number Diff line change
Expand Up @@ -245,6 +245,7 @@ func main() {
dataset.AddCommand(pb.RemoveDatasetCmd)
dataset.AddCommand(pb.ListDatasetCmd)
dataset.AddCommand(pb.StatDatasetCmd)
dataset.AddCommand(pb.SchemaDatasetCmd)

sql.AddCommand(pb.QueryCmd)
sql.AddCommand(pb.SaveSQLCmd)
Expand Down
Loading