Merge branch 'master' into fix-logger

This commit is contained in:
tgruben 2018-10-11 06:58:17 -05:00 • committed by GitHub
commit 1895233ed8
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
17 changed files with 1917 additions and 1342 deletions

78
api.go
View file

@ -102,11 +102,9 @@ func (api *API) Query(ctx context.Context, req *QueryRequest) (QueryResponse, er
return QueryResponse{}, errors.Wrap(err, "validating api method")
}
resp := QueryResponse{}
q, err := pql.NewParser(strings.NewReader(req.Query)).Parse()
if err != nil {
return resp, errors.Wrap(err, "parsing")
return QueryResponse{}, errors.Wrap(err, "parsing")
}
execOpts := &execOptions{
Remote: req.Remote,
@ -114,70 +112,14 @@ func (api *API) Query(ctx context.Context, req *QueryRequest) (QueryResponse, er
ExcludeColumns: req.ExcludeColumns, // NOTE: Kept for Pilosa 1.x compat.
ColumnAttrs: req.ColumnAttrs, // NOTE: Kept for Pilosa 1.x compat.
}
results, err := api.server.executor.Execute(ctx, req.Index, q, req.Shards, execOpts)
resp, err := api.server.executor.Execute(ctx, req.Index, q, req.Shards, execOpts)
if err != nil {
return resp, errors.Wrap(err, "executing")
return QueryResponse{}, errors.Wrap(err, "executing")
}
resp.Results = results
// Fill column attributes if requested.
// execOpts.ColumnAttrs may be set by the Execute method if any of the Calls use Options(columnAttrs=true)
if execOpts.ColumnAttrs {
// Consolidate all column ids across all calls.
var columnIDs []uint64
for _, result := range results {
bm, ok := result.(*Row)
if !ok {
continue
}
columnIDs = uint64Slice(columnIDs).merge(bm.Columns())
}
// Retrieve column attributes across all calls.
columnAttrSets, err := api.readColumnAttrSets(api.holder.Index(req.Index), columnIDs)
if err != nil {
return resp, errors.Wrap(err, "reading column attrs")
}
// Translate column attributes, if necessary.
if api.holder.translateFile != nil {
for _, col := range resp.ColumnAttrSets {
v, err := api.holder.translateFile.TranslateColumnToString(req.Index, col.ID)
if err != nil {
return resp, err
}
col.Key, col.ID = v, 0
}
}
resp.ColumnAttrSets = columnAttrSets
}
return resp, nil
}
// readColumnAttrSets returns a list of column attribute objects by id.
func (api *API) readColumnAttrSets(index *Index, ids []uint64) ([]*ColumnAttrSet, error) {
if index == nil {
return nil, nil
}
ax := make([]*ColumnAttrSet, 0, len(ids))
for _, id := range ids {
// Read attributes for column. Skip column if empty.
attrs, err := index.ColumnAttrStore().Attrs(id)
if err != nil {
return nil, errors.Wrap(err, "getting attrs")
} else if len(attrs) == 0 {
continue
}
// Append column with attributes.
ax = append(ax, &ColumnAttrSet{ID: id, Attrs: attrs})
}
return ax, nil
}
// CreateIndex makes a new Pilosa index.
func (api *API) CreateIndex(_ context.Context, indexName string, options IndexOptions) (*Index, error) {
if err := api.validate(apiCreateIndex); err != nil {
@ -321,13 +263,19 @@ func (api *API) ImportRoaring(ctx context.Context, indexName, fieldName string,
nodes := api.cluster.shardNodes(indexName, shard)
var eg errgroup.Group
field := api.holder.Field(indexName, fieldName)
if field == nil {
return newNotFoundError(ErrFieldNotFound)
}
// only set fields are supported
if field.Type() != FieldTypeSet {
return NewBadRequestError(errors.New("roaring import is only supported for set fields"))
}
for _, node := range nodes {
node := node
if node.ID == api.server.nodeID {
field := api.holder.Field(indexName, fieldName)
if field == nil {
return newNotFoundError(ErrFieldNotFound)
}
// must make a copy of data to operate on locally. field.importRoaring changes data
d2 := make([]byte, len(data))
copy(d2, data)

View file

@ -174,6 +174,11 @@ func (d *diagnosticsCollector) logErr(err error) bool {
return false
}
// EnrichWithCPUInfo adds CPU information to the diagnostics payload.
func (d *diagnosticsCollector) EnrichWithCPUInfo() {
d.Set("CPUArch", d.server.systemInfo.CPUArch())
}
// EnrichWithOSInfo adds OS information to the diagnostics payload.
func (d *diagnosticsCollector) EnrichWithOSInfo() {
uptime, err := d.server.systemInfo.Uptime()
@ -265,6 +270,7 @@ type SystemInfo interface {
MemFree() (uint64, error)
MemTotal() (uint64, error)
MemUsed() (uint64, error)
CPUArch() string
}
// newNopSystemInfo creates a no-op implementation of SystemInfo.
@ -315,3 +321,8 @@ func (n *nopSystemInfo) MemTotal() (uint64, error) {
func (n *nopSystemInfo) MemUsed() (uint64, error) {
return 0, nil
}
// CPUArch returns the CPU architecture, such as amd64
func (n *nopSystemInfo) CPUArch() string {
return ""
}

View file

@ -266,6 +266,40 @@ ClearRow(stargazer=1)
This represents removing the relationship between the user with id=1 and all repositories.
#### Store
**Spec:**
```
Store(<ROW_CALL>, <FIELD>=<ROW>)
```
**Description:**
`Store` writes the results of <ROW_CALL> to the specified row. If the row already exists, it will be replaced. The destination field must be of field type `set`.
**Result Type:** boolean
Upon success, this method always returns `true`. A future version of Pilosa may use this boolean result to indicate whether or not the data in the destination row was changed by the `Store` call.
**Examples:**
Store the contents of stargazer row 1 into stargazer row 2:
```request
Store(Row(stargazer=1), stargazer=2)
```
```response
{"results":[true]}
```
Store the results of the intersection of stargazer rows 10 and 11 into stargazer row 20.
```request
Store(Intersect(Row(stargazer=10), Row(stargazer=11)), stargazer=20)
```
```response
{"results":[true]}
```
### Read Operations
#### Row

View file

@ -79,20 +79,21 @@ func newExecutor(opts ...executorOption) *executor {
}
// Execute executes a PQL query.
func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shards []uint64, opt *execOptions) ([]interface{}, error) {
func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shards []uint64, opt *execOptions) (QueryResponse, error) {
resp := QueryResponse{}
// Verify that an index is set.
if index == "" {
return nil, ErrIndexRequired
return resp, ErrIndexRequired
}
idx := e.Holder.Index(index)
if idx == nil {
return nil, ErrIndexNotFound
return resp, ErrIndexNotFound
}
// Verify that the number of writes do not exceed the maximum.
if e.MaxWritesPerRequest > 0 && q.WriteCallN() > e.MaxWritesPerRequest {
return nil, ErrTooManyWrites
return resp, ErrTooManyWrites
}
// Default options.
@ -105,14 +106,48 @@ func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shar
if !opt.Remote {
for i := range q.Calls {
if err := e.translateCall(index, idx, q.Calls[i]); err != nil {
return nil, err
return resp, err
}
}
}
results, err := e.execute(ctx, index, q, shards, opt)
if err != nil {
return nil, err
return resp, err
}
resp.Results = results
// Fill column attributes if requested.
if opt.ColumnAttrs {
// Consolidate all column ids across all calls.
var columnIDs []uint64
for _, result := range results {
bm, ok := result.(*Row)
if !ok {
continue
}
columnIDs = uint64Slice(columnIDs).merge(bm.Columns())
}
// Retrieve column attributes across all calls.
columnAttrSets, err := e.readColumnAttrSets(e.Holder.Index(index), columnIDs)
if err != nil {
return resp, errors.Wrap(err, "reading column attrs")
}
// Translate column attributes, if necessary.
if idx.Keys() {
for _, col := range columnAttrSets {
v, err := e.Holder.translateFile.TranslateColumnToString(index, col.ID)
if err != nil {
return resp, err
}
col.Key, col.ID = v, 0
}
}
resp.ColumnAttrSets = columnAttrSets
}
// Translate response objects from ids to keys, if necessary.
@ -121,11 +156,35 @@ func (e *executor) Execute(ctx context.Context, index string, q *pql.Query, shar
for i := range results {
results[i], err = e.translateResult(index, idx, q.Calls[i], results[i])
if err != nil {
return nil, err
return resp, err
}
}
}
return results, nil
return resp, nil
}
// readColumnAttrSets returns a list of column attribute objects by id.
func (e *executor) readColumnAttrSets(index *Index, ids []uint64) ([]*ColumnAttrSet, error) {
if index == nil {
return nil, nil
}
ax := make([]*ColumnAttrSet, 0, len(ids))
for _, id := range ids {
// Read attributes for column. Skip column if empty.
attrs, err := index.ColumnAttrStore().Attrs(id)
if err != nil {
return nil, errors.Wrap(err, "getting attrs")
} else if len(attrs) == 0 {
continue
}
// Append column with attributes.
ax = append(ax, &ColumnAttrSet{ID: id, Attrs: attrs})
}
return ax, nil
}
func (e *executor) execute(ctx context.Context, index string, q *pql.Query, shards []uint64, opt *execOptions) ([]interface{}, error) {
@ -184,6 +243,8 @@ func (e *executor) executeCall(ctx context.Context, index string, c *pql.Call, s
return e.executeClearBit(ctx, index, c, opt)
case "ClearRow":
return e.executeClearRow(ctx, index, c, shards, opt)
case "Store":
return e.executeSetRow(ctx, index, c, shards, opt)
case "Count":
e.Holder.Stats.CountWithCustomTags(c.Name, 1, 1.0, []string{indexTag})
return e.executeCount(ctx, index, c, shards, opt)
@ -1213,6 +1274,94 @@ func (e *executor) executeClearRowShard(_ context.Context, index string, c *pql.
return changed, nil
}
// executeSetRow executes a SetRow() call.
func (e *executor) executeSetRow(ctx context.Context, index string, c *pql.Call, shards []uint64, opt *execOptions) (bool, error) {
// Ensure the field type supports SetRow().
fieldName, err := c.FieldArg()
if err != nil {
return false, errors.New("SetRow() argument required: field")
}
field := e.Holder.Field(index, fieldName)
if field == nil {
return false, ErrFieldNotFound
}
if field.Type() != FieldTypeSet {
return false, fmt.Errorf("SetRow() is not supported on %s field types", field.Type())
}
// Execute calls in bulk on each remote node and merge.
mapFn := func(shard uint64) (interface{}, error) {
return e.executeSetRowShard(ctx, index, c, shard)
}
// Merge returned results at coordinating node.
reduceFn := func(prev, v interface{}) interface{} {
val := v.(bool)
if prev == nil {
return val
}
return val || prev.(bool)
}
result, err := e.mapReduce(ctx, index, shards, c, opt, mapFn, reduceFn)
return result.(bool), err
}
// executeSetRowShard executes a SetRow() call for a single shard.
func (e *executor) executeSetRowShard(ctx context.Context, index string, c *pql.Call, shard uint64) (bool, error) {
fieldName, err := c.FieldArg()
if err != nil {
return false, errors.New("SetRow() argument required: field")
}
// Read fields using labels.
rowID, ok, err := c.UintArg(fieldName)
if err != nil {
return false, fmt.Errorf("reading SetRow() row: %v", err)
} else if !ok {
return false, fmt.Errorf("SetRow() row argument '%v' required", rowLabel)
}
field := e.Holder.Field(index, fieldName)
if field == nil {
return false, ErrFieldNotFound
}
// Retrieve source row.
var src *Row
if len(c.Children) == 1 {
row, err := e.executeBitmapCallShard(ctx, index, c.Children[0], shard)
if err != nil {
return false, errors.Wrap(err, "getting source row")
}
src = row
} else {
return false, errors.New("SetRow() requires a source row")
}
// Set the row on the standard view.
changed := false
fragment := e.Holder.fragment(index, fieldName, viewStandard, shard)
if fragment == nil {
// Since the destination fragment doesn't exist, create one.
view, err := field.createViewIfNotExists(viewStandard)
if err != nil {
return false, errors.Wrap(err, "creating view")
}
fragment, err = view.createFragmentIfNotExists(shard)
if err != nil {
return false, errors.Wrapf(err, "creating fragment: %d", shard)
}
}
set, err := fragment.setRow(src, rowID)
if err != nil {
return false, errors.Wrapf(err, "setting row %d on view %s shard %d", rowID, viewStandard, shard)
}
changed = changed || set
return changed, nil
}
// executeSet executes a Set() call.
func (e *executor) executeSet(ctx context.Context, index string, c *pql.Call, opt *execOptions) (bool, error) {
// Read colID.
@ -1708,16 +1857,17 @@ func (e *executor) mapperLocal(ctx context.Context, shards []uint64, mapFn mapFu
func (e *executor) translateCall(index string, idx *Index, c *pql.Call) error {
var colKey, rowKey, fieldName string
if c.Name == "Set" || c.Name == "Clear" || c.Name == "Row" {
switch c.Name {
case "Set", "Clear", "Row", "Range", "SetColumnAttrs":
// Positional args in new PQL syntax require special handling here.
colKey = "_" + columnLabel
fieldName, _ = c.FieldArg()
rowKey = fieldName
} else if c.Name == "SetRowAttrs" {
case "SetRowAttrs":
// Positional args in new PQL syntax require special handling here.
rowKey = "_" + rowLabel
fieldName = callArgString(c, "_field")
} else {
default:
colKey = "col"
fieldName = callArgString(c, "field")
rowKey = "row"

View file

@ -1020,6 +1020,58 @@ func TestExecutor_Execute_Range(t *testing.T) {
})
}
// Ensure a range query with keys can be executed.
func TestExecutor_Execute_Range_WithKeys(t *testing.T) {
c := test.MustRunCluster(t, 1)
defer c.Close()
hldr := test.Holder{Holder: c[0].Server.Holder()}
// Create index.
index := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{})
// Create field.
if _, err := index.CreateFieldIfNotExists("f", pilosa.OptFieldTypeTime(pilosa.TimeQuantum("YMDH")), pilosa.OptFieldKeys()); err != nil {
t.Fatal(err)
}
// Set columns.
cc := `
Set(2, f="foo", 1999-12-31T00:00)
Set(3, f="foo", 2000-01-01T00:00)
Set(4, f="foo", 2000-01-02T00:00)
Set(5, f="foo", 2000-02-01T00:00)
Set(6, f="foo", 2001-01-01T00:00)
Set(7, f="foo", 2002-01-01T02:00)
Set(2, f="foo", 1999-12-30T00:00)
Set(2, f="foo", 2002-02-01T00:00)
Set(2, f="bar", 2001-01-01T00:00)
`
if _, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: cc}); err != nil {
t.Fatal(err)
}
t.Run("Standard", func(t *testing.T) {
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Range(f="foo", 1999-12-31T00:00, 2002-01-01T03:00)`}); err != nil {
t.Fatal(err)
} else if columns := res.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(columns, []uint64{2, 3, 4, 5, 6, 7}) {
t.Fatalf("unexpected columns: %+v", columns)
}
})
t.Run("Clear", func(t *testing.T) {
if _, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Clear( 2, f="foo")`}); err != nil {
t.Fatal(err)
}
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Range(f="foo", 1999-12-31T00:00, 2002-01-01T03:00)`}); err != nil {
t.Fatal(err)
} else if columns := res.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(columns, []uint64{3, 4, 5, 6, 7}) {
t.Fatalf("unexpected columns: %+v", columns)
}
})
}
// Ensure a Range(bsiGroup) query can be executed.
func TestExecutor_Execute_BSIGroupRange(t *testing.T) {
c := test.MustRunCluster(t, 1)
@ -1509,6 +1561,36 @@ func TestExecutor_QueryCall(t *testing.T) {
}
})
t.Run("columnAttrsWithKeys", func(t *testing.T) {
c := test.MustRunCluster(t, 1)
defer c.Close()
hldr := test.Holder{Holder: c[0].Server.Holder()}
// Set columns for rows 0, 10, & 20 across two shards.
if idx, err := hldr.CreateIndex("i", pilosa.IndexOptions{Keys: true}); err != nil {
t.Fatal(err)
} else if _, err := idx.CreateField("f", pilosa.OptFieldKeys()); err != nil {
t.Fatal(err)
} else if _, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `
Set("one-hundred", f="ten")
SetColumnAttrs("one-hundred", foo="bar")
`}); err != nil {
t.Fatal(err)
}
targetColAttrSets := []*pilosa.ColumnAttrSet{
{Key: "one-hundred", Attrs: map[string]interface{}{"foo": "bar"}},
}
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Options(Row(f="ten"), columnAttrs=true)`}); err != nil {
t.Fatal(err)
} else if keys := res.Results[0].(*pilosa.Row).Keys; !reflect.DeepEqual(keys, []string{"one-hundred"}) {
t.Fatalf("unexpected keys: %+v", keys)
} else if attrs := res.ColumnAttrSets; !reflect.DeepEqual(attrs, targetColAttrSets) {
t.Fatalf("unexpected attrs: %s", spew.Sdump(attrs))
}
})
t.Run("shards", func(t *testing.T) {
c := test.MustRunCluster(t, 1)
defer c.Close()
@ -1913,6 +1995,145 @@ func TestExecutor_Execute_ClearRow(t *testing.T) {
})
}
// Ensure a row can be set.
func TestExecutor_Execute_SetRow(t *testing.T) {
t.Run("Set_NewRow", func(t *testing.T) {
c := test.MustRunCluster(t, 1)
defer c.Close()
hldr := test.Holder{Holder: c[0].Server.Holder()}
index := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{TrackExistence: true})
if _, err := index.CreateField("f", pilosa.OptFieldTypeDefault()); err != nil {
t.Fatal(err)
}
if _, err := index.CreateField("tmp", pilosa.OptFieldTypeDefault()); err != nil {
t.Fatal(err)
}
// Set bits.
if _, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `` +
fmt.Sprintf("Set(%d, f=%d)\n", 3, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth-1, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth+1, 10),
}); err != nil {
t.Fatal(err)
}
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Row(f=10)`}); err != nil {
t.Fatal(err)
} else if bits := res.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(bits, []uint64{3, ShardWidth - 1, ShardWidth + 1}) {
t.Fatalf("unexpected columns: %+v", bits)
}
// Store row 10 into a different row.
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Store(Row(f=10), tmp=20)`}); err != nil {
t.Fatal(err)
} else if res := res.Results[0].(bool); !res {
t.Fatalf("unexpected set row result: %+v", res)
}
// Ensure the row was populated.
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Row(tmp=20)`}); err != nil {
t.Fatal(err)
} else if bits := res.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(bits, []uint64{3, ShardWidth - 1, ShardWidth + 1}) {
t.Fatalf("unexpected columns: %+v", bits)
}
})
t.Run("Set_NoSource", func(t *testing.T) {
c := test.MustRunCluster(t, 1)
defer c.Close()
hldr := test.Holder{Holder: c[0].Server.Holder()}
index := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{TrackExistence: true})
_, err := index.CreateField("f", pilosa.OptFieldTypeDefault())
if err != nil {
t.Fatal(err)
}
// Set bits.
if _, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `` +
fmt.Sprintf("Set(%d, f=%d)\n", 3, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth-1, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth+1, 10),
}); err != nil {
t.Fatal(err)
}
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Row(f=10)`}); err != nil {
t.Fatal(err)
} else if bits := res.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(bits, []uint64{3, ShardWidth - 1, ShardWidth + 1}) {
t.Fatalf("unexpected columns: %+v", bits)
}
// Store row 9 (which doesn't exist) into a different row.
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Store(Row(f=9), f=20)`}); err != nil {
t.Fatal(err)
} else if res := res.Results[0].(bool); !res {
t.Fatalf("unexpected set row result: %+v", res)
}
// Ensure the row was populated.
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Row(f=20)`}); err != nil {
t.Fatal(err)
} else if bits := res.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(bits, []uint64{}) {
t.Fatalf("unexpected columns: %+v", bits)
}
// Store row 9 (which doesn't exist) into a row that does exist.
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Store(Row(f=9), f=10)`}); err != nil {
t.Fatal(err)
} else if res := res.Results[0].(bool); !res {
t.Fatalf("unexpected set row result: %+v", res)
}
// Ensure the row was populated.
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Row(f=10)`}); err != nil {
t.Fatal(err)
} else if bits := res.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(bits, []uint64{}) {
t.Fatalf("unexpected columns: %+v", bits)
}
})
t.Run("Set_ExistingDestination", func(t *testing.T) {
c := test.MustRunCluster(t, 1)
defer c.Close()
hldr := test.Holder{Holder: c[0].Server.Holder()}
index := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{TrackExistence: true})
_, err := index.CreateField("f", pilosa.OptFieldTypeDefault())
if err != nil {
t.Fatal(err)
}
// Set bits.
if _, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `` +
fmt.Sprintf("Set(%d, f=%d)\n", 3, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth-1, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth+1, 10) +
fmt.Sprintf("Set(%d, f=%d)\n", 1, 20) +
fmt.Sprintf("Set(%d, f=%d)\n", ShardWidth+1, 20),
}); err != nil {
t.Fatal(err)
}
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Row(f=20)`}); err != nil {
t.Fatal(err)
} else if bits := res.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(bits, []uint64{1, ShardWidth + 1}) {
t.Fatalf("unexpected columns: %+v", bits)
}
// Store row 10 into an existing row.
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Store(Row(f=10), f=20)`}); err != nil {
t.Fatal(err)
} else if res := res.Results[0].(bool); !res {
t.Fatalf("unexpected set row result: %+v", res)
}
// Ensure the row was populated.
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Row(f=20)`}); err != nil {
t.Fatal(err)
} else if bits := res.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(bits, []uint64{3, ShardWidth - 1, ShardWidth + 1}) {
t.Fatalf("unexpected columns: %+v", bits)
}
})
}
func benchmarkExistence(nn bool, b *testing.B) {
c := test.MustRunCluster(b, 1)
defer c.Close()

View file

@ -352,11 +352,11 @@ func (f *fragment) unprotectedRow(rowID uint64) *Row {
}
// Only use a subset of the containers.
// NOTE: The start & end ranges must be divisible by
// NOTE: The start & end ranges must be divisible by container width.
data := f.storage.OffsetRange(f.shard*ShardWidth, rowID*ShardWidth, (rowID+1)*ShardWidth)
// Reference bitmap subrange in storage.
// We Clone() data because otherwise row will contains pointers to containers in storage.
// We Clone() data because otherwise row will contain pointers to containers in storage.
// This causes unexpected results when we cache the row and try to use it later.
row := &Row{
segments: []rowSegment{{
@ -491,6 +491,56 @@ func (f *fragment) unprotectedClearBit(rowID, columnID uint64) (changed bool, er
return changed, nil
}
// setRow replaces an existing row (specified by rowID) with the given
// Row. This updates both the on-disk storage and the in-cache bitmap.
func (f *fragment) setRow(row *Row, rowID uint64) (bool, error) {
f.mu.Lock()
defer f.mu.Unlock()
return f.unprotectedSetRow(row, rowID)
}
func (f *fragment) unprotectedSetRow(row *Row, rowID uint64) (changed bool, err error) {
// TODO: In order to return `changed`, we need to first compare
// the existing row with the given row. Determine if the overhead
// of this is worth having `changed`.
// For now we will assume changed is always true.
changed = true
// First container of the row in storage.
headContainerKey := rowID << shardVsContainerExponent
// Remove every existing container in the row.
for i := uint64(0); i < (1 << shardVsContainerExponent); i++ {
f.storage.Containers.Remove(headContainerKey + i)
}
// From the given row, get the rowSegment for this shard.
seg := row.segment(f.shard)
if seg == nil {
return changed, nil
}
// Put each container from rowSegment to fragment storage.
citer, _ := seg.data.Containers.Iterator(f.shard << shardVsContainerExponent)
for citer.Next() {
k, c := citer.Value()
f.storage.Containers.Put(headContainerKey+(k%(1<<shardVsContainerExponent)), c)
}
// Update the row in cache.
n := f.storage.CountRange(rowID*ShardWidth, (rowID+1)*ShardWidth)
f.cache.BulkAdd(rowID, n)
// Snapshot storage.
if err := f.snapshot(); err != nil {
return false, errors.Wrap(err, "snapshotting")
}
f.stats.Count("setRow", 1, 1.0)
return changed, nil
}
// ClearRow clears a row for a given rowID within the fragment.
// This updates both the on-disk storage and the in-cache bitmap.
func (f *fragment) clearRow(rowID uint64) (bool, error) {

View file

@ -122,6 +122,54 @@ func TestFragment_ClearRow(t *testing.T) {
}
}
// Ensure a fragment can set a row.
func TestFragment_SetRow(t *testing.T) {
f := mustOpenFragment("i", "f", viewStandard, 7, "")
defer f.Close()
rowID := uint64(1000)
// Set bits on the fragment.
if _, err := f.setBit(rowID, 8000001); err != nil {
t.Fatal(err)
} else if _, err := f.setBit(rowID, 8065536); err != nil {
t.Fatal(err)
}
// Verify data on row.
if cols := f.row(rowID).Columns(); !reflect.DeepEqual(cols, []uint64{8000001, 8065536}) {
t.Fatalf("unexpected columns: %+v", cols)
}
// Verify count on row.
if n := f.row(rowID).Count(); n != 2 {
t.Fatalf("unexpected count: %d", n)
}
// Set row (overwrite existing data).
row := NewRow(8000002, 8065537, 8131074)
if changed, err := f.unprotectedSetRow(row, rowID); err != nil {
t.Fatal(err)
} else if !changed {
t.Fatalf("expected changed value: %v", changed)
}
// Verify data on row.
if cols := f.row(rowID).Columns(); !reflect.DeepEqual(cols, []uint64{8000002, 8065537, 8131074}) {
t.Fatalf("unexpected columns after set row: %+v", cols)
}
// Verify count on row.
if n := f.row(rowID).Count(); n != 3 {
t.Fatalf("unexpected count after set row: %d", n)
}
// Close and reopen the fragment & verify the data.
if err := f.reopen(); err != nil {
t.Fatal(err)
} else if n := f.row(rowID).Count(); n != 3 {
t.Fatalf("unexpected count (reopen): %d", n)
}
}
// Ensure a fragment can set & read a value.
func TestFragment_SetValue(t *testing.T) {
t.Run("OK", func(t *testing.T) {

View file

@ -15,6 +15,8 @@
package gopsutil
import (
"runtime"
"github.com/pilosa/pilosa"
"github.com/shirou/gopsutil/host"
"github.com/shirou/gopsutil/mem"
@ -109,6 +111,11 @@ func (s *systemInfo) KernelVersion() (string, error) {
return host.KernelVersion()
}
// CPUArch returns the CPU architecture, such as amd64
func (s *systemInfo) CPUArch() string {
return runtime.GOARCH
}
// NewSystemInfo is a constructor for the gopsutil implementation of SystemInfo.
func NewSystemInfo() *systemInfo {
return &systemInfo{}

View file

@ -15,7 +15,6 @@
package gopsutil_test
import (
"log"
"testing"
"github.com/pilosa/pilosa"
@ -25,15 +24,6 @@ import (
func TestSystemInfo(t *testing.T) {
var systemInfo pilosa.SystemInfo = gopsutil.NewSystemInfo()
// Uptime()(uint64, error)
// Platform()(string, error)
// Family()(string, error)
// OSVersion()(string, error)
// KernelVersion()(string, error)
// MemFree()(uint64, error)
// MemTotal()(uint64, error)
// MemUsed()(uint64, error)
//
uptime, err := systemInfo.Uptime()
if err != nil || uptime == 0 {
t.Fatalf("Error collecting uptime (error: %v)", err)
@ -70,8 +60,12 @@ func TestSystemInfo(t *testing.T) {
}
memtotal, err := systemInfo.MemTotal()
log.Println(memtotal)
if err != nil {
t.Fatalf("Error getting memtotal. (memtotal: %v, error: %v)", memtotal, err)
}
cpuArch := systemInfo.CPUArch()
if cpuArch == "" {
t.Fatalf("Error getting CPU arch.")
}
}

View file

@ -170,13 +170,34 @@ func (h *Handler) Close() error {
func (h *Handler) populateValidators() {
h.validators = map[string]*queryValidationSpec{}
h.validators["GetFragmentNodes"] = queryValidationSpecRequired("shard", "index")
h.validators["GetShardMax"] = queryValidationSpecRequired()
h.validators["PostQuery"] = queryValidationSpecRequired().Optional("shards", "columnAttrs", "excludeRowAttrs", "excludeColumns")
h.validators["Home"] = queryValidationSpecRequired()
h.validators["PostClusterResizeAbort"] = queryValidationSpecRequired()
h.validators["PostClusterResizeRemoveNode"] = queryValidationSpecRequired()
h.validators["PostClusterResizeSetCoordinator"] = queryValidationSpecRequired()
h.validators["GetExport"] = queryValidationSpecRequired("index", "field", "shard")
h.validators["GetFragmentData"] = queryValidationSpecRequired("index", "field", "shard")
h.validators["PostFragmentData"] = queryValidationSpecRequired("index", "field", "shard")
h.validators["GetIndexes"] = queryValidationSpecRequired()
h.validators["GetIndex"] = queryValidationSpecRequired()
h.validators["PostIndex"] = queryValidationSpecRequired()
h.validators["DeleteIndex"] = queryValidationSpecRequired()
h.validators["PostField"] = queryValidationSpecRequired()
h.validators["DeleteField"] = queryValidationSpecRequired()
h.validators["PostImport"] = queryValidationSpecRequired()
h.validators["PostImportRoaring"] = queryValidationSpecRequired().Optional("remote")
h.validators["PostQuery"] = queryValidationSpecRequired().Optional("shards", "columnAttrs", "excludeRowAttrs", "excludeColumns")
h.validators["GetInfo"] = queryValidationSpecRequired()
h.validators["RecalculateCaches"] = queryValidationSpecRequired()
h.validators["GetSchema"] = queryValidationSpecRequired()
h.validators["GetStatus"] = queryValidationSpecRequired()
h.validators["GetVersion"] = queryValidationSpecRequired()
h.validators["PostClusterMessage"] = queryValidationSpecRequired()
h.validators["GetFragmentBlockData"] = queryValidationSpecRequired()
h.validators["GetFragmentBlocks"] = queryValidationSpecRequired("index", "field", "view", "shard")
h.validators["GetFragmentNodes"] = queryValidationSpecRequired("shard", "index")
h.validators["PostIndexAttrDiff"] = queryValidationSpecRequired()
h.validators["PostFieldAttrDiff"] = queryValidationSpecRequired()
h.validators["GetNodes"] = queryValidationSpecRequired()
h.validators["GetShardMax"] = queryValidationSpecRequired()
h.validators["GetTranslateData"] = queryValidationSpecRequired("offset")
}
func (h *Handler) queryArgValidator(next http.Handler) http.Handler {
@ -203,40 +224,40 @@ func (h *Handler) queryArgValidator(next http.Handler) http.Handler {
// newRouter creates a new mux http router.
func newRouter(handler *Handler) *mux.Router {
router := mux.NewRouter()
router.HandleFunc("/", handler.handleHome).Methods("GET")
router.HandleFunc("/cluster/resize/abort", handler.handlePostClusterResizeAbort).Methods("POST")
router.HandleFunc("/cluster/resize/remove-node", handler.handlePostClusterResizeRemoveNode).Methods("POST")
router.HandleFunc("/cluster/resize/set-coordinator", handler.handlePostClusterResizeSetCoordinator).Methods("POST")
router.HandleFunc("/", handler.handleHome).Methods("GET").Name("Home")
router.HandleFunc("/cluster/resize/abort", handler.handlePostClusterResizeAbort).Methods("POST").Name("PostClusterResizeAbort")
router.HandleFunc("/cluster/resize/remove-node", handler.handlePostClusterResizeRemoveNode).Methods("POST").Name("PostClusterResizeRemoveNode")
router.HandleFunc("/cluster/resize/set-coordinator", handler.handlePostClusterResizeSetCoordinator).Methods("POST").Name("PostClusterResizeSetCoordinator")
router.PathPrefix("/debug/pprof/").Handler(http.DefaultServeMux).Methods("GET")
router.Handle("/debug/vars", expvar.Handler()).Methods("GET")
router.HandleFunc("/export", handler.handleGetExport).Methods("GET").Name("GetExport")
router.HandleFunc("/index", handler.handleGetIndexes).Methods("GET")
router.HandleFunc("/index/{index}", handler.handleGetIndex).Methods("GET")
router.HandleFunc("/index/{index}", handler.handlePostIndex).Methods("POST")
router.HandleFunc("/index/{index}", handler.handleDeleteIndex).Methods("DELETE")
router.HandleFunc("/index", handler.handleGetIndexes).Methods("GET").Name("GetIndexes")
router.HandleFunc("/index/{index}", handler.handleGetIndex).Methods("GET").Name("GetIndex")
router.HandleFunc("/index/{index}", handler.handlePostIndex).Methods("POST").Name("PostIndex")
router.HandleFunc("/index/{index}", handler.handleDeleteIndex).Methods("DELETE").Name("DeleteIndex")
//router.HandleFunc("/index/{index}/field", handler.handleGetFields).Methods("GET") // Not implemented.
router.HandleFunc("/index/{index}/field/{field}", handler.handlePostField).Methods("POST")
router.HandleFunc("/index/{index}/field/{field}", handler.handleDeleteField).Methods("DELETE")
router.HandleFunc("/index/{index}/field/{field}/import", handler.handlePostImport).Methods("POST")
router.HandleFunc("/index/{index}/field/{field}/import-roaring/{shard}", handler.handlePostImportRoaring).Methods("POST")
router.HandleFunc("/index/{index}/field/{field}", handler.handlePostField).Methods("POST").Name("PostField")
router.HandleFunc("/index/{index}/field/{field}", handler.handleDeleteField).Methods("DELETE").Name("DeleteField")
router.HandleFunc("/index/{index}/field/{field}/import", handler.handlePostImport).Methods("POST").Name("PostImport")
router.HandleFunc("/index/{index}/field/{field}/import-roaring/{shard}", handler.handlePostImportRoaring).Methods("POST").Name("PostImportRoaring")
router.HandleFunc("/index/{index}/query", handler.handlePostQuery).Methods("POST").Name("PostQuery")
router.HandleFunc("/info", handler.handleGetInfo).Methods("GET")
router.HandleFunc("/recalculate-caches", handler.handleRecalculateCaches).Methods("POST")
router.HandleFunc("/schema", handler.handleGetSchema).Methods("GET")
router.HandleFunc("/status", handler.handleGetStatus).Methods("GET")
router.HandleFunc("/version", handler.handleGetVersion).Methods("GET")
router.HandleFunc("/info", handler.handleGetInfo).Methods("GET").Name("GetInfo")
router.HandleFunc("/recalculate-caches", handler.handleRecalculateCaches).Methods("POST").Name("RecalculateCaches")
router.HandleFunc("/schema", handler.handleGetSchema).Methods("GET").Name("GetSchema")
router.HandleFunc("/status", handler.handleGetStatus).Methods("GET").Name("GetStatus")
router.HandleFunc("/version", handler.handleGetVersion).Methods("GET").Name("GetVersion")
// /internal endpoints are for internal use only; they may change at any time.
// DO NOT rely on these for external applications!
router.HandleFunc("/internal/cluster/message", handler.handlePostClusterMessage).Methods("POST")
router.HandleFunc("/internal/fragment/block/data", handler.handleGetFragmentBlockData).Methods("GET")
router.HandleFunc("/internal/cluster/message", handler.handlePostClusterMessage).Methods("POST").Name("PostClusterMessage")
router.HandleFunc("/internal/fragment/block/data", handler.handleGetFragmentBlockData).Methods("GET").Name("GetFragmentBlockData")
router.HandleFunc("/internal/fragment/blocks", handler.handleGetFragmentBlocks).Methods("GET").Name("GetFragmentBlocks")
router.HandleFunc("/internal/fragment/nodes", handler.handleGetFragmentNodes).Methods("GET").Name("GetFragmentNodes")
router.HandleFunc("/internal/index/{index}/attr/diff", handler.handlePostIndexAttrDiff).Methods("POST")
router.HandleFunc("/internal/index/{index}/field/{field}/attr/diff", handler.handlePostFieldAttrDiff).Methods("POST")
router.HandleFunc("/internal/index/{index}/attr/diff", handler.handlePostIndexAttrDiff).Methods("POST").Name("PostIndexAttrDiff")
router.HandleFunc("/internal/index/{index}/field/{field}/attr/diff", handler.handlePostFieldAttrDiff).Methods("POST").Name("PostFieldAttrDiff")
router.HandleFunc("/internal/nodes", handler.handleGetNodes).Methods("GET").Name("GetNodes")
router.HandleFunc("/internal/shards/max", handler.handleGetShardsMax).Methods("GET") // TODO: deprecate, but it's being used by the client
router.HandleFunc("/internal/translate/data", handler.handleGetTranslateData).Methods("GET")
router.HandleFunc("/internal/shards/max", handler.handleGetShardsMax).Methods("GET").Name("GetShardsMax") // TODO: deprecate, but it's being used by the client
router.HandleFunc("/internal/translate/data", handler.handleGetTranslateData).Methods("GET").Name("GetTranslateData")
// TODO: Apply MethodNotAllowed statuses to all endpoints.
// Ideally this would be automatic, as described in this (wontfix) ticket:
@ -521,7 +542,12 @@ func (p *postIndexRequest) UnmarshalJSON(b []byte) error {
return err
}
// Unmarshal expected values.
var _p _postIndexRequest
_p := _postIndexRequest{
Options: pilosa.IndexOptions{
Keys: false,
TrackExistence: true,
},
}
if err := json.Unmarshal(b, &_p); err != nil {
return errors.Wrap(err, "unmarshalling expected values")
}
@ -597,7 +623,12 @@ func (h *Handler) handlePostIndex(w http.ResponseWriter, r *http.Request) {
resp := successResponse{}
// Decode request.
var req postIndexRequest
req := postIndexRequest{
Options: pilosa.IndexOptions{
Keys: false,
TrackExistence: true,
},
}
err := json.NewDecoder(r.Body).Decode(&req)
if err != nil && err != io.EOF {
resp.write(w, err)
@ -1473,11 +1504,16 @@ func (h *Handler) handlePostImportRoaring(w http.ResponseWriter, r *http.Request
return
}
resp := &pilosa.ImportResponse{}
// TODO give meaningful stats for import
err = h.api.ImportRoaring(r.Context(), urlVars["index"], urlVars["field"], shard, remote, body)
resp := &pilosa.ImportResponse{}
if err != nil {
resp.Err = err.Error()
if _, ok := err.(pilosa.BadRequestError); ok {
w.WriteHeader(http.StatusBadRequest)
} else {
w.WriteHeader(http.StatusInternalServerError)
}
}
// Marshal response object.
buf, err := h.api.Serializer.Marshal(resp)

View file

@ -31,8 +31,9 @@ func TestPostIndexRequestUnmarshalJSON(t *testing.T) {
expected postIndexRequest
err string
}{
{json: `{"options": {}}`, expected: postIndexRequest{Options: pilosa.IndexOptions{}}},
{json: `{"options": {"keys": true}}`, expected: postIndexRequest{Options: pilosa.IndexOptions{Keys: true}}},
{json: `{"options": {}}`, expected: postIndexRequest{Options: pilosa.IndexOptions{TrackExistence: true}}},
{json: `{"options": {"trackExistence": false}}`, expected: postIndexRequest{Options: pilosa.IndexOptions{TrackExistence: false}}},
{json: `{"options": {"keys": true}}`, expected: postIndexRequest{Options: pilosa.IndexOptions{Keys: true, TrackExistence: true}}},
{json: `{"options": 4}`, err: "options is not map[string]interface{}"},
{json: `{"option": {}}`, err: "Unknown key: option:map[]"},
{json: `{"options": {"badKey": "test"}}`, err: "Unknown key: badKey:test"},
@ -53,7 +54,7 @@ func TestPostIndexRequestUnmarshalJSON(t *testing.T) {
if test.err == "" {
if !reflect.DeepEqual(*actual, test.expected) {
t.Errorf("expected: %v, but got: %v", test.expected, *actual)
t.Errorf("expected: %v, but got: %v for JSON: %s", test.expected, *actual, test.json)
}
}

View file

@ -69,9 +69,10 @@ func NewIndex(path, name string) (*Index, error) {
newAttrStore: newNopAttrStore,
columnAttrs: nopStore,
broadcaster: NopBroadcaster,
Stats: NopStatsClient,
logger: NopLogger,
broadcaster: NopBroadcaster,
Stats: NopStatsClient,
logger: NopLogger,
trackExistence: true,
}, nil
}

View file

@ -117,7 +117,7 @@ var nameRegexp = regexp.MustCompile(`^[a-z][a-z0-9_-]{0,63}$`)
// ColumnAttrSet represents a set of attributes for a vertical column in an index.
// Can have a set of attributes attached to it.
type ColumnAttrSet struct {
ID uint64 `json:"id"`
ID uint64 `json:"id,omitempty"`
Key string `json:"key,omitempty"`
Attrs map[string]interface{} `json:"attrs,omitempty"`
}

View file

@ -11,6 +11,7 @@ Call <- 'Set' {p.startCall("Set")} open col comma args (comma timestamp)? close
/ 'SetColumnAttrs' {p.startCall("SetColumnAttrs")} open col comma args close {p.endCall()}
/ 'Clear' {p.startCall("Clear")} open col comma args close {p.endCall()}
/ 'ClearRow' {p.startCall("ClearRow")} open arg close {p.endCall()}
/ 'Store' {p.startCall("Store")} open Call comma arg close {p.endCall()}
/ 'TopN' {p.startCall("TopN")} open posfield (comma allargs)? close {p.endCall()}
/ 'Range' {p.startCall("Range")} open (timerange / conditional / arg) close {p.endCall()}
/ < IDENT > { p.startCall(buffer[begin:end] ) } open allargs comma? close { p.endCall() }

File diff suppressed because it is too large Load diff

View file

@ -676,6 +676,7 @@ func (s *Server) monitorDiagnostics() {
s.diagnostics.Set("NumCPU", runtime.NumCPU())
s.diagnostics.Set("NodeID", s.nodeID)
s.diagnostics.Set("ClusterID", s.cluster.id)
s.diagnostics.EnrichWithCPUInfo()
s.diagnostics.EnrichWithOSInfo()
// Flush the diagnostics metrics at startup, then on each tick interval

View file

@ -104,6 +104,22 @@ func TestHandler_Endpoints(t *testing.T) {
})
t.Run("ImportRoaringFieldTypeFail", func(t *testing.T) {
// Roaring import into a non-set field should fail.
if _, err := i0.CreateFieldIfNotExists("int-field", pilosa.OptFieldTypeInt(0, 1)); err != nil {
t.Fatal(err)
}
w := httptest.NewRecorder()
roaringData, _ := hex.DecodeString("3B3001000100000900010000000100010009000100")
req := test.MustNewHTTPRequest("POST", "/index/i0/field/int-field/import-roaring/0", bytes.NewBuffer(roaringData))
req.Header.Set("Content-Type", "application/x-binary")
h.ServeHTTP(w, req)
if w.Code != gohttp.StatusBadRequest {
t.Fatalf("unexpected status code: %d", w.Code)
}
})
t.Run("Status", func(t *testing.T) {
w := httptest.NewRecorder()
h.ServeHTTP(w, test.MustNewHTTPRequest("GET", "/status", nil))