Merge branch 'master' into bounds-check

This commit is contained in:
tgruben 2018-09-24 12:59:10 -05:00 • committed by GitHub
commit e62be47667
No known key found for this signature in database
GPG key ID: 4AEE18F83AFDEB23
27 changed files with 643 additions and 63 deletions

View file

@ -68,6 +68,7 @@ jobs:
docker:
- image: circleci/python:2.7-jessie
steps:
- run: '[[ -v CIRCLE_PR_NUMBER ]] && circleci step halt || true' # Skip job if this is a PR
- *fast-checkout
- run: sudo pip install awscli
- run: make prerelease-upload

View file

@ -42,7 +42,7 @@ func TestServerConfig(t *testing.T) {
tests := []commandTest{
// TEST 0
{
args: []string{"server", "--data-dir", actualDataDir, "--cluster.hosts", "localhost:10111,localhost:10110", "--bind", "localhost:10111"},
args: []string{"server", "--data-dir", actualDataDir, "--cluster.hosts", "localhost:10111,localhost:10110", "--bind", "localhost:10111", "--translation.map-size", "100000"},
env: map[string]string{"PILOSA_DATA_DIR": "/tmp/myEnvDatadir", "PILOSA_CLUSTER_LONG_QUERY_TIME": "1m30s", "PILOSA_MAX_WRITES_PER_REQUEST": "2000"},
cfgFileContent: `
data-dir = "/tmp/myFileDatadir"
@ -65,13 +65,14 @@ func TestServerConfig(t *testing.T) {
v.Check(cmd.Server.Config.Cluster.Hosts, []string{"localhost:10111", "localhost:10110"})
v.Check(cmd.Server.Config.Cluster.LongQueryTime, toml.Duration(time.Second*90))
v.Check(cmd.Server.Config.MaxWritesPerRequest, 2000)
v.Check(cmd.Server.Config.Translation.MapSize, 100000)
return v.Error()
},
},
// TEST 1
{
args: []string{"server", "--anti-entropy.interval", "9m0s"},
env: map[string]string{"PILOSA_CLUSTER_HOSTS": "localhost:1110,localhost:1111", "PILOSA_BIND": "localhost:1110"},
env: map[string]string{"PILOSA_CLUSTER_HOSTS": "localhost:1110,localhost:1111", "PILOSA_BIND": "localhost:1110", "PILOSA_TRANSLATION_MAP_SIZE": "100000"},
cfgFileContent: `
bind = "localhost:0"
data-dir = "` + actualDataDir + `"
@ -85,12 +86,13 @@ func TestServerConfig(t *testing.T) {
v := validator{}
v.Check(cmd.Server.Config.Cluster.Hosts, []string{"localhost:1110", "localhost:1111"})
v.Check(cmd.Server.Config.AntiEntropy.Interval, toml.Duration(time.Minute*9))
v.Check(cmd.Server.Config.Translation.MapSize, 100000)
return v.Error()
},
},
// TEST 2
{
args: []string{"server", "--log-path", logFile.Name(), "--cluster.disabled", "true"},
args: []string{"server", "--log-path", logFile.Name(), "--cluster.disabled", "true", "--translation.map-size", "100000"},
env: map[string]string{},
cfgFileContent: `
bind = "localhost:19444"

View file

@ -279,3 +279,49 @@ func TestImportCommand_BugOverwriteValue(t *testing.T) {
t.Fatalf("Import Run with values doesn't work: %s", err)
}
}
// Ensure that import into bool field runs.
func TestImportCommand_RunBool(t *testing.T) {
buf := bytes.Buffer{}
stdin, stdout, stderr := GetIO(buf)
cm := NewImportCommand(stdin, stdout, stderr)
ctx := context.Background()
cmd := test.MustRunCluster(t, 1)[0]
cm.Host = cmd.API.Node().URI.HostPort()
http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i", strings.NewReader("")))
http.DefaultClient.Do(MustNewHTTPRequest("POST", "http://"+cm.Host+"/index/i/field/f", strings.NewReader(`{"options":{"type": "bool"}}`)))
cm.Index = "i"
cm.Field = "f"
t.Run("Valid", func(t *testing.T) {
file, err := ioutil.TempFile("", "import-bool.csv")
if err != nil {
t.Fatal(err)
}
file.Write([]byte("0,1\n1,2\n1,3"))
cm.Paths = []string{file.Name()}
err = cm.Run(ctx)
if err != nil {
t.Fatalf("Import Run to bool field doesn't work: %s", err)
}
})
// Ensure that invalid bool values return an error.
t.Run("Invalid", func(t *testing.T) {
file, err := ioutil.TempFile("", "import-invalid-bool.csv")
if err != nil {
t.Fatal(err)
}
file.Write([]byte("0,1\n1,2\n1,3\n2,4"))
cm.Paths = []string{file.Name()}
err = cm.Run(ctx)
if !strings.Contains(err.Error(), "bool field imports only support values 0 and 1") {
t.Fatalf("expect error: bool field imports only support values 0 and 1, actual: %s", err)
}
})
}

View file

@ -45,6 +45,7 @@ func BuildServerFlags(cmd *cobra.Command, srv *server.Command) {
// Translation
flags.StringVarP(&srv.Config.Translation.PrimaryURL, "translation.primary-url", "", srv.Config.Translation.PrimaryURL, "DEPRECATED: URL for primary translation node for replication.")
flags.IntVarP(&srv.Config.Translation.MapSize, "translation.map-size", "", srv.Config.Translation.MapSize, "Size in bytes of mmap to allocate for key translation.")
// Gossip
flags.StringVarP(&srv.Config.Gossip.Port, "gossip.port", "", srv.Config.Gossip.Port, "Port to which pilosa should bind for internal state sharing.")

View file

@ -64,6 +64,20 @@ If you are using [integer](../data-model/#bsi-range-encoding) field values, the
pilosa import -i project -f stargazer-counts project-stargazer-counts.csv
```
##### Importing Boolean Values
If you are using a [boolean](../data-model/#boolean) field, the CSV file should be in the format `Boolean,Value`, where `Boolean` is either `0` (false) or `1` (true).
For example, importing a file with the following contents will result in columns 3 and 9 being set in the `false` row, and columns 1, 2, 4, and 8 being set in the `true` row.
```
0,3
0,9
1,1
1,2
1,4
1,8
```
<div class="note">
<p>Note that you must first create a field. View <a href="../api-reference/#create-field">Create Field</a> for more details. The `-e` flag can create the necessary schema when using a field of type "set".</p>
</div>

View file

@ -108,6 +108,8 @@ The request payload is in JSON, and may contain the `options` field. The `option
* `int`
* `min` (int): Minimum integer value allowed for the field.
* `max` (int): Maximum integer value allowed for the field.
* `bool`
* (boolean fields take no arguments)
* `time`
* `timeQuantum` (string): [Time Quantum](../data-model/#time-quantum) for this field.
* `mutex`

View file

@ -307,6 +307,18 @@ The config file is in the [toml format](https://github.com/toml-lang/toml) and h
skip-verify = true
```
#### Translation Map Size
* Description: Size in bytes of mmap to allocate for key translation
* Flag: `translation.map-size`
* Env: `PILOSA_TRANSLATION_MAP_SIZE`
* Config:
```toml
[translation]
map-size = 10737418240
```
### Example Cluster Configuration
A three node cluster running on different hosts could be minimally configured as follows:

View file

@ -117,7 +117,7 @@ Query operations run in parallel, and they are evenly distributed across a clust
### Field Type
Upon creation, fields are configured to be of a certain type. Pilosa supports the following field types: `set`, `int`, `time`, and `mutex`.
Upon creation, fields are configured to be of a certain type. Pilosa supports the following field types: `set`, `int`, `bool`, `time`, and `mutex`.
#### Set
@ -192,3 +192,7 @@ Set(3, A=8, 2017-05-19T00:00)
#### Mutex
Mutex fields are similar to `set` fields, with the distinction of requiring the row value for each column to be mutually exclusive. In other words, each column can only have a single value for the field. If the field value for a column is updated on a `mutex` field, then the previous field value for that column will be cleared. This field type is like a field in an RDBMS table where every record contains a single value for a particular field.
#### Boolean
A boolean field is similar to a `mutex` field tracking only two values: `true` and `false`. Boolean fields do not maintain a sorted cache, nor do they support key values.

View file

@ -22,7 +22,7 @@ nav = []
<strong id="fragment">Fragment:</strong> A Fragment is the intersection of a [field](#field) and a [shard](#shard) in an [index](#index).
<strong id="field">[Field](../data-model/#field):</strong> Fields are used to group [rows](#row) into different categories. Row IDs are namespaced by field such that the same row ID in a different field refers to a different row. For [ranked](#topn) fields, rows are kept in sorted order within the field. Fields are one of four types: set, [int](#bsi), time, and mutex. For more information, see [data model](../data-model/) and [Creating fields](../api-reference/#create-field).
<strong id="field">[Field](../data-model/#field):</strong> Fields are used to group [rows](#row) into different categories. Row IDs are namespaced by field such that the same row ID in a different field refers to a different row. For [ranked](#topn) fields, rows are kept in sorted order within the field. Fields are one of four types: set, [int](#bsi), bool, time, and mutex. For more information, see [data model](../data-model/) and [Creating fields](../api-reference/#create-field).
<strong id="frame">[Frame](../data-model/#field):</strong> Prior to Pilosa 1.0, fields were known as frames.

View file

@ -483,7 +483,7 @@ curl localhost:10101/index/patients/field/tcells \
{"success":true}
```
Next, let's populate our fields with data. There are two ways to get data into fields: use the `SetFieldValue()` PQL function to set fields individually, or use the `pilosa import` command to import many values at once. First, let's set some field data using PQL.
Next, let's populate our fields with data. There are two ways to get data into fields: use the `Set()` PQL function to set fields individually, or use the `pilosa import` command to import many values at once. First, let's set some field data using PQL.
The following queries set the age, weight, and t-cell count for the patient with ID `1` in our system:
``` request

View file

@ -734,8 +734,7 @@ func (e *executor) executeBitmapShard(_ context.Context, index string, c *pql.Ca
rowID, rowOK, rowErr := c.UintArg(fieldName)
if rowErr != nil {
return nil, fmt.Errorf("Row() error with arg for row: %v", rowErr)
}
if !rowOK {
} else if !rowOK {
return nil, fmt.Errorf("Row() must specify %v", rowLabel)
}
@ -1140,7 +1139,6 @@ func (e *executor) executeClearBitField(ctx context.Context, index string, c *pq
// executeSet executes a Set() call.
func (e *executor) executeSet(ctx context.Context, index string, c *pql.Call, opt *execOptions) (bool, error) {
// Read colID.
colID, ok, err := c.UintArg("_" + columnLabel)
if err != nil {
@ -1172,8 +1170,9 @@ func (e *executor) executeSet(ctx context.Context, index string, c *pql.Call, op
}
}
// Int field.
if f.Type() == FieldTypeInt {
// Read remaining fields using labels.
// Read row value.
rowVal, ok, err := c.IntArg(fieldName)
if err != nil {
return false, fmt.Errorf("reading Set() row: %v", err)
@ -1182,27 +1181,27 @@ func (e *executor) executeSet(ctx context.Context, index string, c *pql.Call, op
}
return e.executeSetValueField(ctx, index, c, f, colID, rowVal, opt)
} else {
// Read remaining fields using labels.
rowID, ok, err := c.UintArg(fieldName)
if err != nil {
return false, fmt.Errorf("reading Set() row: %v", err)
} else if !ok {
return false, fmt.Errorf("Set() row argument '%v' required", rowLabel)
}
var timestamp *time.Time
sTimestamp, ok := c.Args["_timestamp"].(string)
if ok {
t, err := time.Parse(TimeFormat, sTimestamp)
if err != nil {
return false, fmt.Errorf("invalid date: %s", sTimestamp)
}
timestamp = &t
}
return e.executeSetBitField(ctx, index, c, f, colID, rowID, timestamp, opt)
}
// Read row ID.
rowID, ok, err := c.UintArg(fieldName)
if err != nil {
return false, fmt.Errorf("reading Set() row: %v", err)
} else if !ok {
return false, fmt.Errorf("Set() row argument '%v' required", rowLabel)
}
var timestamp *time.Time
sTimestamp, ok := c.Args["_timestamp"].(string)
if ok {
t, err := time.Parse(TimeFormat, sTimestamp)
if err != nil {
return false, fmt.Errorf("invalid date: %s", sTimestamp)
}
timestamp = &t
}
return e.executeSetBitField(ctx, index, c, f, colID, rowID, timestamp, opt)
}
// executeSetBitField executes a Set() call for a specific field.
@ -1675,7 +1674,21 @@ func (e *executor) translateCall(index string, idx *Index, c *pql.Call) error {
// will raise an error downstream when it's used.
return nil
}
if field.keys() {
// Bool field keys do not use the translator because there
// are only two possible values. Instead, they are handled
// directly.
if field.Type() == FieldTypeBool {
boolVal, err := callArgBool(c, rowKey)
if err != nil {
return errors.Wrap(err, "getting bool key")
}
rowID := falseRowID
if boolVal {
rowID = trueRowID
}
c.Args[rowKey] = rowID
} else if field.keys() {
if c.Args[rowKey] != nil && !isString(c.Args[rowKey]) {
return errors.New("row value must be a string when field 'keys' option enabled")
}
@ -1832,6 +1845,18 @@ func (vc *ValCount) larger(other ValCount) ValCount {
}
}
func callArgBool(call *pql.Call, key string) (bool, error) {
value, ok := call.Args[key]
if !ok {
return false, errors.New("missing bool argument")
}
b, ok := value.(bool)
if !ok {
return false, fmt.Errorf("invalid bool argument type: %T", value)
}
return b, nil
}
func callArgString(call *pql.Call, key string) string {
value, ok := call.Args[key]
if !ok {

View file

@ -376,6 +376,78 @@ func TestExecutor_Execute_SetBit(t *testing.T) {
})
}
// Ensure a set query can be executed on a bool field.
func TestExecutor_Execute_SetBool(t *testing.T) {
t.Run("Basic", func(t *testing.T) {
c := test.MustRunCluster(t, 1)
defer c.Close()
hldr := test.Holder{Holder: c[0].Server.Holder()}
// Create fields.
index := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{})
if _, err := index.CreateFieldIfNotExists("f", pilosa.OptFieldTypeBool()); err != nil {
t.Fatal(err)
}
// Set a true bit.
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(100, f=true)`}); err != nil {
t.Fatal(err)
} else if !res.Results[0].(bool) {
t.Fatalf("expected column changed")
}
// Set the same bit to true again verify nothing changed.
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(100, f=true)`}); err != nil {
t.Fatal(err)
} else if res.Results[0].(bool) {
t.Fatalf("expected column to be unchanged")
}
// Set the same bit to false.
if res, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(100, f=false)`}); err != nil {
t.Fatal(err)
} else if !res.Results[0].(bool) {
t.Fatalf("expected column changed")
}
// Ensure that the false row is set.
if result, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Row(f=false)`}); err != nil {
t.Fatal(err)
} else if columns := result.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(columns, []uint64{100}) {
t.Fatalf("unexpected colums: %+v", columns)
}
// Ensure that the true row is empty.
if result, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Row(f=true)`}); err != nil {
t.Fatal(err)
} else if columns := result.Results[0].(*pilosa.Row).Columns(); !reflect.DeepEqual(columns, []uint64{}) {
t.Fatalf("unexpected colums: %+v", columns)
}
})
t.Run("Error", func(t *testing.T) {
c := test.MustRunCluster(t, 1)
defer c.Close()
hldr := test.Holder{Holder: c[0].Server.Holder()}
// Create fields.
index := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{})
if _, err := index.CreateFieldIfNotExists("f", pilosa.OptFieldTypeBool()); err != nil {
t.Fatal(err)
}
// Set bool using a string value.
if _, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(100, f="true")`}); err == nil {
t.Fatalf("expected invalid bool type error")
}
// Set bool using an integer.
if _, err := c[0].API.Query(context.Background(), &pilosa.QueryRequest{Index: "i", Query: `Set(100, f=1)`}); err == nil {
t.Fatalf("expected invalid bool type error")
}
})
}
// Ensure old PQL syntax doesn't break anything too badly.
func TestExecutor_Execute_OldPQL(t *testing.T) {
c := test.MustRunCluster(t, 1)
@ -397,7 +469,7 @@ func TestExecutor_Execute_SetValue(t *testing.T) {
defer c.Close()
hldr := test.Holder{Holder: c[0].Server.Holder()}
// Create felds.
// Create fields.
index := hldr.MustCreateIndexIfNotExists("i", pilosa.IndexOptions{})
if _, err := index.CreateFieldIfNotExists("f", pilosa.OptFieldTypeInt(0, 50)); err != nil {
t.Fatal(err)

View file

@ -52,6 +52,7 @@ const (
FieldTypeInt = "int"
FieldTypeTime = "time"
FieldTypeMutex = "mutex"
FieldTypeBool = "bool"
)
// Field represents a container for views.
@ -155,6 +156,16 @@ func OptFieldTypeMutex(cacheType string, cacheSize uint32) FieldOption {
}
}
func OptFieldTypeBool() FieldOption {
return func(fo *FieldOptions) error {
if fo.Type != "" {
return errors.Errorf("field type is already set to: %s", fo.Type)
}
fo.Type = FieldTypeBool
return nil
}
}
// NewField returns a new instance of field.
func NewField(path, index, name string, opts FieldOption) (*Field, error) {
err := validateName(name)
@ -441,6 +452,14 @@ func (f *Field) applyOptions(opt FieldOptions) error {
f.options.Max = 0
f.options.TimeQuantum = ""
f.options.Keys = opt.Keys
case FieldTypeBool:
f.options.Type = FieldTypeBool
f.options.CacheType = CacheTypeNone
f.options.CacheSize = 0
f.options.Min = 0
f.options.Max = 0
f.options.TimeQuantum = ""
f.options.Keys = false
default:
return errors.New("invalid field type")
}
@ -681,6 +700,10 @@ func (f *Field) deleteView(name string) error {
}
// Row returns a row of the standard view.
// It seems this method is only being used by the test
// package, and the fact that it's only allowed on
// `set` fields is odd. This may be considered for
// deprecation in a future version.
func (f *Field) Row(rowID uint64) (*Row, error) {
if f.Type() != FieldTypeSet {
return nil, errors.Errorf("row method unsupported for field type: %s", f.Type())
@ -954,10 +977,18 @@ func (f *Field) Import(rowIDs, columnIDs []uint64, timestamps []*time.Time) erro
return errors.New("time quantum not set in field")
}
fieldType := f.Type()
// Split import data by fragment.
dataByFragment := make(map[importKey]importData)
for i := range rowIDs {
rowID, columnID := rowIDs[i], columnIDs[i]
// Bool-specific data validation.
if fieldType == FieldTypeBool && rowID > 1 {
return errors.New("bool field imports only support values 0 and 1")
}
var timestamp *time.Time
if len(timestamps) > i {
timestamp = timestamps[i]
@ -1191,6 +1222,12 @@ func (o *FieldOptions) MarshalJSON() ([]byte, error) {
o.CacheSize,
o.Keys,
})
case FieldTypeBool:
return json.Marshal(struct {
Type string `json:"type"`
}{
o.Type,
})
}
return nil, errors.New("invalid field type")
}

View file

@ -72,6 +72,10 @@ const (
// defaultFragmentMaxOpN is the default value for Fragment.MaxOpN.
defaultFragmentMaxOpN = 2000
// Row ids used for boolean fields.
falseRowID = uint64(0)
trueRowID = uint64(1)
)
// fragment represents the intersection of a field and shard in an index.
@ -1330,16 +1334,29 @@ func (f *fragment) bulkImport(rowIDs, columnIDs []uint64) error {
return fmt.Errorf("mismatch of row/column len: %d != %d", len(rowIDs), len(columnIDs))
}
if f.mutexVector != nil {
return f.bulkImportMutex(rowIDs, columnIDs)
}
return f.bulkImportStandard(rowIDs, columnIDs)
}
// bulkImportStandard performs a bulk import on a standard fragment.
func (f *fragment) bulkImportStandard(rowIDs, columnIDs []uint64) error {
// Create a temporary bitmap which will be populated by rowIDs and columnIDs
// and then merged into the existing fragment's bitmap.
localBitmap := roaring.NewBitmap()
// Disconnect op writer so we don't append updates.
localBitmap.OpWriter = nil
// Process every bit.
// If an error occurs then reopen the storage.
lastID := uint64(0)
// rowSet maintains the set of rowIDs present in this import.
// It allows the cache to be updated once per row, instead of once
// per bit.
rowSet := make(map[uint64]struct{})
lastRowID := uint64(0)
// Process every bit by writing to a local bitmap,
// to be merged with fragment storage next.
for i := range rowIDs {
rowID, columnID := rowIDs[i], columnIDs[i]
@ -1349,18 +1366,18 @@ func (f *fragment) bulkImport(rowIDs, columnIDs []uint64) error {
return err
}
// Write to storage.
// Write to local storage.
_, err = localBitmap.Add(pos)
if err != nil {
return err
}
// Reduce the StatsD rate for high volume stats
f.stats.Count("ImportBit", 1, 0.0001)
// import optimization to avoid linear foreach calls
// slight risk of concurrent cache counter being off but
// no real danger
if i == 0 || rowID != lastID {
lastID = rowID
// Add row to rowSet.
if i == 0 || rowID != lastRowID {
lastRowID = rowID
rowSet[rowID] = struct{}{}
}
@ -1389,6 +1406,96 @@ func (f *fragment) bulkImport(rowIDs, columnIDs []uint64) error {
return unprotectedWriteToFragment(f, results)
}
// bulkImportMutex performs a bulk import on a fragment while ensuring
// mutex restrictions. Because the mutex requirements must be checked
// against storage, this method must acquire a write lock on the fragment
// during the entire process, and it handles every bit independently.
func (f *fragment) bulkImportMutex(rowIDs, columnIDs []uint64) error {
f.mu.Lock()
defer f.mu.Unlock()
// Disconnect op writer so we don't append updates.
f.storage.OpWriter = nil
// If an error occurs then reopen the storage.
if err := func() error {
// rowSet maintains the set of rowIDs present in this import.
// It allows the cache to be updated once per row, instead of once
// per bit.
rowSet := make(map[uint64]struct{})
lastRowID := uint64(0)
// Process every bit.
for i := range rowIDs {
rowID, columnID := rowIDs[i], columnIDs[i]
// Handle mutex vector (i.e. clear an existing row).
if existingRowID, found := f.mutexVector.Get(columnID); found && existingRowID != rowID {
// Determine the position of the bit in the storage.
pos, err := f.pos(existingRowID, columnID)
if err != nil {
return err
}
// Clear storage.
_, err = f.storage.Remove(pos)
if err != nil {
return err
}
rowSet[existingRowID] = struct{}{}
}
// Determine the position of the bit in the storage.
pos, err := f.pos(rowID, columnID)
if err != nil {
return err
}
// Write to storage.
_, err = f.storage.Add(pos)
if err != nil {
return err
}
// Reduce the StatsD rate for high volume stats
f.stats.Count("ImportBit", 1, 0.0001)
// Add row to rowSet.
if i == 0 || rowID != lastRowID {
lastRowID = rowID
rowSet[rowID] = struct{}{}
}
// Invalidate block checksum.
delete(f.checksums, int(rowID/HashBlockSize))
}
// Update cache counts for all rows.
for rowID := range rowSet {
// Import should ALWAYS have row() load a new bm from fragment.storage
// because the row that's in rowCache hasn't been updated with
// this import's data.
f.cache.BulkAdd(rowID, f.unprotectedRow(rowID).Count())
}
f.cache.Invalidate()
return nil
}(); err != nil {
_ = f.closeStorage()
_ = f.openStorage()
return err
}
// Write the storage to disk and reload.
if err := f.snapshot(); err != nil {
return err
}
return nil
}
// importValue bulk imports a set of range-encoded values.
func (f *fragment) importValue(columnIDs, values []uint64, bitDepth uint) error {
f.mu.Lock()
@ -2103,3 +2210,33 @@ func (v *rowsVector) Get(colID uint64) (uint64, bool) {
// Set is not used for rowsVector.
func (v *rowsVector) Set(colID, rowID uint64) {}
// boolVector implements the vector interface by looking
// at data in rows 0 and 1.
type boolVector struct {
f *fragment
}
// newBoolVector returns a boolVector for a given fragment.
func newBoolVector(f *fragment) *boolVector {
return &boolVector{
f: f,
}
}
// Get returns the rowID associated to the given colID.
// Additionally, it returns true if a value was found,
// otherwise it returns false.
func (v *boolVector) Get(colID uint64) (uint64, bool) {
rows := v.f.rowsForColumn(colID)
if len(rows) == 1 {
switch rows[0] {
case falseRowID, trueRowID:
return rows[0], true
}
}
return 0, false
}
// Set is not used for boolVector.
func (v *boolVector) Set(colID, rowID uint64) {}

View file

@ -1178,6 +1178,146 @@ func TestFragment_SetMutex(t *testing.T) {
}
}
// Ensure a fragment can import mutually exclusive values.
func TestFragment_ImportMutex(t *testing.T) {
tests := []struct {
rowIDs []uint64
colIDs []uint64
exp map[uint64][]uint64
}{
{
[]uint64{1, 1, 1, 1},
[]uint64{0, 1, 2, 3},
map[uint64][]uint64{
1: {0, 1, 2, 3},
},
},
{
[]uint64{1, 1, 1, 1, 2, 2, 2, 2},
[]uint64{0, 1, 2, 3, 0, 1, 2, 3},
map[uint64][]uint64{
1: {},
2: {0, 1, 2, 3},
},
},
{
[]uint64{1, 1, 1, 1, 2},
[]uint64{0, 1, 2, 3, 1},
map[uint64][]uint64{
1: {0, 2, 3},
2: {1},
},
},
{
[]uint64{1, 1, 1, 1, 2, 2, 1},
[]uint64{0, 1, 2, 3, 1, 8, 1},
map[uint64][]uint64{
1: {0, 1, 2, 3},
2: {8},
},
},
{
[]uint64{1, 2, 3},
[]uint64{8, 8, 8},
map[uint64][]uint64{
1: {},
2: {},
3: {8},
},
},
}
for i, test := range tests {
t.Run(fmt.Sprintf("importmutex%d", i), func(t *testing.T) {
f := mustOpenMutexFragment("i", "f", viewStandard, 0, "")
defer f.Close()
err := f.bulkImport(test.rowIDs, test.colIDs)
if err != nil {
t.Fatalf("bulk importing ids: %v", err)
}
// Check for expected results.
for k, v := range test.exp {
cols := f.row(k).Columns()
if !reflect.DeepEqual(cols, v) {
t.Fatalf("expected: %v, but got: %v", v, cols)
}
}
})
}
}
// Ensure a fragment can import bool values.
func TestFragment_ImportBool(t *testing.T) {
tests := []struct {
rowIDs []uint64
colIDs []uint64
exp map[uint64][]uint64
}{
{
[]uint64{1, 1, 1, 1},
[]uint64{0, 1, 2, 3},
map[uint64][]uint64{
1: {0, 1, 2, 3},
},
},
{
[]uint64{0, 0, 0, 0, 1, 1, 1, 1},
[]uint64{0, 1, 2, 3, 0, 1, 2, 3},
map[uint64][]uint64{
0: {},
1: {0, 1, 2, 3},
},
},
{
[]uint64{0, 0, 0, 0, 1},
[]uint64{0, 1, 2, 3, 1},
map[uint64][]uint64{
0: {0, 2, 3},
1: {1},
},
},
{
[]uint64{1, 1, 1, 1, 0, 0, 1},
[]uint64{0, 1, 2, 3, 1, 8, 1},
map[uint64][]uint64{
0: {8},
1: {0, 1, 2, 3},
},
},
{
[]uint64{0, 1, 2},
[]uint64{8, 8, 8},
map[uint64][]uint64{
0: {},
1: {}, // This isn't {8} because fragment doesn't validate bool values.
2: {8},
},
},
}
for i, test := range tests {
t.Run(fmt.Sprintf("importmutex%d", i), func(t *testing.T) {
f := mustOpenBoolFragment("i", "f", viewStandard, 0, "")
defer f.Close()
err := f.bulkImport(test.rowIDs, test.colIDs)
if err != nil {
t.Fatalf("bulk importing ids: %v", err)
}
// Check for expected results.
for k, v := range test.exp {
cols := f.row(k).Columns()
if !reflect.DeepEqual(cols, v) {
t.Fatalf("expected: %v, but got: %v", v, cols)
}
}
})
}
}
func BenchmarkFragment_Snapshot(b *testing.B) {
if *FragmentPath == "" {
b.Skip("no fragment specified")
@ -1302,6 +1442,13 @@ func mustOpenMutexFragment(index, field, view string, shard uint64, cacheType st
return frag
}
// mustOpenBoolFragment returns a new instance of Fragment for a bool field.
func mustOpenBoolFragment(index, field, view string, shard uint64, cacheType string) *fragment {
frag := mustOpenFragment(index, field, view, shard, cacheType)
frag.mutexVector = newBoolVector(frag)
return frag
}
// Reopen closes the fragment and reopens it as a new instance.
func (f *fragment) reopen() error {
if err := f.Close(); err != nil {

View file

@ -687,6 +687,8 @@ func (h *Handler) handlePostField(w http.ResponseWriter, r *http.Request) {
fos = append(fos, pilosa.OptFieldTypeTime(*req.Options.TimeQuantum))
case pilosa.FieldTypeMutex:
fos = append(fos, pilosa.OptFieldTypeMutex(*req.Options.CacheType, *req.Options.CacheSize))
case pilosa.FieldTypeBool:
fos = append(fos, pilosa.OptFieldTypeBool())
}
if req.Options.Keys != nil {
if *req.Options.Keys {
@ -778,6 +780,20 @@ func (o *fieldOptions) validate() error {
} else if o.TimeQuantum != nil {
return pilosa.NewBadRequestError(errors.New("timeQuantum does not apply to field type mutex"))
}
case pilosa.FieldTypeBool:
if o.CacheType != nil {
return pilosa.NewBadRequestError(errors.New("cacheType does not apply to field type bool"))
} else if o.CacheSize != nil {
return pilosa.NewBadRequestError(errors.New("cacheSize does not apply to field type bool"))
} else if o.Min != nil {
return pilosa.NewBadRequestError(errors.New("min does not apply to field type bool"))
} else if o.Max != nil {
return pilosa.NewBadRequestError(errors.New("max does not apply to field type bool"))
} else if o.TimeQuantum != nil {
return pilosa.NewBadRequestError(errors.New("timeQuantum does not apply to field type bool"))
} else if o.Keys != nil {
return pilosa.NewBadRequestError(errors.New("keys does not apply to field type bool"))
}
default:
return errors.Errorf("invalid field type: %s", o.Type)
}

View file

@ -248,6 +248,23 @@ func (c *Call) FieldArg() (string, error) {
return "", fmt.Errorf("No field argument specified")
}
// BoolArg is for reading the value at key from call.Args as a bool. If the
// key is not in Call.Args, the value of the returned bool will be false, and
// the error will be nil. The value is assumed to be a bool. An error is
// returned if the value is not a bool.
func (c *Call) BoolArg(key string) (bool, bool, error) {
val, ok := c.Args[key]
if !ok {
return false, false, nil
}
switch tval := val.(type) {
case bool:
return tval, true, nil
default:
return false, true, fmt.Errorf("could not convert %v of type %T to bool in Call.BoolArg", tval, tval)
}
}
// UintArg is for reading the value at key from call.Args as a uint64. If the
// key is not in Call.Args, the value of the returned bool will be false, and
// the error will be nil. The value is assumed to be a uint64 or an int64 and

View file

@ -234,6 +234,13 @@ func OptServerClusterHasher(h Hasher) ServerOption {
}
}
func OptServerTranslateFileMapSize(mapSize int) ServerOption {
return func(s *Server) error {
s.holder.translateFile = NewTranslateFile(OptTranslateFileMapSize(mapSize))
return nil
}
}
// NewServer returns a new instance of Server.
func NewServer(opts ...ServerOption) (*Server, error) {
s := &Server{

View file

@ -71,8 +71,9 @@ type Config struct {
// Gossip config is based around memberlist.Config.
Gossip gossip.Config `toml:"gossip"`
// DEPRECATED: Translation config supports translation store replication.
Translation struct {
MapSize int `toml:"map-size"`
// DEPRECATED: Translation config supports translation store replication.
PrimaryURL string `toml:"primary-url"`
} `toml:"translation"`

View file

@ -283,6 +283,14 @@ func (m *Command) SetupServer() error {
coordinatorOpt,
}
if m.Config.Translation.MapSize > 0 {
serverOptions = append(
serverOptions,
pilosa.OptServerTranslateFileMapSize(
m.Config.Translation.MapSize,
),
)
}
serverOptions = append(serverOptions, m.serverOptions...)
m.Server, err = pilosa.NewServer(serverOptions...)

View file

@ -404,6 +404,7 @@ func TestClusteringNodesReplica1(t *testing.T) {
// Create new main with the same config.
config := cluster[2].Command.Config
config.Translation.MapSize = 100000
// config.Bind = cluster[2].API.Node().URI.HostPort()
// this isn't necessary, but makes the test run way faster
@ -477,6 +478,7 @@ func TestClusteringNodesReplica2(t *testing.T) {
// Create new main with the same config.
config := cluster[2].Command.Config
config.Translation.MapSize = 100000
// config.Bind = cluster[2].API.Node().URI.HostPort()
// this isn't necessary, but makes the test run way faster
@ -497,6 +499,7 @@ func TestClusteringNodesReplica2(t *testing.T) {
// Create new main with the same config.
config = cluster[1].Command.Config
// config.Bind = cluster[1].API.Node().URI.HostPort()
config.Translation.MapSize = 100000
// this isn't necessary, but makes the test run way faster
config.Gossip.Port = strconv.Itoa(int(cluster[1].Command.GossipTransport().URI.Port))

View file

@ -55,11 +55,20 @@ func newCommand(opts ...server.CommandOption) *Command {
panic(err)
}
// set aggressive close timeout by default to avoid hanging tests. This was
// Set aggressive close timeout by default to avoid hanging tests. This was
// a problem with PDK tests which used go-pilosa as well. We put it at the
// beginning of the option slice so that it can be overridden by user-passed
// options.
opts = append([]server.CommandOption{server.OptCommandCloseTimeout(time.Millisecond * 2)}, opts...)
// Also set TranslateFile MapSize to a smaller number so memory allocation
// does not fail on 32-bit systems.
opts = append([]server.CommandOption{
server.OptCommandCloseTimeout(time.Millisecond * 2),
server.OptCommandServerOptions(
pilosa.OptServerTranslateFileMapSize(
2 << 25,
),
),
}, opts...)
m := &Command{commandOptions: opts}
m.Command = server.NewCommand(bytes.NewReader(nil), ioutil.Discard, ioutil.Discard, opts...)
m.Config.DataDir = path

View file

@ -5,7 +5,6 @@ import (
"bytes"
"context"
"encoding/binary"
"errors"
"fmt"
"io"
"io/ioutil"
@ -17,6 +16,7 @@ import (
"time"
"github.com/cespare/xxhash"
"github.com/pkg/errors"
)
const (
@ -81,9 +81,29 @@ type TranslateFile struct {
replicationRetryInterval time.Duration
}
// TranslateFileOption is a functional option type for pilosa.TranslateFile
type TranslateFileOption func(f *TranslateFile) error
func OptTranslateFileMapSize(mapSize int) TranslateFileOption {
return func(f *TranslateFile) error {
f.mapSize = mapSize
return nil
}
}
// NewTranslateFile returns a new instance of TranslateFile.
func NewTranslateFile() *TranslateFile {
return &TranslateFile{
func NewTranslateFile(opts ...TranslateFileOption) *TranslateFile {
var defaultMapSize64 int64 = 10 * (1 << 30)
var defaultMapSize int
if ^uint(0)>>32 > 0 {
// 10GB default map size
defaultMapSize = int(defaultMapSize64)
} else {
// Use 2GB default map size on 32-bit systems
defaultMapSize = (1 << 31) - 1
}
f := &TranslateFile{
writeNotify: make(chan struct{}),
closing: make(chan struct{}),
cols: make(map[string]*index),
@ -96,6 +116,16 @@ func NewTranslateFile() *TranslateFile {
replicationRetryInterval: defaultReplicationRetryInterval,
}
for _, opt := range opts {
err := opt(f)
if err != nil {
// TODO (2.0): Change func signature to return error
panic(errors.Wrap(err, "applying option"))
}
}
return f
}
func (s *TranslateFile) Open() (err error) {

View file

@ -1,5 +0,0 @@
package pilosa
// defaultMapSize is the default size of mapped memory for the translate store.
// It is passed as an int to syscall.Mmap and so must be < 2^31
const defaultMapSize = (1 << 31) - 1 // 2GB

View file

@ -1,8 +0,0 @@
// +build !386
package pilosa
// defaultMapSize is the default size of mapped memory for the translate store.
// It is passed as an int to syscall.Mmap and so can only be larger than 2^31 on
// 64bit systems.
const defaultMapSize = 10 * (1 << 30) // 10GB

View file

@ -367,7 +367,7 @@ func TestPrintTranslateFile(t *testing.T) {
}
f.Close()
s := pilosa.NewTranslateFile()
s := pilosa.NewTranslateFile(pilosa.OptTranslateFileMapSize(2 << 25))
s.Path = f.Name()
err = s.Open()
if err != nil {
@ -809,7 +809,7 @@ func NewTranslateFile() *TranslateFile {
}
f.Close()
s := &TranslateFile{TranslateFile: pilosa.NewTranslateFile()}
s := &TranslateFile{TranslateFile: pilosa.NewTranslateFile(pilosa.OptTranslateFileMapSize(2 << 25))}
s.Path = f.Name()
return s
}
@ -847,7 +847,7 @@ func (s *TranslateFile) Reopen() error {
}
s.lock.Lock()
s.TranslateFile = pilosa.NewTranslateFile()
s.TranslateFile = pilosa.NewTranslateFile(pilosa.OptTranslateFileMapSize(2 << 25))
s.lock.Unlock()
s.Path = prev.Path
s.SetPrimaryStore("restored-primary", prev.PrimaryTranslateStore)

View file

@ -244,6 +244,8 @@ func (v *view) newFragment(path string, shard uint64) *fragment {
frag.stats = v.stats.WithTags(fmt.Sprintf("shard:%d", shard))
if v.fieldType == FieldTypeMutex {
frag.mutexVector = newRowsVector(frag)
} else if v.fieldType == FieldTypeBool {
frag.mutexVector = newBoolVector(frag)
}
return frag
}