diff --git a/.gitignore b/.gitignore index 484111e7c..8cc080522 100644 --- a/.gitignore +++ b/.gitignore @@ -23,7 +23,7 @@ dist/ .DS_Store demo/sling_commands_demo.workflow *.screenstudio -./sling +/sling core/dbio/filesys/test/dataset1M.csv core/dbio/filesys/test/dataset100k.csv tests/suite/ diff --git a/.goreleaser.mac.yaml b/.goreleaser.mac.yaml index 6c8bba09e..25fed998f 100644 --- a/.goreleaser.mac.yaml +++ b/.goreleaser.mac.yaml @@ -54,4 +54,11 @@ brews: branch: main homepage: https://slingdata.io/ - description: "Data Integration made simple, from the command line. Extract and load data from popular data sources to destinations with high performance and ease." \ No newline at end of file + description: "Data Integration made simple, from the command line. Extract and load data from popular data sources to destinations with high performance and ease." + service: | + run [opt_bin/"sling", "serve", "workbench", "--no-browser"] + keep_alive true + working_dir Dir.home + log_path var/"log/sling-workbench.log" + error_log_path var/"log/sling-workbench.log" + environment_variables PATH: std_service_path_env \ No newline at end of file diff --git a/core/dbio/api/spec.go b/core/dbio/api/spec.go index f9c5a9efb..ec6ba7e09 100644 --- a/core/dbio/api/spec.go +++ b/core/dbio/api/spec.go @@ -1435,6 +1435,7 @@ type SingleRequest struct { id string `yaml:"-" json:"-"` timestamp int64 `yaml:"-" json:"-"` durationMs int64 `yaml:"-" json:"-"` + index int `yaml:"-" json:"-"` // 1-based request number of the endpoint iter *Iteration `yaml:"-" json:"-"` // the iteration that the req belongs to state StateMap `yaml:"-" json:"-"` // copy of iteration state for request (prevents mutation) endpoint *Endpoint `yaml:"-" json:"-"` @@ -1461,6 +1462,7 @@ func NewSingleRequest(iter *Iteration) *SingleRequest { return &SingleRequest{ id: id, timestamp: time.Now().UnixMilli(), + index: iter.endpoint.totalReqs, endpoint: iter.endpoint, iter: iter, state: state, @@ -1492,38 +1494,99 @@ func (lrs *SingleRequest) Map() map[string]any { return vars } -// SpecEventChn, when non-nil, receives structured spec test events as JSON-safe maps. -// The LSP layer creates this channel before a test run and drains it in a goroutine. -var SpecEventChn chan map[string]any +// Spec event types. The strings are the contract for spec inspectors +// (LSP, MCP, workbench). +const ( + SpecEventTypeEndpointStart = "endpoint-start" + SpecEventTypeRequestComplete = "request-complete" + SpecEventTypeRecords = "records" + SpecEventTypeEndpointDone = "endpoint-done" + SpecEventTypeError = "error" +) -// FireSpecEvent sends an event to SpecEventChn if it is non-nil. -func FireSpecEvent(event map[string]any) { - if SpecEventChn != nil { - SpecEventChn <- event +// SpecEvent is one structured event of a spec test run. It carries the +// request/response, the iteration state before and after the request, and the +// records the request pulled. +type SpecEvent struct { + Type string `json:"type"` + Endpoint string `json:"endpoint,omitempty"` + RequestIndex int `json:"request_index,omitempty"` + Request map[string]any `json:"request,omitempty"` + Response map[string]any `json:"response,omitempty"` + StateBefore map[string]any `json:"state_before,omitempty"` + StateAfter map[string]any `json:"state_after,omitempty"` + Records []any `json:"records,omitempty"` + Error string `json:"error,omitempty"` + DurationMs int64 `json:"duration_ms,omitempty"` + + // legacy fields, read by released spec inspectors (VS Code extension) + ReqID string `json:"req_id,omitempty"` + Timestamp int64 `json:"timestamp,omitempty"` + IterID string `json:"iter_id,omitempty"` + IterSequence int `json:"iter_sequence,omitempty"` + SizeBytes int `json:"size_bytes,omitempty"` // request-complete: response body size + RecordCount int `json:"record_count,omitempty"` // endpoint-done: records pulled +} + +// specEventCtxKey carries a spec test's event handler through the request +// context, so its events go to its own consumer instead of a package global. +type specEventCtxKey struct{} + +// WithSpecEventHandler returns a context carrying fn as the spec event +// handler. The API client picks it up from the request context, so a test run +// is isolated and its cancel stops it. +func WithSpecEventHandler(ctx context.Context, fn func(SpecEvent)) context.Context { + if fn == nil { + return ctx + } + return context.WithValue(ctx, specEventCtxKey{}, fn) +} + +// fireSpecEvent sends one event to the handler carried by ctx, if any. +func fireSpecEvent(ctx context.Context, event SpecEvent) { + if ctx == nil { + return + } + if fn, ok := ctx.Value(specEventCtxKey{}).(func(SpecEvent)); ok && fn != nil { + fn(event) } } -// ToSpecEvent builds a JSON-safe map with all request/response details -// for the spec inspector. Includes unexported fields (id, timestamp, -// endpoint name, iteration id) that don't normally marshal. -func (req *SingleRequest) ToSpecEvent() map[string]any { - event := g.M( - "type", "request-complete", - "req_id", req.id, - "timestamp", req.timestamp, - "endpoint", req.endpoint.Name, - "duration_ms", req.durationMs, - ) +// recordsToAny adapts records for SpecEvent.Records ([]any). +func recordsToAny(records []map[string]any) []any { + if len(records) == 0 { + return nil + } + out := make([]any, len(records)) + for i, rec := range records { + out[i] = rec + } + return out +} + +// ToSpecEvent builds the request-complete event for the spec inspector. +// It includes the request index, the iteration state before/after the request +// and its duration, which don't normally marshal. +func (req *SingleRequest) ToSpecEvent() SpecEvent { + event := SpecEvent{ + Type: SpecEventTypeRequestComplete, + ReqID: req.id, + Timestamp: req.timestamp, + Endpoint: req.endpoint.Name, + RequestIndex: req.index, + DurationMs: req.durationMs, + StateBefore: maps.Clone(req.state), + } // iteration context if req.iter != nil { - event["iter_id"] = req.iter.id - event["iter_sequence"] = req.iter.sequence + event.IterID = req.iter.id + event.IterSequence = req.iter.sequence } // request state if req.Request != nil { - event["request"] = g.M( + event.Request = g.M( "method", req.Request.Method, "url", req.Request.URL, "headers", req.Request.Headers, @@ -1534,16 +1597,27 @@ func (req *SingleRequest) ToSpecEvent() map[string]any { // response state if req.Response != nil { - event["size_bytes"] = len(req.Response.Text) - event["response"] = g.M( + event.SizeBytes = len(req.Response.Text) + event.Response = g.M( "status", req.Response.Status, "headers", req.Response.Headers, "body", req.Response.Text, + "size_bytes", len(req.Response.Text), "record_count", len(req.Response.Records), "records", req.Response.Records, ) } + // state after the request: processors may have moved the iteration state. + // Lock iter.context: other goroutines write iter.state concurrently. + if req.iter != nil { + stateAfter := StateMap{} + req.iter.context.Lock() + maps.Copy(stateAfter, req.iter.state) + req.iter.context.Unlock() + event.StateAfter = stateAfter + } + return event } diff --git a/core/dbio/connection/connection_discover.go b/core/dbio/connection/connection_discover.go index 58dfc1c82..d73ce4ef0 100644 --- a/core/dbio/connection/connection_discover.go +++ b/core/dbio/connection/connection_discover.go @@ -18,9 +18,79 @@ import ( "github.com/spf13/cast" ) +// TestOptions configures a connection test. +type TestOptions struct { + Endpoints []string // endpoint names to test; empty means all + Limit int // records per request (default 10) + MaxRequests int // requests per endpoint (default 2) + Context map[string]any // spec test context (store, range, mode) + SpecFile string // overlay this spec file on the connection's spec + Trace bool // trace-level logging + OnEvent func(api.SpecEvent) +} + +// testOptionsFromEnv builds TestOptions from the SLING_TEST_* environment +// variables, the way `sling conns test` has always configured a test. +func testOptionsFromEnv() TestOptions { + opts := TestOptions{} + + if val := os.Getenv("SLING_TEST_ENDPOINTS"); val != "" { + opts.Endpoints = strings.Split(val, ",") + } + + limit := cast.ToInt(g.Getenv("SLING_TEST_ENDPOINT_LIMIT", "10")) + if limit > 1000 { + limit = 1000 // let's set the max limit to 1000 for testing + } + if g.Getenv("SLING_TEST_ENDPOINT_LIMIT") == "" { + g.Debug(env.MagentaString(g.F("testing endpoints with a record limit: %d. Set env var SLING_TEST_ENDPOINT_LIMIT to modify.", limit))) + } + opts.Limit = limit + + maxRequests := cast.ToInt(g.Getenv("SLING_TEST_ENDPOINT_MAX_REQUESTS", "2")) + if maxRequests == 0 { + maxRequests = 3 + } + if g.Getenv("SLING_TEST_ENDPOINT_MAX_REQUESTS") == "" { + g.Debug(env.MagentaString(g.F("testing endpoints with a max requests: %d. Set env var SLING_TEST_ENDPOINT_MAX_REQUESTS to modify.", maxRequests))) + } + opts.MaxRequests = maxRequests + + if val := g.Getenv("SLING_TEST_ENDPOINT_CONTEXT"); val != "" { + contextMap := g.M() + if err := g.Unmarshal(val, &contextMap); err != nil { + g.Warn("could not set context for spec testing: %s", err.Error()) + } + opts.Context = contextMap + } + + return opts +} + +// Test keeps its signature: it builds the options from the SLING_TEST_* env +// vars and calls TestWithOptions, so the CLI does not change. func (c *Connection) Test() (ok bool, err error) { + return c.TestWithOptions(context.Background(), testOptionsFromEnv()) +} + +// TestWithOptions tests a connection. For API connections it tests the +// requested (or all) endpoints, emitting a spec event per step to opts.OnEvent +// and stopping promptly when ctx is cancelled. +func (c *Connection) TestWithOptions(ctx context.Context, opts TestOptions) (ok bool, err error) { os.Setenv("SLING_TEST_MODE", "true") + if ctx == nil { + ctx = context.Background() + } + if opts.OnEvent != nil { + ctx = api.WithSpecEventHandler(ctx, opts.OnEvent) + } + if opts.Trace { + level := g.GetLogLevel() + g.SetLogLevel(g.TraceLevel) + defer g.SetLogLevel(level) + } + switch { case c.Type.IsDb(): dbConn, err := c.AsDatabase(AsConnOptions{UseCache: c.GetType() == dbio.TypeDbDuckDb}) @@ -37,9 +107,9 @@ func (c *Connection) Test() (ok bool, err error) { return ok, g.Error(err, "could not initiate %s", c.Name) } - ctx, cancel := context.WithTimeout(context.Background(), 25*time.Second) + fileCtx, cancel := context.WithTimeout(ctx, 25*time.Second) defer cancel() - err = fileClient.Init(ctx) + err = fileClient.Init(fileCtx) if err != nil { return ok, g.Error(err, "could not connect to %s", c.Name) } @@ -56,7 +126,7 @@ func (c *Connection) Test() (ok bool, err error) { g.Debug(g.Marshal(nodes.Paths())) } case c.Type.IsAPI(): - apiClient, err := c.AsAPI(AsConnOptions{UseCache: false}) + apiClient, err := c.AsAPIContext(ctx, AsConnOptions{UseCache: false}) if err != nil { return ok, g.Error(err, "could not initiate %s", c.Name) } @@ -68,40 +138,46 @@ func (c *Connection) Test() (ok bool, err error) { return ok, g.Error(err, "could not authenticate to %s", c.Name) } - var testEndpoints, testedEndpoints []string - if val := os.Getenv("SLING_TEST_ENDPOINTS"); val != "" { - testEndpoints = strings.Split(os.Getenv("SLING_TEST_ENDPOINTS"), ",") - } + testEndpoints := opts.Endpoints endpoints, err := apiClient.ListEndpoints() if err != nil { return ok, g.Error(err, "could not list endpoints") } - limit := cast.ToInt(g.Getenv("SLING_TEST_ENDPOINT_LIMIT", "10")) + limit := opts.Limit + if limit <= 0 { + limit = 10 + } if limit > 1000 { limit = 1000 // let's set the max limit to 1000 for testing } - if g.Getenv("SLING_TEST_ENDPOINT_LIMIT") == "" { - g.Debug(env.MagentaString(g.F("testing endpoints with a record limit: %d. Set env var SLING_TEST_ENDPOINT_LIMIT to modify.", limit))) - } - maxRequests := cast.ToInt(g.Getenv("SLING_TEST_ENDPOINT_MAX_REQUESTS", "2")) - if maxRequests == 0 { - maxRequests = 3 + maxRequests := opts.MaxRequests + if maxRequests <= 0 { + maxRequests = 2 } - apiClient.Context.Map.Set("max_requests", maxRequests) - if g.Getenv("SLING_TEST_ENDPOINT_MAX_REQUESTS") == "" { - g.Debug(env.MagentaString(g.F("testing endpoints with a max requests: %d. Set env var SLING_TEST_ENDPOINT_MAX_REQUESTS to modify.", maxRequests))) - } // obtain the best endpoint for testing one (for connectivity/authentication) + // (legacy env path, kept for callers that use Test) if cast.ToBool(g.Getenv("SLING_TEST_SINGLE_ENDPOINT")) { testEndpoints = []string{apiClient.GetTestEndpoint()} } + emit := func(event api.SpecEvent) { + if opts.OnEvent != nil { + opts.OnEvent(event) + } + } + + var testedEndpoints []string for _, endpoint := range endpoints { + // a cancelled test stops before the next endpoint + if err := ctx.Err(); err != nil { + return false, err + } + // check for match to test (if provided) allowTest := len(testEndpoints) == 0 for _, testEndpoint := range testEndpoints { @@ -115,22 +191,16 @@ func (c *Connection) Test() (ok bool, err error) { println() g.Info("testing endpoint: %#v", endpoint.Name) - api.FireSpecEvent(g.M("type", "endpoint-start", "endpoint", endpoint.Name)) + emit(api.SpecEvent{Type: api.SpecEventTypeEndpointStart, Endpoint: endpoint.Name}) testedEndpoints = append(testedEndpoints, endpoint.Name) // set limits for testing options := api.APIStreamConfig{Flatten: 1, Limit: limit} // set context if provided - contextPayload := cast.ToString(g.Getenv("SLING_TEST_ENDPOINT_CONTEXT")) - if contextPayload != "" { - contextMap := g.M() - if err := g.Unmarshal(contextPayload, &contextMap); err != nil { - g.Warn("could not set context for spec testing: %s", err.Error()) - } - + if len(opts.Context) > 0 { // set store - if store, ok := contextMap["store"]; ok && store != "" { + if store, ok := opts.Context["store"]; ok && store != "" { storeMap, err := g.UnmarshalMap(cast.ToString(store)) if err != nil { g.Warn("could not unmarshal context store: %s", err.Error()) @@ -139,22 +209,29 @@ func (c *Connection) Test() (ok bool, err error) { } // set range & mode - options.Range = cast.ToString(contextMap["range"]) - options.Mode = cast.ToString(contextMap["mode"]) + options.Range = cast.ToString(opts.Context["range"]) + options.Mode = cast.ToString(opts.Context["mode"]) } df, err := apiClient.ReadDataflow(endpoint.Name, options) if err != nil { + if ctxErr := ctx.Err(); ctxErr != nil { + return false, ctxErr + } + emit(api.SpecEvent{Type: api.SpecEventTypeError, Endpoint: endpoint.Name, Error: err.Error()}) return ok, g.Error(err, "error testing endpoint: %s", endpoint.Name) } data, err := df.Collect() if err != nil { + if ctxErr := ctx.Err(); ctxErr != nil { + return false, ctxErr + } + emit(api.SpecEvent{Type: api.SpecEventTypeError, Endpoint: endpoint.Name, Error: err.Error()}) return ok, g.Error(err, "could collect data from endpoint: %s", endpoint.Name) } g.Debug(" got %d records from endpoint: %s", len(data.Rows), endpoint.Name) - api.FireSpecEvent(g.M("type", "endpoint-done", "endpoint", endpoint.Name, "record_count", len(data.Rows))) records := data.Records(false) if len(records) > 0 { @@ -162,6 +239,12 @@ func (c *Connection) Test() (ok bool, err error) { g.Debug(" columns = %s", g.Marshal(lo.Keys(record))) } + emit(api.SpecEvent{ + Type: api.SpecEventTypeEndpointDone, + Endpoint: endpoint.Name, + RecordCount: len(data.Rows), + }) + } for _, testEndpoint := range testEndpoints { @@ -176,6 +259,11 @@ func (c *Connection) Test() (ok bool, err error) { } } + // a cancelled test reports cancellation, not success + if err := ctx.Err(); err != nil { + return false, err + } + } return true, nil diff --git a/core/dbio/connection/connection_local.go b/core/dbio/connection/connection_local.go index 979f03e6a..bdf693e0e 100644 --- a/core/dbio/connection/connection_local.go +++ b/core/dbio/connection/connection_local.go @@ -1,6 +1,7 @@ package connection import ( + "context" "encoding/base64" "encoding/json" "github.com/flarco/g" @@ -15,6 +16,7 @@ import ( "github.com/spf13/cast" "gopkg.in/yaml.v2" "os" + "path/filepath" "sort" "strings" "time" @@ -64,16 +66,71 @@ func (ce ConnEntries) Discover(name string, opt *DiscoverOptions) (nodes filesys return } +// Test keeps its signature: it builds the options from the SLING_TEST_* env +// vars and calls TestWithOptions, so the CLI does not change. func (ce ConnEntries) Test(name string) (ok bool, err error) { + return ce.TestWithOptions(context.Background(), name, testOptionsFromEnv()) +} + +// TestWithOptions tests the named connection with opts. When opts.SpecFile is +// set, it overlays that spec file on the connection's spec for this test only. +func (ce ConnEntries) TestWithOptions(ctx context.Context, name string, opts TestOptions) (ok bool, err error) { + if opts.SpecFile != "" { + entries, err := ce.withSpecFile(name, opts.SpecFile) + if err != nil { + return false, err + } + ce = entries + } + conn := ce.Get(name) if conn.Name == "" { return ok, g.Error("Invalid Connection name: %s. Make sure it is created. See https://docs.slingdata.io/sling-cli/environment", name) } defer conn.Connection.Close() - ok, err = conn.Connection.Test() + ok, err = conn.Connection.TestWithOptions(ctx, opts) return } +// withSpecFile returns a copy of entries where the named connection's spec +// points at specFile (a relative path resolves against the working directory). +// The original entries stay untouched. +func (ce ConnEntries) withSpecFile(name, specFile string) (ConnEntries, error) { + absSpec := specFile + if !filepath.IsAbs(absSpec) { + wd, err := os.Getwd() + if err != nil { + return nil, g.Error(err, "could not resolve spec file path: %s", specFile) + } + absSpec = filepath.Join(wd, absSpec) + } + if _, err := os.Stat(absSpec); err != nil { + return nil, g.Error(err, "spec file not found: %s", absSpec) + } + + out := make(ConnEntries, len(ce)) + copy(out, ce) + for i := range out { + if !strings.EqualFold(out[i].Name, name) { + continue + } + + data := make(map[string]any, len(out[i].Connection.Data)+1) + for k, v := range out[i].Connection.Data { + data[k] = v + } + data["spec"] = "file://" + absSpec + + conn, err := NewConnection(out[i].Connection.Name, out[i].Connection.Type, data) + if err != nil { + return nil, g.Error(err, "could not overlay spec file on connection %s", name) + } + out[i].Connection = conn + return out, nil + } + return nil, g.Error("Invalid Connection name: %s. Make sure it is created.", name) +} + var ( localConns ConnEntries localConnsTs time.Time diff --git a/core/dbio/connection/connection_test.go b/core/dbio/connection/connection_test.go index a7a0037b0..6b0b38424 100644 --- a/core/dbio/connection/connection_test.go +++ b/core/dbio/connection/connection_test.go @@ -1,12 +1,22 @@ package connection import ( + "context" + "errors" + "fmt" + "net/http" + "net/http/httptest" + "os" + "path/filepath" "strings" + "sync" "testing" + "time" "github.com/flarco/g" "github.com/microsoft/go-mssqldb/msdsn" "github.com/slingdata-io/sling-cli/core/dbio" + "github.com/slingdata-io/sling-cli/core/dbio/api" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" ) @@ -556,3 +566,223 @@ func TestDynamoDBConnectionURL(t *testing.T) { require.NoError(t, err) assert.Equal(t, "ap-southeast-1", c.DataS(true)["aws_region"]) } + +const testOptionsSpecYAML = ` +name: "Test Options API" +defaults: + state: + page: 1 +endpoints: + items: + request: + url: "%s/items" + parameters: + page: "{state.page}" + response: + records: + jmespath: "data[]" + pagination: + stop_condition: "false" + next_state: + page: "state.page + 1" +` + +func testOptionsEntries(t *testing.T, serverURL string) ConnEntries { + t.Helper() + + specPath := filepath.Join(t.TempDir(), "spec.yaml") + require.NoError(t, os.WriteFile(specPath, []byte(fmt.Sprintf(testOptionsSpecYAML, serverURL)), 0644)) + + conn, err := NewConnection("TEST_OPTIONS_API", dbio.TypeApi, map[string]any{ + "type": "api", + "spec": "file://" + specPath, + }) + require.NoError(t, err) + + return ConnEntries{{Name: "TEST_OPTIONS_API", Connection: conn}} +} + +func TestConnEntriesTestWithOptions(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + fmt.Fprint(w, `{"data":[{"id":1},{"id":2}]}`) + })) + defer server.Close() + + var events []api.SpecEvent + levelBefore := g.GetLogLevel() + ok, err := testOptionsEntries(t, server.URL).TestWithOptions(context.Background(), "TEST_OPTIONS_API", TestOptions{ + Endpoints: []string{"items"}, + Limit: 10, + MaxRequests: 2, + Trace: true, + OnEvent: func(e api.SpecEvent) { events = append(events, e) }, + }) + require.NoError(t, err) + require.True(t, ok) + assert.Equal(t, levelBefore, g.GetLogLevel(), "trace log level not restored") + + types := []string{} + for _, e := range events { + types = append(types, e.Type) + } + assert.Equal(t, []string{ + api.SpecEventTypeEndpointStart, + api.SpecEventTypeRequestComplete, api.SpecEventTypeRecords, + api.SpecEventTypeRequestComplete, api.SpecEventTypeRecords, + api.SpecEventTypeEndpointDone, + }, types, "unexpected event sequence: %#v", types) + + completes, records := 0, 0 + for _, e := range events { + switch e.Type { + case api.SpecEventTypeRequestComplete: + completes++ + assert.Equal(t, "items", e.Endpoint) + assert.Equal(t, completes, e.RequestIndex) + assert.NotNil(t, e.Request) + assert.NotNil(t, e.Response) + assert.NotNil(t, e.StateBefore, "state_before missing") + assert.NotNil(t, e.StateAfter, "state_after missing") + assert.EqualValues(t, 2, e.Response["record_count"]) + assert.EqualValues(t, 200, e.Response["status"]) + // legacy wire fields + assert.NotEmpty(t, e.ReqID) + assert.NotZero(t, e.Timestamp) + assert.NotEmpty(t, e.IterID) + assert.Equal(t, len(`{"data":[{"id":1},{"id":2}]}`), e.SizeBytes) + case api.SpecEventTypeEndpointDone: + assert.Equal(t, 4, e.RecordCount) + case api.SpecEventTypeRecords: + records++ + assert.Len(t, e.Records, 2) + } + } + assert.Equal(t, 2, completes, "max_requests not honored") + assert.Equal(t, 2, records, "one records event per request expected") +} + +func TestConnEntriesTestErrorEvent(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {})) + deadURL := server.URL + server.Close() // no listener: requests fail to connect + + var events []api.SpecEvent + ok, err := testOptionsEntries(t, deadURL).TestWithOptions(context.Background(), "TEST_OPTIONS_API", TestOptions{ + Endpoints: []string{"items"}, + OnEvent: func(e api.SpecEvent) { events = append(events, e) }, + }) + require.Error(t, err) + assert.False(t, ok) + require.NotEmpty(t, events) + last := events[len(events)-1] + assert.Equal(t, api.SpecEventTypeError, last.Type, "events: %#v", events) + assert.Equal(t, "items", last.Endpoint) + assert.NotEmpty(t, last.Error) +} + +func TestConnEntriesTestFromEnv(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + fmt.Fprint(w, `{"data":[{"id":1}]}`) + })) + defer server.Close() + + t.Setenv("SLING_TEST_ENDPOINTS", "items") + t.Setenv("SLING_TEST_ENDPOINT_LIMIT", "5") + t.Setenv("SLING_TEST_ENDPOINT_MAX_REQUESTS", "1") + + ok, err := testOptionsEntries(t, server.URL).Test("TEST_OPTIONS_API") + require.NoError(t, err) + assert.True(t, ok) +} + +func TestConnEntriesTestSpecFileOverlay(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + fmt.Fprint(w, `{"data":[{"id":1}]}`) + })) + defer server.Close() + + // the connection's own spec has no "items" endpoint + connSpec := filepath.Join(t.TempDir(), "conn.yaml") + require.NoError(t, os.WriteFile(connSpec, []byte(fmt.Sprintf(` +name: "Conn Spec" +endpoints: + others: + request: + url: "%s/others" + response: + records: + jmespath: "data[]" +`, server.URL)), 0644)) + + conn, err := NewConnection("TEST_OPTIONS_API", dbio.TypeApi, map[string]any{ + "type": "api", + "spec": "file://" + connSpec, + }) + require.NoError(t, err) + entries := ConnEntries{{Name: "TEST_OPTIONS_API", Connection: conn}} + + // the draft spec replaces the connection's spec for this test only + draftSpec := filepath.Join(t.TempDir(), "draft.yaml") + require.NoError(t, os.WriteFile(draftSpec, []byte(fmt.Sprintf(testOptionsSpecYAML, server.URL)), 0644)) + + var events []api.SpecEvent + ok, err := entries.TestWithOptions(context.Background(), "TEST_OPTIONS_API", TestOptions{ + SpecFile: draftSpec, + Endpoints: []string{"items"}, + OnEvent: func(e api.SpecEvent) { events = append(events, e) }, + }) + require.NoError(t, err) + assert.True(t, ok) + require.NotEmpty(t, events) + assert.Equal(t, "items", events[0].Endpoint) + + // the connection entry itself is untouched + assert.Equal(t, "file://"+connSpec, entries.Get("TEST_OPTIONS_API").Connection.Data["spec"]) + + // a missing spec file is reported + _, err = entries.TestWithOptions(context.Background(), "TEST_OPTIONS_API", TestOptions{ + SpecFile: filepath.Join(t.TempDir(), "nope.yaml"), + }) + require.Error(t, err) + assert.Contains(t, err.Error(), "spec file not found") +} + +func TestConnEntriesTestCancel(t *testing.T) { + blocked := make(chan struct{}) + var once sync.Once + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + once.Do(func() { close(blocked) }) + <-r.Context().Done() // returns when the client aborts the request + })) + defer server.Close() + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + errCh := make(chan error, 1) + go func() { + _, err := testOptionsEntries(t, server.URL).TestWithOptions(ctx, "TEST_OPTIONS_API", TestOptions{ + Endpoints: []string{"items"}, + MaxRequests: 2, + }) + errCh <- err + }() + + select { + case <-blocked: + case <-time.After(10 * time.Second): + t.Fatal("request never reached the server") + } + cancel() + + select { + case err := <-errCh: + require.Error(t, err) + assert.True(t, errors.Is(err, context.Canceled), "expected context.Canceled, got %v", err) + case <-time.After(10 * time.Second): + t.Fatal("test did not stop on cancel") + } +} diff --git a/core/env/envfile.go b/core/env/envfile.go index f967ec516..ca474e17e 100644 --- a/core/env/envfile.go +++ b/core/env/envfile.go @@ -28,12 +28,36 @@ type EnvFile struct { Connections map[string]map[string]any `json:"connections,omitempty" yaml:"connections,omitempty"` Env map[string]any `json:"env,omitempty" yaml:"env,omitempty"` Variables map[string]any `json:"variables,omitempty" yaml:"variables,omitempty"` // legacy + Workbench *WorkbenchConfig `json:"workbench,omitempty" yaml:"workbench,omitempty"` Path string `json:"-" yaml:"-"` TopComment string `json:"-" yaml:"-"` Body string `json:"-" yaml:"-"` } +// WorkbenchConfig is the `workbench:` block of env.yaml: the settings of +// `sling serve workbench`. A command-line flag wins over the file. +type WorkbenchConfig struct { + // Host is the listen address. The default is 127.0.0.1. + Host string `json:"host,omitempty" yaml:"host,omitempty"` + // Port is the listen port. The default is 7879; 0 picks a free port. + Port int `json:"port,omitempty" yaml:"port,omitempty"` + // Token is required when Host is not loopback. + Token string `json:"token,omitempty" yaml:"token,omitempty"` + // ProjectsRoot limits the Open-folder dialog to one folder tree. + ProjectsRoot string `json:"projects_root,omitempty" yaml:"projects_root,omitempty"` + // Shell allows interactive shell terminals. The default is true. + Shell *bool `json:"shell,omitempty" yaml:"shell,omitempty"` + // WorkerIdle stops a project worker that has no sessions and no running + // work after this duration, for example "15m". The default is 15m. + WorkerIdle string `json:"worker_idle,omitempty" yaml:"worker_idle,omitempty"` + // PathExtra is prepended to PATH for shells, runs and agent CLIs, so a + // launchd or systemd service finds the tools the user's shell does. + PathExtra string `json:"path_extra,omitempty" yaml:"path_extra,omitempty"` + // Env holds extra environment variables for workers. + Env map[string]any `json:"env,omitempty" yaml:"env,omitempty"` +} + func (ef *EnvFile) WriteEnvFile() (err error) { output, err := ef.marshalEnvFileBytes() if err != nil { @@ -123,7 +147,7 @@ func (ef *EnvFile) structToRootNode(original *yaml.Node) (*yaml.Node, error) { } managed := map[string]struct{}{ - "connections": {}, "variables": {}, "env": {}, + "connections": {}, "variables": {}, "env": {}, "workbench": {}, } if original != nil && len(original.Content) > 0 && original.Content[0].Kind == yaml.MappingNode { newMap := doc.Content[0] diff --git a/go.mod b/go.mod index 0c94b3549..f51b5cf4f 100644 --- a/go.mod +++ b/go.mod @@ -99,6 +99,7 @@ require ( github.com/tidwall/jsonc v0.3.3 github.com/tidwall/sjson v1.2.5 github.com/timeplus-io/proton-go-driver/v2 v2.0.19 + github.com/trebi-ai/agent-wire v0.2.1 github.com/trinodb/trino-go-client v0.328.0 github.com/twpayne/go-geom v1.6.1 github.com/valentin-kaiser/go-dbase v1.14.4 @@ -194,6 +195,7 @@ require ( github.com/containerd/errdefs v1.0.0 // indirect github.com/containerd/errdefs/pkg v0.3.0 // indirect github.com/coreos/go-oidc/v3 v3.21.0 // indirect + github.com/creack/pty v1.1.24 // indirect github.com/creasty/defaults v1.8.0 // indirect github.com/cyphar/filepath-securejoin v0.2.4 // indirect github.com/danieljoos/wincred v1.2.2 // indirect @@ -212,7 +214,7 @@ require ( github.com/exasol/error-reporting-go v0.2.0 // indirect github.com/felixge/httpsnoop v1.1.0 // indirect github.com/francoispqt/gojay v1.2.13 // indirect - github.com/fsnotify/fsnotify v1.9.0 // indirect + github.com/fsnotify/fsnotify v1.10.1 // indirect github.com/gabriel-vasile/mimetype v1.4.7 // indirect github.com/ganigeorgiev/fexpr v0.4.1 // indirect github.com/go-faster/city v1.0.1 // indirect @@ -232,6 +234,7 @@ require ( github.com/goccy/go-json v0.10.6 // indirect github.com/goccy/go-yaml v1.17.1 // indirect github.com/godbus/dbus v0.0.0-20190726142602-4481cbc300e2 // indirect + github.com/gofrs/flock v0.13.0 // indirect github.com/golang-jwt/jwt/v4 v4.5.2 // indirect github.com/golang-jwt/jwt/v5 v5.3.0 // indirect github.com/golang-sql/civil v0.0.0-20220223132316-b832511892a9 // indirect diff --git a/go.sum b/go.sum index e28595d30..8f1f33181 100644 --- a/go.sum +++ b/go.sum @@ -195,8 +195,6 @@ github.com/aws/aws-sdk-go-v2/service/internal/checksum v1.9.15/go.mod h1:e3IzZvQ github.com/aws/aws-sdk-go-v2/service/internal/endpoint-discovery v1.11.14 h1:3exo28cClRTVnxdj/LULxkESZSSv74RUIjZ7tfHXfWQ= github.com/aws/aws-sdk-go-v2/service/internal/endpoint-discovery v1.11.14/go.mod h1:yLon9pByjyB6JZq5IAmwnjE3ObIhD0QibfRWH7tUhLU= github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.9.27/go.mod h1:EOwBD4J4S5qYszS5/3DpkejfuK+Z5/1uzICfPaZLtqw= -github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.23 h1:pbrxO/kuIwgEsOPLkaHu0O+m4fNgLU8B3vxQ+72jTPw= -github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.23/go.mod h1:/CMNUqoj46HpS3MNRDEDIwcgEnrtZlKRaHNaHxIFpNA= github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.31 h1:w2SIhW92DZPFrSL4ksVCr8IYff5OZwIcxg8+95tzvAI= github.com/aws/aws-sdk-go-v2/service/internal/presigned-url v1.13.31/go.mod h1:wAhpCQbkov+IcvjozJbd2xRCoZybUEHNkcFunssNACg= github.com/aws/aws-sdk-go-v2/service/internal/s3shared v1.19.23 h1:03xatSQO4+AM1lTAbnRg5OK528EUg744nW7F73U8DKw= @@ -411,8 +409,8 @@ github.com/fsnotify/fsevents v0.2.0 h1:BRlvlqjvNTfogHfeBOFvSC9N0Ddy+wzQCQukyoD7o github.com/fsnotify/fsevents v0.2.0/go.mod h1:B3eEk39i4hz8y1zaWS/wPrAP4O6wkIl7HQwKBr1qH/w= github.com/fsnotify/fsnotify v1.4.7/go.mod h1:jwhsz4b93w/PPRr/qN1Yymfu8t87LnFCMoQvtojpjFo= github.com/fsnotify/fsnotify v1.5.4/go.mod h1:OVB6XrOHzAwXMpEM7uPOzcehqUV2UqJxmVXmkdnm1bU= -github.com/fsnotify/fsnotify v1.9.0 h1:2Ml+OJNzbYCTzsxtv8vKSFD9PbJjmhYF14k/jKC7S9k= -github.com/fsnotify/fsnotify v1.9.0/go.mod h1:8jBTzvmWwFyi3Pb8djgCCO5IBqzKJ/Jwo8TRcHyHii0= +github.com/fsnotify/fsnotify v1.10.1 h1:b0/UzAf9yR5rhf3RPm9gf3ehBPpf0oZKIjtpKrx59Ho= +github.com/fsnotify/fsnotify v1.10.1/go.mod h1:TLheqan6HD6GBK6PrDWyDPBaEV8LspOxvPSjC+bVfgo= github.com/fvbommel/sortorder v1.1.0 h1:fUmoe+HLsBTctBDoaBwpQo5N+nrCp8g/BjKb/6ZQmYw= github.com/fvbommel/sortorder v1.1.0/go.mod h1:uk88iVf1ovNn1iLfgUVU2F9o5eO30ui720w+kxuqRs0= github.com/gabriel-vasile/mimetype v1.4.7 h1:SKFKl7kD0RiPdbht0s7hFtjl489WcQ1VyPW8ZzUMYCA= @@ -1107,6 +1105,10 @@ github.com/tonistiigi/units v0.0.0-20180711220420-6950e57a87ea h1:SXhTLE6pb6eld/ github.com/tonistiigi/units v0.0.0-20180711220420-6950e57a87ea/go.mod h1:WPnis/6cRcDZSUvVmezrxJPkiO87ThFYsoUiMwWNDJk= github.com/tonistiigi/vt100 v0.0.0-20240514184818-90bafcd6abab h1:H6aJ0yKQ0gF49Qb2z5hI1UHxSQt4JMyxebFR15KnApw= github.com/tonistiigi/vt100 v0.0.0-20240514184818-90bafcd6abab/go.mod h1:ulncasL3N9uLrVann0m+CDlJKWsIAP34MPcOJF6VRvc= +github.com/trebi-ai/agent-wire v0.1.0 h1:iEVwL8Ns8eAk0+pum9JR3fBt7jEBNiLDXNSnkASRwIg= +github.com/trebi-ai/agent-wire v0.1.0/go.mod h1:K0WMWYNBepCzwbhnVZf6p5PNpkmnH0q4t1mM3FRcYJc= +github.com/trebi-ai/agent-wire v0.2.1 h1:qcQ+atPGAltkRg7qJQayEurPkmq0QNxmgWYtBorUCs4= +github.com/trebi-ai/agent-wire v0.2.1/go.mod h1:K0WMWYNBepCzwbhnVZf6p5PNpkmnH0q4t1mM3FRcYJc= github.com/trinodb/trino-go-client v0.328.0 h1:X6hrGGysA3nvyVcz8kJbBS98srLNTNsnNYwRkMC1atA= github.com/trinodb/trino-go-client v0.328.0/go.mod h1:e/nck9W6hy+9bbyZEpXKFlNsufn3lQGpUgDL1d5f1FI= github.com/twmb/avro v1.7.2 h1:cmrEBRSbELRqsg/dRkQvVWuOaR2EfGifHIt/2iJ9lfI= diff --git a/tests/suite.cli.yaml b/tests/suite.cli.yaml index 8b2b3741c..aa8672de8 100644 --- a/tests/suite.cli.yaml +++ b/tests/suite.cli.yaml @@ -3093,6 +3093,41 @@ - 'PRODUCT DESCRIPTION' - execution succeeded +# The workbench server (plan 6.17, 3.15): the command's help, and a smoke test +# that starts it on a free port, asks /healthz and stops it. The temp home keeps +# the run out of the real ~/.sling. +- id: 620 + name: 'Serve workbench help' + run: sling serve workbench --help + output_contains: + - 'Run the Sling workbench' + - '--token' + +- id: 621 + name: 'Serve workbench starts, answers healthz and stops' + run: | + HOME_DIR=$(mktemp -d) + LOG=$(mktemp) + SLING_HOME_DIR="$HOME_DIR" sling serve workbench --no-browser --port 0 > "$LOG" 2>&1 & + PID=$! + PORT="" + for _ in $(seq 1 60); do + PORT=$(sed -n 's|.*workbench listening on http://127.0.0.1:\([0-9]*\).*|\1|p' "$LOG" | head -1) + [ -n "$PORT" ] && break + sleep 0.5 + done + if [ -z "$PORT" ]; then + cat "$LOG" + kill "$PID" 2>/dev/null || true + exit 1 + fi + BODY=$(curl -fsS "http://127.0.0.1:$PORT/healthz") + kill "$PID" 2>/dev/null || true + wait "$PID" 2>/dev/null || true + echo "healthz=$BODY" + output_contains: + - 'healthz=ok' + # Nested CLI suites. Loader in sling_cli_test.go merges these in. - suite: suite.cli.assist.yaml - suite: suite.cli.build.yaml