From 7a32745b28e3e5ddd48d65e7ab75950b84e0c987 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Wed, 19 Sep 2018 12:35:00 -0500 Subject: [PATCH 1/8] Make translate map size configurable. --- ctl/server.go | 1 + holder.go | 23 +++++++++++++++++-- server.go | 12 ++++++++++ server/config.go | 3 ++- server/server.go | 10 +++++++++ test/pilosa.go | 15 +++++++++++-- translate.go | 33 +++++++++++++++++++++++++--- translate_mapsize_386.go | 5 ----- translate_mapsize_all64bitsystems.go | 8 ------- 9 files changed, 89 insertions(+), 21 deletions(-) delete mode 100644 translate_mapsize_386.go delete mode 100644 translate_mapsize_all64bitsystems.go diff --git a/ctl/server.go b/ctl/server.go index 90fdfb46e..194633a48 100644 --- a/ctl/server.go +++ b/ctl/server.go @@ -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 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.") diff --git a/holder.go b/holder.go index 794104e53..5cffeea9b 100644 --- a/holder.go +++ b/holder.go @@ -77,9 +77,19 @@ type Holder struct { Logger Logger } +// HolderOption is a functional option type for pilosa.Holder +type HolderOption func(f *Holder) error + +func OptHolderTranslateFileMapSize(mapSize int) HolderOption { + return func(h *Holder) error { + h.translateFile = NewTranslateFile(OptTranslateFileMapSize(mapSize)) + return nil + } +} + // NewHolder returns a new instance of Holder. -func NewHolder() *Holder { - return &Holder{ +func NewHolder(opts ...HolderOption) *Holder { + h := &Holder{ indexes: make(map[string]*Index), closing: make(chan struct{}), @@ -97,6 +107,15 @@ func NewHolder() *Holder { Logger: NopLogger, } + + for _, opt := range opts { + err := opt(h) + if err != nil { + // TODO (2.0): Change func signature to return error + panic(errors.Wrap(err, "applying option")) + } + } + return h } // Open initializes the root data directory for the holder. diff --git a/server.go b/server.go index cde1dd406..708fd853b 100644 --- a/server.go +++ b/server.go @@ -234,6 +234,18 @@ func OptServerClusterHasher(h Hasher) ServerOption { } } +func OptServerHolderOptions(opts ...HolderOption) ServerOption { + return func(s *Server) error { + for _, opt := range opts { + err := opt(s.holder) + if err != nil { + return errors.Wrap(err, "applying option") + } + } + return nil + } +} + // NewServer returns a new instance of Server. func NewServer(opts ...ServerOption) (*Server, error) { s := &Server{ diff --git a/server/config.go b/server/config.go index 8367eb835..d255e4bb6 100644 --- a/server/config.go +++ b/server/config.go @@ -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"` diff --git a/server/server.go b/server/server.go index 4fdd711fe..61942522f 100644 --- a/server/server.go +++ b/server/server.go @@ -283,6 +283,16 @@ func (m *Command) SetupServer() error { coordinatorOpt, } + if m.Config.Translation.MapSize > 0 { + serverOptions = append( + serverOptions, + pilosa.OptServerHolderOptions( + pilosa.OptHolderTranslateFileMapSize( + m.Config.Translation.MapSize, + ), + ), + ) + } serverOptions = append(serverOptions, m.serverOptions...) m.Server, err = pilosa.NewServer(serverOptions...) diff --git a/test/pilosa.go b/test/pilosa.go index 9dbd78057..835291cbd 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -55,11 +55,22 @@ 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.OptServerHolderOptions( + pilosa.OptHolderTranslateFileMapSize( + 2 << 26, + ), + ), + ), + }, opts...) m := &Command{commandOptions: opts} m.Command = server.NewCommand(bytes.NewReader(nil), ioutil.Discard, ioutil.Discard, opts...) m.Config.DataDir = path diff --git a/translate.go b/translate.go index 41f3b4852..f1eddfbfe 100644 --- a/translate.go +++ b/translate.go @@ -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,26 @@ 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 { + // 10GB default map size + defaultMapSize := 10 * (1 << 30) + // Use 2GB default map size on 32-bit systems + if 32<<(^uint(0)>>32&1) == 32 { + defaultMapSize = (1 << 31) - 1 // 2GB + } + + f := &TranslateFile{ writeNotify: make(chan struct{}), closing: make(chan struct{}), cols: make(map[string]*index), @@ -96,6 +113,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) { diff --git a/translate_mapsize_386.go b/translate_mapsize_386.go deleted file mode 100644 index 259b735ac..000000000 --- a/translate_mapsize_386.go +++ /dev/null @@ -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 diff --git a/translate_mapsize_all64bitsystems.go b/translate_mapsize_all64bitsystems.go deleted file mode 100644 index 605d6270a..000000000 --- a/translate_mapsize_all64bitsystems.go +++ /dev/null @@ -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 From 2394d3610848aac94355c7ebeb81dfc9ccaec303 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Wed, 19 Sep 2018 13:30:22 -0500 Subject: [PATCH 2/8] Adjust tests, fix 32-bit config --- cmd/server_test.go | 8 +++++--- server/server_test.go | 2 ++ test/pilosa.go | 2 +- translate.go | 14 +++++++++----- translate_test.go | 6 +++--- 5 files changed, 20 insertions(+), 12 deletions(-) diff --git a/cmd/server_test.go b/cmd/server_test.go index e58f8af4d..84b2297d4 100644 --- a/cmd/server_test.go +++ b/cmd/server_test.go @@ -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" diff --git a/server/server_test.go b/server/server_test.go index ab860dc1b..a15bb65e4 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -403,6 +403,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 @@ -476,6 +477,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 diff --git a/test/pilosa.go b/test/pilosa.go index 835291cbd..a2e6c1c58 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -66,7 +66,7 @@ func newCommand(opts ...server.CommandOption) *Command { server.OptCommandServerOptions( pilosa.OptServerHolderOptions( pilosa.OptHolderTranslateFileMapSize( - 2 << 26, + 2 << 25, ), ), ), diff --git a/translate.go b/translate.go index f1eddfbfe..fc4a9d3c7 100644 --- a/translate.go +++ b/translate.go @@ -93,11 +93,15 @@ func OptTranslateFileMapSize(mapSize int) TranslateFileOption { // NewTranslateFile returns a new instance of TranslateFile. func NewTranslateFile(opts ...TranslateFileOption) *TranslateFile { - // 10GB default map size - defaultMapSize := 10 * (1 << 30) - // Use 2GB default map size on 32-bit systems - if 32<<(^uint(0)>>32&1) == 32 { - defaultMapSize = (1 << 31) - 1 // 2GB + var defaultMapSize64 int64 = 1 << 33 + 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{ diff --git a/translate_test.go b/translate_test.go index 706fe365f..1aea26cfe 100644 --- a/translate_test.go +++ b/translate_test.go @@ -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) From ebc4d7fc13458e11e33a2f725cccbb0ec42d4dfa Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Wed, 19 Sep 2018 13:37:42 -0500 Subject: [PATCH 3/8] Set low map size for 32-bit --- server/server_test.go | 1 + 1 file changed, 1 insertion(+) diff --git a/server/server_test.go b/server/server_test.go index a15bb65e4..faca9c9d0 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -498,6 +498,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)) From 42293c58a5ee088e88d51aba32be80f1bf8fa497 Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Wed, 19 Sep 2018 14:11:33 -0500 Subject: [PATCH 4/8] Do not run prerelease in CI if this is a pull request --- .circleci/config.yml | 1 + 1 file changed, 1 insertion(+) diff --git a/.circleci/config.yml b/.circleci/config.yml index c5bf059e2..b42b2aa64 100644 --- a/.circleci/config.yml +++ b/.circleci/config.yml @@ -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 From 995a24d0afbe556dbfd00a316c695f5f07c88c11 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Wed, 19 Sep 2018 15:05:16 -0500 Subject: [PATCH 5/8] ensure mutex imports unset previous columns --- fragment.go | 121 +++++++++++++++++++++++++++++++++++--- fragment_internal_test.go | 70 ++++++++++++++++++++++ 2 files changed, 182 insertions(+), 9 deletions(-) diff --git a/fragment.go b/fragment.go index 7f651b135..7e3cd26d9 100644 --- a/fragment.go +++ b/fragment.go @@ -1330,16 +1330,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 +1362,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 +1402,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() diff --git a/fragment_internal_test.go b/fragment_internal_test.go index dd2d00b41..e0865c606 100644 --- a/fragment_internal_test.go +++ b/fragment_internal_test.go @@ -1178,6 +1178,76 @@ 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) + } + } + }) + } +} + func BenchmarkFragment_Snapshot(b *testing.B) { if *FragmentPath == "" { b.Skip("no fragment specified") From 52ba336461be7d5ec2f4caba98cf4b38ce05ddda Mon Sep 17 00:00:00 2001 From: Cody Soyland Date: Thu, 20 Sep 2018 16:05:28 -0500 Subject: [PATCH 6/8] Address code review feedback --- ctl/server.go | 2 +- docs/configuration.md | 12 ++++++++++++ holder.go | 23 ++--------------------- server.go | 9 ++------- server/server.go | 6 ++---- test/pilosa.go | 6 ++---- translate.go | 3 +-- 7 files changed, 22 insertions(+), 39 deletions(-) diff --git a/ctl/server.go b/ctl/server.go index 194633a48..7f384ec7b 100644 --- a/ctl/server.go +++ b/ctl/server.go @@ -45,7 +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 of mmap to allocate for key translation.") + 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.") diff --git a/docs/configuration.md b/docs/configuration.md index 22a67ddcf..5750441ef 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -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: diff --git a/holder.go b/holder.go index 5cffeea9b..794104e53 100644 --- a/holder.go +++ b/holder.go @@ -77,19 +77,9 @@ type Holder struct { Logger Logger } -// HolderOption is a functional option type for pilosa.Holder -type HolderOption func(f *Holder) error - -func OptHolderTranslateFileMapSize(mapSize int) HolderOption { - return func(h *Holder) error { - h.translateFile = NewTranslateFile(OptTranslateFileMapSize(mapSize)) - return nil - } -} - // NewHolder returns a new instance of Holder. -func NewHolder(opts ...HolderOption) *Holder { - h := &Holder{ +func NewHolder() *Holder { + return &Holder{ indexes: make(map[string]*Index), closing: make(chan struct{}), @@ -107,15 +97,6 @@ func NewHolder(opts ...HolderOption) *Holder { Logger: NopLogger, } - - for _, opt := range opts { - err := opt(h) - if err != nil { - // TODO (2.0): Change func signature to return error - panic(errors.Wrap(err, "applying option")) - } - } - return h } // Open initializes the root data directory for the holder. diff --git a/server.go b/server.go index 708fd853b..ce6319f45 100644 --- a/server.go +++ b/server.go @@ -234,14 +234,9 @@ func OptServerClusterHasher(h Hasher) ServerOption { } } -func OptServerHolderOptions(opts ...HolderOption) ServerOption { +func OptServerTranslateFileMapSize(mapSize int) ServerOption { return func(s *Server) error { - for _, opt := range opts { - err := opt(s.holder) - if err != nil { - return errors.Wrap(err, "applying option") - } - } + s.holder.translateFile = NewTranslateFile(OptTranslateFileMapSize(mapSize)) return nil } } diff --git a/server/server.go b/server/server.go index 61942522f..b4ac595fa 100644 --- a/server/server.go +++ b/server/server.go @@ -286,10 +286,8 @@ func (m *Command) SetupServer() error { if m.Config.Translation.MapSize > 0 { serverOptions = append( serverOptions, - pilosa.OptServerHolderOptions( - pilosa.OptHolderTranslateFileMapSize( - m.Config.Translation.MapSize, - ), + pilosa.OptServerTranslateFileMapSize( + m.Config.Translation.MapSize, ), ) } diff --git a/test/pilosa.go b/test/pilosa.go index a2e6c1c58..f8f60eebe 100644 --- a/test/pilosa.go +++ b/test/pilosa.go @@ -64,10 +64,8 @@ func newCommand(opts ...server.CommandOption) *Command { opts = append([]server.CommandOption{ server.OptCommandCloseTimeout(time.Millisecond * 2), server.OptCommandServerOptions( - pilosa.OptServerHolderOptions( - pilosa.OptHolderTranslateFileMapSize( - 2 << 25, - ), + pilosa.OptServerTranslateFileMapSize( + 2 << 25, ), ), }, opts...) diff --git a/translate.go b/translate.go index fc4a9d3c7..2ef289f45 100644 --- a/translate.go +++ b/translate.go @@ -93,7 +93,7 @@ func OptTranslateFileMapSize(mapSize int) TranslateFileOption { // NewTranslateFile returns a new instance of TranslateFile. func NewTranslateFile(opts ...TranslateFileOption) *TranslateFile { - var defaultMapSize64 int64 = 1 << 33 + var defaultMapSize64 int64 = 10 * (1 << 30) var defaultMapSize int if ^uint(0)>>32 > 0 { @@ -103,7 +103,6 @@ func NewTranslateFile(opts ...TranslateFileOption) *TranslateFile { // Use 2GB default map size on 32-bit systems defaultMapSize = (1 << 31) - 1 } - f := &TranslateFile{ writeNotify: make(chan struct{}), closing: make(chan struct{}), From 4f17dbdbf1872742720c890baea3542eed7abe47 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Tue, 18 Sep 2018 10:34:31 -0500 Subject: [PATCH 7/8] add support for Bool fields prevent import of non-boolean row values to bool fields --- ctl/import_test.go | 46 +++++++++++++++++++++++ docs/administration.md | 14 +++++++ docs/api-reference.md | 2 + docs/data-model.md | 6 ++- docs/glossary.md | 2 +- docs/tutorials.md | 2 +- executor.go | 75 +++++++++++++++++++++++++------------- executor_test.go | 74 ++++++++++++++++++++++++++++++++++++- field.go | 37 +++++++++++++++++++ fragment.go | 34 +++++++++++++++++ fragment_internal_test.go | 77 +++++++++++++++++++++++++++++++++++++++ http/handler.go | 16 ++++++++ pql/ast.go | 17 +++++++++ view.go | 2 + 14 files changed, 375 insertions(+), 29 deletions(-) diff --git a/ctl/import_test.go b/ctl/import_test.go index cd75a5d0a..2913005fb 100644 --- a/ctl/import_test.go +++ b/ctl/import_test.go @@ -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) + } + }) +} diff --git a/docs/administration.md b/docs/administration.md index cd4e30184..e626cc0e5 100644 --- a/docs/administration.md +++ b/docs/administration.md @@ -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 +``` +

