From 38e419f1bf394ce54b9ce391de60033ce16f58d4 Mon Sep 17 00:00:00 2001 From: Fritz Larco Date: Wed, 23 Sep 2026 14:19:14 -0300 Subject: [PATCH] Add context-aware connection tests and workbench env config - Replace the global SpecEventChn with a typed SpecEvent sent through the request context, so each spec test run is isolated, cancellable, and reports the request index plus iteration state before and after each request. Legacy wire fields stay for released spec inspectors. - Add Connection.TestWithOptions and ConnEntries.TestWithOptions (endpoints, limit, max requests, context, trace, OnEvent) plus a spec-file overlay that tests a draft spec without mutating the connection entry. Test() still reads the SLING_TEST_* env vars, so the CLI is unchanged. - Emit error events and return context.Canceled promptly on cancellation instead of reporting success. - Add the workbench block to env.yaml, run `sling serve workbench --no-browser` as a brew service, and cover the command with CLI smoke tests. --- .gitignore | 2 +- .goreleaser.mac.yaml | 9 +- core/dbio/api/spec.go | 120 ++++++++-- core/dbio/connection/connection_discover.go | 148 ++++++++++--- core/dbio/connection/connection_local.go | 59 ++++- core/dbio/connection/connection_test.go | 230 ++++++++++++++++++++ core/env/envfile.go | 26 ++- go.mod | 5 +- go.sum | 10 +- tests/suite.cli.yaml | 35 +++ 10 files changed, 582 insertions(+), 62 deletions(-) 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