Merge branch 'master' into fb-1188-ttl

This commit is contained in:
tgruben 2022-03-04 11:30:46 -06:00 • committed by GitHub
commit 0e819b7d7a
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
5 changed files with 176 additions and 52 deletions

8
api.go
View file

@ -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
}

View file

@ -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()

View file

@ -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")

View file

@ -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

View file

@ -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
}