diff --git a/api.go b/api.go index e0c4236f9..6149c6aee 100644 --- a/api.go +++ b/api.go @@ -328,11 +328,17 @@ func (api *API) CreateField(ctx context.Context, indexName string, fieldName str } // Create field. - field, err := index.CreateFieldAndBroadcast(cfm) + field, err := index.CreateField(fieldName, opts...) if err != nil { return nil, errors.Wrap(err, "creating field") } + // Send the create field message to all nodes. We do this *outside* the + // CreateField logic so we're not blocking on it. + if err := api.holder.sendOrSpool(cfm); err != nil { + return nil, errors.Wrap(err, "sending CreateField message") + } + api.holder.Stats.CountWithCustomTags(MetricCreateField, 1, 1.0, []string{fmt.Sprintf("index:%s", indexName)}) return field, nil } diff --git a/api_test.go b/api_test.go index 27e6de36f..0ebf16b5f 100644 --- a/api_test.go +++ b/api_test.go @@ -29,6 +29,8 @@ import ( "github.com/molecula/featurebase/v3/shardwidth" "github.com/molecula/featurebase/v3/test" . "github.com/molecula/featurebase/v3/vprint" // nolint:staticcheck + + "golang.org/x/sync/errgroup" ) func TestAPI_Import(t *testing.T) { @@ -1403,6 +1405,42 @@ func TestVariousApiTranslateCalls(t *testing.T) { } } +func TestAPI_CreateField(t *testing.T) { + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + c := test.MustRunCluster(t, 3) + defer c.Close() + + nodes := make([]*test.Command, 3) + for i := range nodes { + nodes[i] = c.GetNode(i) + } + + if _, err := nodes[0].API.CreateIndex(ctx, "i", pilosa.IndexOptions{}); err != nil { + t.Fatal(err) + } + eg, ctx := errgroup.WithContext(context.Background()) + for _, n := range nodes { + node := n + eg.Go(func() error { + for i := 0; i < 10; i++ { + _, err := node.API.CreateField(ctx, "i", fmt.Sprintf("f%d", i)) + if err != nil && !errors.Is(err, pilosa.ErrFieldExists) { + return err + } + } + return nil + }) + } + err := eg.Wait() + if err != nil { + if errors.Is(err, pilosa.ErrFieldExists) { + t.Fatalf("conflict error: %v", err) + } + t.Fatalf("unexpected error: %T %v", err, err) + } +} + func TestAPI_RBFDebugInfo(t *testing.T) { ctx, cancel := context.WithCancel(context.Background()) defer cancel() diff --git a/index.go b/index.go index 66932cd32..12e0aee68 100644 --- a/index.go +++ b/index.go @@ -503,12 +503,21 @@ func (i *Index) CreateField(name string, opts ...FieldOption) (*Field, error) { return nil, errors.Wrap(err, "validating name") } - i.mu.Lock() - defer i.mu.Unlock() + // Grab lock, check for field existing, release lock. We don't want + // to stay holding the lock, but we might care about the ErrFieldExists + // part of this. + err = func() error { + i.mu.Lock() + defer i.mu.Unlock() - // Ensure field doesn't already exist. - if i.fields[name] != nil { - return nil, newConflictError(ErrFieldExists) + // Ensure field doesn't already exist. + if i.fields[name] != nil { + return newConflictError(ErrFieldExists) + } + return nil + }() + if err != nil { + return nil, err } // Apply and validate functional options. @@ -524,37 +533,26 @@ func (i *Index) CreateField(name string, opts ...FieldOption) (*Field, error) { Meta: fo, } - // Create the field in etcd as the system of record. + // Create the field in etcd as the system of record. We do this without + // the lock held because it can take an arbitrary amount of time... if err := i.persistField(context.Background(), cfm); err != nil { return nil, errors.Wrap(err, "persisting field") } - return i.createField(cfm, false) -} - -// CreateFieldAndBroadcast creates a field locally, then broadcasts the -// creation to other nodes so they can create locally as well. An error is -// returned if the field already exists. -func (i *Index) CreateFieldAndBroadcast(cfm *CreateFieldMessage) (*Field, error) { - err := ValidateName(cfm.Field) - if err != nil { - return nil, errors.Wrap(err, "validating name") - } - + // This is identical to the previous check, because we could get super + // unlucky and have the persist-field thing happen, and somehow the field + // gets created, before we get to run again, and the specific nature of + // the error can matter to the backend. i.mu.Lock() defer i.mu.Unlock() // Ensure field doesn't already exist. - if i.fields[cfm.Field] != nil { + if i.fields[name] != nil { return nil, newConflictError(ErrFieldExists) } - // Create the field in etcd as the system of record. - if err := i.persistField(context.Background(), cfm); err != nil { - return nil, errors.Wrap(err, "persisting field") - } - - return i.createField(cfm, true) + // Actually do the internal bookkeeping. + return i.createField(cfm) } // CreateFieldIfNotExists creates a field with the given options if it doesn't exist. @@ -594,7 +592,7 @@ func (i *Index) CreateFieldIfNotExists(name string, opts ...FieldOption) (*Field return nil, errors.Wrap(err, "persisting field") } - return i.createField(cfm, false) + return i.createField(cfm) } // CreateFieldIfNotExistsWithOptions is a method which I created because I @@ -632,7 +630,7 @@ func (i *Index) CreateFieldIfNotExistsWithOptions(name string, opt *FieldOptions return nil, errors.Wrap(err, "persisting field") } - return i.createField(cfm, false) + return i.createField(cfm) } // persistField stores the field information in etcd. @@ -667,14 +665,14 @@ func (i *Index) createFieldIfNotExists(cfm *CreateFieldMessage) (*Field, error) return f, nil } - return i.createField(cfm, false) + return i.createField(cfm) } -// createField, in addition to creating a new Field, calls Field.Open which -// potentially aquires a lock on Index. So until/unless we refactor the -// Index.createField() function call path, we cannot call Index.createField -// while holding an Index lock. -func (i *Index) createField(cfm *CreateFieldMessage, broadcast bool) (*Field, error) { +// createField does the internal field creation logic, creating the in-memory +// data structure, and kicking translation sync if appropriate. It does not +// notify other nodes; that's done from the API's initial CreateField call +// now. +func (i *Index) createField(cfm *CreateFieldMessage) (*Field, error) { opt := cfm.Meta if opt == nil { opt = &FieldOptions{} @@ -711,13 +709,6 @@ func (i *Index) createField(cfm *CreateFieldMessage, broadcast bool) (*Field, er // enable Txf to find the index in field_test.go TestField_SetValue f.idx = i - if broadcast { - // Send the create field message to all nodes. - if err := i.holder.sendOrSpool(cfm); err != nil { - return nil, errors.Wrap(err, "sending CreateField message") - } - } - // Kick off the field's translation sync process. if err := i.translationSyncer.Reset(); err != nil { return nil, errors.Wrap(err, "resetting translation syncer") diff --git a/pilosa.go b/pilosa.go index 9cf4f715f..218e37491 100644 --- a/pilosa.go +++ b/pilosa.go @@ -111,6 +111,12 @@ func newConflictError(err error) ConflictError { return ConflictError{err} } +// Unwrap makes it so that a ConflictError wrapping ErrFieldExists gets a +// true from errors.Is(ErrFieldExists). +func (c ConflictError) Unwrap() error { + return c.error +} + // NotFoundError wraps an error value to signify that a resource was not found // such that in an HTTP scenario, http.StatusNotFound would be returned. type NotFoundError error diff --git a/server/grpc_test.go b/server/grpc_test.go index 782063859..1124f332f 100644 --- a/server/grpc_test.go +++ b/server/grpc_test.go @@ -534,9 +534,10 @@ func TestQuerySQL(t *testing.T) { {"color", "[]string"}, {"height", "int64"}, {"score", "int64"}, + {"timestamp", "timestamp"}, }, rows: []row{ - {[]columnResponse{uint64(2), int64(16), []string{"blue"}, int64(30), int64(-8)}}, + {[]columnResponse{uint64(2), int64(16), []string{"blue"}, int64(30), int64(-8), "2011-01-02T12:32:00Z"}}, }, }, eq: equal, @@ -551,18 +552,19 @@ func TestQuerySQL(t *testing.T) { {"color", "[]string"}, {"height", "int64"}, {"score", "int64"}, + {"timestamp", "timestamp"}, }, rows: []row{ - {[]columnResponse{uint64(1), int64(27), []string{"blue"}, int64(20), int64(-10)}}, - {[]columnResponse{uint64(2), int64(16), []string{"blue"}, int64(30), int64(-8)}}, - {[]columnResponse{uint64(3), int64(19), []string{"red"}, int64(40), int64(6)}}, - {[]columnResponse{uint64(4), int64(27), []string{"green"}, int64(50), int64(0)}}, - {[]columnResponse{uint64(5), int64(16), []string{"blue"}, int64(60), int64(-2)}}, - {[]columnResponse{uint64(6), int64(34), []string{"blue"}, int64(70), int64(100)}}, - {[]columnResponse{uint64(7), int64(27), []string{"blue"}, int64(80), int64(0)}}, - {[]columnResponse{uint64(8), int64(16), []string{}, int64(90), int64(-13)}}, - {[]columnResponse{uint64(9), int64(16), []string{"red"}, int64(100), int64(80)}}, - {[]columnResponse{uint64(10), int64(31), []string{"red"}, int64(110), int64(-2)}}, + {[]columnResponse{uint64(1), int64(27), []string{"blue"}, int64(20), int64(-10), "2011-04-02T12:32:00Z"}}, + {[]columnResponse{uint64(2), int64(16), []string{"blue"}, int64(30), int64(-8), "2011-01-02T12:32:00Z"}}, + {[]columnResponse{uint64(3), int64(19), []string{"red"}, int64(40), int64(6), "2012-01-02T12:32:00Z"}}, + {[]columnResponse{uint64(4), int64(27), []string{"green"}, int64(50), int64(0), "2013-09-02T12:32:00Z"}}, + {[]columnResponse{uint64(5), int64(16), []string{"blue"}, int64(60), int64(-2), "2014-01-02T12:32:00Z"}}, + {[]columnResponse{uint64(6), int64(34), []string{"blue"}, int64(70), int64(100), "2010-05-02T12:32:00Z"}}, + {[]columnResponse{uint64(7), int64(27), []string{"blue"}, int64(80), int64(0), "2016-08-02T12:32:00Z"}}, + {[]columnResponse{uint64(8), int64(16), []string{}, int64(90), int64(-13), "2020-01-02T12:32:00Z"}}, + {[]columnResponse{uint64(9), int64(16), []string{"red"}, int64(100), int64(80), "2000-03-02T12:32:00Z"}}, + {[]columnResponse{uint64(10), int64(31), []string{"red"}, int64(110), int64(-2), "2018-01-02T12:32:00Z"}}, }, }, eq: equal, @@ -837,6 +839,64 @@ func TestQuerySQL(t *testing.T) { }, eq: equal, }, + { + // GroupBy(Rows(field='age'),Rows(field='height'),filter=Intersect(Row(timestamp>"2017-09-02T12:32:00Z"),Row(height>40))) + sql: "select age, height from grouper where timestamp > '2017-09-02T12:32:00Z' and height > 40 group by age, height", + exp: tableResponse{ + headers: []columnInfo{ + {"age", "int64"}, + {"height", "int64"}, + }, + rows: []row{ + {[]columnResponse{int64(16), int64(90)}}, + {[]columnResponse{int64(31), int64(110)}}, + }, + }, + eq: equalUnordered, + }, + { + // Extract(Union(Row(timestamp>"2017-09-02T12:32:00Z"),Row(height>90)),Rows(age), Rows(height)) + sql: "select age, height from grouper where timestamp > '2017-09-02T12:32:00Z' or height > 90", + exp: tableResponse{ + headers: []columnInfo{ + {"age", "int64"}, + {"height", "int64"}, + }, + rows: []row{ + {[]columnResponse{int64(16), int64(90)}}, + {[]columnResponse{int64(16), int64(100)}}, + {[]columnResponse{int64(31), int64(110)}}, + }, + }, + eq: equalUnordered, + }, + { + //Extract(Intersect(Row(timestamp>"2017-09-02T12:32:00Z"),Row(timestamp<"2019-09-02T12:32:00Z")),Rows(age), Rows(height)) + sql: "select age, height from grouper where timestamp > '2017-09-02T12:32:00Z' and timestamp < '2019-09-02T12:32:00Z'", + exp: tableResponse{ + headers: []columnInfo{ + {"age", "int64"}, + {"height", "int64"}, + }, + rows: []row{ + {[]columnResponse{int64(31), int64(110)}}, + }, + }, + eq: equalUnordered, + }, + { + //Distinct(Row(timestamp>"2019-09-02T12:32:00Z"), index='grouper',field='age') + sql: "select distinct age from grouper where timestamp > '2019-09-02T12:32:00Z'", + exp: tableResponse{ + headers: []columnInfo{ + {"age", "int64"}, + }, + rows: []row{ + {[]columnResponse{int64(16)}}, + }, + }, + eq: equalUnordered, + }, { sql: "show tables", exp: tableResponse{ @@ -865,6 +925,7 @@ func TestQuerySQL(t *testing.T) { {[]columnResponse{"color", "keyed-set"}}, {[]columnResponse{"height", "int"}}, {[]columnResponse{"score", "int"}}, + {[]columnResponse{"timestamp", "timestamp"}}, }, }, eq: equal, @@ -1558,6 +1619,26 @@ func setUpTestQuerySQLUnary(ctx context.Context, t *testing.T) (gh *server.GRPCH t.Fatal(err) } } + m.MustCreateField(t, grouper.Name(), "timestamp", pilosa.OptFieldTypeTimestamp(pilosa.DefaultEpoch, pilosa.TimeUnitSeconds)) + for id, timestamp := range map[int]string{ + 1: "2011-04-02T12:32:00Z", + 2: "2011-01-02T12:32:00Z", + 3: "2012-01-02T12:32:00Z", + 4: "2013-09-02T12:32:00Z", + 5: "2014-01-02T12:32:00Z", + 6: "2010-05-02T12:32:00Z", + 7: "2016-08-02T12:32:00Z", + 8: "2020-01-02T12:32:00Z", + 9: "2000-03-02T12:32:00Z", + 10: "2018-01-02T12:32:00Z", + } { + if _, err := gh.QueryPQLUnary(ctx, &pb.QueryPQLRequest{ + Index: grouper.Name(), + Pql: fmt.Sprintf("Set(%d, timestamp=\"%s\")", id, timestamp), + }); err != nil { + t.Fatal(err) + } + } // joiner joiner := m.MustCreateIndex(t, "joiner", pilosa.IndexOptions{TrackExistence: true}) @@ -1657,6 +1738,8 @@ func toTableResponse(resp *pb.TableResponse) tableResponse { tr.rows[i].columns[j] = v.Float64Val case *pb.ColumnResponse_DecimalVal: tr.rows[i].columns[j] = pql.NewDecimal(v.DecimalVal.Value, v.DecimalVal.Scale) + case *pb.ColumnResponse_TimestampVal: + tr.rows[i].columns[j] = v.TimestampVal default: tr.rows[i].columns[j] = nil }