From cdfcd6db5ec822080ca7efbf0d84da5daf028f9e Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Mon, 24 Sep 2018 14:56:29 -0500 Subject: [PATCH] store remoteAvailableShards to file --- field.go | 61 ++++++++++++++++++++++++++++++++++-- field_internal_test.go | 20 ++++++++++++ server.go | 8 +++-- server/server_test.go | 71 ++++++++++++++++++++++++++++++++++++++++++ 4 files changed, 156 insertions(+), 4 deletions(-) diff --git a/field.go b/field.go index 822522bc4..6d90a57ce 100644 --- a/field.go +++ b/field.go @@ -15,6 +15,7 @@ package pilosa import ( + "bytes" "encoding/json" "fmt" "io/ioutil" @@ -230,8 +231,17 @@ func (f *Field) AvailableShards() *roaring.Bitmap { return b } -// addRemoteAvailableShards merges the set of available shards into the current known set. -func (f *Field) addRemoteAvailableShards(b *roaring.Bitmap) { +// addRemoteAvailableShards merges the set of available shards into the current known set +// and saves the set to a file. +func (f *Field) addRemoteAvailableShards(b *roaring.Bitmap) error { + f.mergeRemoteAvailableShards(b) + + // Save the updated bitmap to the data store. + return f.saveAvailableShards() +} + +// mergeRemoteAvailableShards merges the set of available shards into the current known set. +func (f *Field) mergeRemoteAvailableShards(b *roaring.Bitmap) { f.mu.Lock() defer f.mu.Unlock() f.remoteAvailableShards = f.remoteAvailableShards.Union(b) @@ -291,6 +301,10 @@ func (f *Field) Open() error { return errors.Wrap(err, "loading meta") } + if err := f.loadAvailableShards(); err != nil { + return errors.Wrap(err, "loading available shards") + } + // Apply the field options loaded from meta. if err := f.applyOptions(f.options); err != nil { return errors.Wrap(err, "applying options") @@ -467,6 +481,49 @@ func (f *Field) applyOptions(opt FieldOptions) error { return nil } +// loadAvailableShards reads remoteAvailableShards data for the field, if any. +func (f *Field) loadAvailableShards() error { + bm := roaring.NewBitmap() + + // Read data from meta file. + buf, err := ioutil.ReadFile(filepath.Join(f.path, ".available.shards")) + if os.IsNotExist(err) { + return nil + } else if err != nil { + return errors.Wrap(err, "reading available shards") + } else { + if err := bm.UnmarshalBinary(buf); err != nil { + return errors.Wrap(err, "unmarshaling") + } + } + + // Merge bitmap from file into field. + f.mergeRemoteAvailableShards(bm) + + return nil +} + +// saveAvailableShards writes remoteAvailableShards data for the field. +func (f *Field) saveAvailableShards() error { + f.mu.RLock() + defer f.mu.RUnlock() + + // Write available shards to buffer. + var buf bytes.Buffer + if n, err := f.remoteAvailableShards.WriteTo(&buf); err != nil { + return errors.Wrap(err, "writing bitmap to buffer") + } else if n != int64(buf.Len()) { + return fmt.Errorf("buffer size mismatch: %d != %d", n, buf.Len()) + } + + // Write buffer to file. + if err := ioutil.WriteFile(filepath.Join(f.path, ".available.shards"), buf.Bytes(), 0666); err != nil { + return errors.Wrap(err, "writing available shards") + } + + return nil +} + // Close closes the field and its views. func (f *Field) Close() error { f.mu.Lock() diff --git a/field_internal_test.go b/field_internal_test.go index 2d37ea25c..d72a3fd6a 100644 --- a/field_internal_test.go +++ b/field_internal_test.go @@ -22,6 +22,7 @@ import ( "time" "github.com/pilosa/pilosa/pql" + "github.com/pilosa/pilosa/roaring" ) // Ensure a bsiGroup can adjust to its baseValue. @@ -341,3 +342,22 @@ func TestField_RowTime(t *testing.T) { } } + +func TestField_PersistAvailableShards(t *testing.T) { + f := MustOpenField(OptFieldTypeDefault()) + + // bm represents remote available shards. + bm := roaring.NewBitmap(1, 2, 3) + + if err := f.addRemoteAvailableShards(bm); err != nil { + t.Fatal(err) + } + + // Reload field and verify that shard data is persisted. + if err := f.Reopen(); err != nil { + t.Fatal(err) + } else if !reflect.DeepEqual(f.remoteAvailableShards.Slice(), bm.Slice()) { + t.Fatalf("unexpected available shards (reopen). expected: %v, but got: %v", bm.Slice(), f.remoteAvailableShards.Slice()) + } + +} diff --git a/server.go b/server.go index 81fc434f9..6cd313dc0 100644 --- a/server.go +++ b/server.go @@ -481,7 +481,9 @@ func (s *Server) receiveMessage(m Message) error { if f == nil { return fmt.Errorf("Local field not found: %s/%s", obj.Index, obj.Field) } - f.addRemoteAvailableShards(roaring.NewBitmap(obj.Shard)) + if err := f.addRemoteAvailableShards(roaring.NewBitmap(obj.Shard)); err != nil { + return errors.Wrap(err, "adding remote available shards") + } case *CreateIndexMessage: opt := obj.Meta _, err := s.holder.CreateIndex(obj.Index, *opt) @@ -646,7 +648,9 @@ func (s *Server) mergeRemoteStatus(ns *NodeStatus) error { s.logger.Printf("Local Field not found: %s/%s", is.Name, fs.Name) continue } - f.addRemoteAvailableShards(fs.AvailableShards) + if err := f.addRemoteAvailableShards(fs.AvailableShards); err != nil { + return errors.Wrap(err, "adding remote available shards") + } } } diff --git a/server/server_test.go b/server/server_test.go index ed73975d7..aa202a54c 100644 --- a/server/server_test.go +++ b/server/server_test.go @@ -621,4 +621,75 @@ func TestMain_ImportTimestamp(t *testing.T) { } } +func TestClusterQueriesAfterRestart(t *testing.T) { + cluster := test.MustRunCluster(t, 3) + defer cluster.Close() + cmd1 := cluster[1] + + cmd1.MustCreateIndex(t, "testidx", pilosa.IndexOptions{}) + cmd1.MustCreateField(t, "testidx", "testfield", pilosa.OptFieldTypeSet(pilosa.CacheTypeRanked, 10)) + + // build a query to set the first bit in 100 shards + query := strings.Builder{} + for i := 0; i < 100; i++ { + query.WriteString(fmt.Sprintf("Set(%d, testfield=0)", i*pilosa.ShardWidth)) + } + _, err := cmd1.API.Query(context.Background(), &pilosa.QueryRequest{ + Index: "testidx", + Query: query.String(), + }) + if err != nil { + t.Fatalf("setting 100 bits in 100 shards: %v", err) + } + + results, err := cmd1.API.Query(context.Background(), &pilosa.QueryRequest{ + Index: "testidx", + Query: "Count(Row(testfield=0))", + }) + if err != nil { + t.Fatalf("counting row: %v", err) + } + if results.Results[0].(uint64) != 100 { + t.Fatalf("Count should be 100, but got %v of type %[1]T", results.Results[0]) + } + + err = cmd1.Command.Close() + if err != nil { + t.Fatalf("closing node0: %v", err) + } + + // confirm that cluster stops accepting queries after one node closes + if _, err := cluster[0].API.Query(context.Background(), &pilosa.QueryRequest{}); !strings.Contains(err.Error(), "not allowed in state STARTING") { + t.Fatalf("got unexpected error querying an incomplete cluster: %v", err) + } + + // Create new main with the same config. + config := cmd1.Command.Config + config.Bind = cmd1.API.Node().URI.HostPort() + + // this isn't necessary, but makes the test run way faster + config.Gossip.Port = strconv.Itoa(int(cmd1.Command.GossipTransport().URI.Port)) + cmd1.Command = server.NewCommand(cmd1.Stdin, cmd1.Stdout, cmd1.Stderr) + cmd1.Command.Config = config + err = cmd1.Start() + if err != nil { + t.Fatalf("reopening node 0: %v", err) + } + + for cmd1.API.State() != pilosa.ClusterStateNormal { + time.Sleep(time.Millisecond) + } + + results, err = cmd1.API.Query(context.Background(), &pilosa.QueryRequest{ + Index: "testidx", + Query: "Count(Row(testfield=0))", + }) + if err != nil { + t.Fatalf("counting row: %v", err) + } + if results.Results[0].(uint64) != 100 { + t.Fatalf("Count should be 100, but got %v of type %[1]T", results.Results[0]) + } +} + // TODO: confirm that things keep working if a node is hard-closed (no nodeLeave event) and immediately restarted with a different address.