Note that you must first create a field. View Create Field for more details. The `-e` flag can create the necessary schema when using a field of type "set".

diff --git a/docs/api-reference.md b/docs/api-reference.md index 890b37941..216a01bca 100644 --- a/docs/api-reference.md +++ b/docs/api-reference.md @@ -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` diff --git a/docs/data-model.md b/docs/data-model.md index ba7c12024..bababdd34 100644 --- a/docs/data-model.md +++ b/docs/data-model.md @@ -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. diff --git a/docs/glossary.md b/docs/glossary.md index c39a0bf6f..fccbe0ae8 100644 --- a/docs/glossary.md +++ b/docs/glossary.md @@ -22,7 +22,7 @@ nav = [] Fragment: A Fragment is the intersection of a [field](#field) and a [shard](#shard) in an [index](#index). -[Field](../data-model/#field): 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). +[Field](../data-model/#field): 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). [Frame](../data-model/#field): Prior to Pilosa 1.0, fields were known as frames. diff --git a/docs/tutorials.md b/docs/tutorials.md index 040d85a1c..47344db4f 100644 --- a/docs/tutorials.md +++ b/docs/tutorials.md @@ -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 diff --git a/executor.go b/executor.go index b0083a514..4ae5ac81e 100644 --- a/executor.go +++ b/executor.go @@ -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 translater 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 { diff --git a/executor_test.go b/executor_test.go index ee919f001..828a40b09 100644 --- a/executor_test.go +++ b/executor_test.go @@ -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) diff --git a/field.go b/field.go index d038b606c..822522bc4 100644 --- a/field.go +++ b/field.go @@ -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") } diff --git a/fragment.go b/fragment.go index 7e3cd26d9..312b4af10 100644 --- a/fragment.go +++ b/fragment.go @@ -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. @@ -2206,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) {} diff --git a/fragment_internal_test.go b/fragment_internal_test.go index e0865c606..87f4e6137 100644 --- a/fragment_internal_test.go +++ b/fragment_internal_test.go @@ -1248,6 +1248,76 @@ func TestFragment_ImportMutex(t *testing.T) { } } +// 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") @@ -1372,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 { diff --git a/http/handler.go b/http/handler.go index bc108568a..820bc6521 100644 --- a/http/handler.go +++ b/http/handler.go @@ -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) } diff --git a/pql/ast.go b/pql/ast.go index cd483e030..0e2ff6aef 100644 --- a/pql/ast.go +++ b/pql/ast.go @@ -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 diff --git a/view.go b/view.go index 2feac8b44..5fc0f5eb8 100644 --- a/view.go +++ b/view.go @@ -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 } From 4c1709b8bd36acfb1110b8fc7aff6ea3251c560a Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Fri, 21 Sep 2018 09:24:48 -0500 Subject: [PATCH 8/8] fix spelling mistake --- executor.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/executor.go b/executor.go index 4ae5ac81e..ff3db909d 100644 --- a/executor.go +++ b/executor.go @@ -1675,7 +1675,7 @@ func (e *executor) translateCall(index string, idx *Index, c *pql.Call) error { return nil } - // Bool field keys do not use the translater because there + // Bool field keys do not use the translator because there // are only two possible values. Instead, they are handled // directly. if field.Type() == FieldTypeBool {