From 30059f9e7e3a30f1d77c93c35ccdec2dc996599b Mon Sep 17 00:00:00 2001 From: Todd Gruben Date: Wed, 17 Oct 2018 12:06:54 -0500 Subject: [PATCH 1/3] force truncate of available shards file to avoid corruption --- field.go | 87 ++++++++++++++++++++++-------------------- field_internal_test.go | 39 +++++++++++++++++++ 2 files changed, 84 insertions(+), 42 deletions(-) diff --git a/field.go b/field.go index 27de57d8f..e3be986f8 100644 --- a/field.go +++ b/field.go @@ -15,6 +15,7 @@ package pilosa import ( + "bufio" "encoding/json" "fmt" "io/ioutil" @@ -234,7 +235,6 @@ func (f *Field) AvailableShards() *roaring.Bitmap { // 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() } @@ -246,6 +246,50 @@ func (f *Field) mergeRemoteAvailableShards(b *roaring.Bitmap) { f.remoteAvailableShards = f.remoteAvailableShards.Union(b) } +// loadAvailableShards reads remoteAvailableShards data for the field, if any. +func (f *Field) loadAvailableShards() error { + bm := roaring.NewBitmap() + // Read data from meta file. + path := filepath.Join(f.path, ".available.shards") + buf, err := ioutil.ReadFile(path) + 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 { + // Open or create file. + f.mu.Lock() + defer f.mu.Unlock() + path := filepath.Join(f.path, ".available.shards") + + file, err := os.OpenFile(path, os.O_WRONLY|os.O_CREATE|os.O_TRUNC, 0666) + if err != nil { + return errors.Wrap(err, "opening available shards file") + } + defer file.Close() + + // Write available shards to file. + bw := bufio.NewWriter(file) + if _, err = f.remoteAvailableShards.WriteTo(bw); err != nil { + return errors.Wrap(err, "writing bitmap to buffer") + } + bw.Flush() + + return nil +} + // Type returns the field type. func (f *Field) Type() string { f.mu.RLock() @@ -480,47 +524,6 @@ 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 { - // Open or create file. - file, err := os.OpenFile(filepath.Join(f.path, ".available.shards"), os.O_WRONLY|os.O_CREATE, 0666) - if err != nil { - return errors.Wrap(err, "opening available shards file") - } - - f.mu.RLock() - defer f.mu.RUnlock() - - // Write available shards to file. - if _, err := f.remoteAvailableShards.WriteTo(file); err != nil { - return errors.Wrap(err, "writing bitmap to buffer") - } - - 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 d72a3fd6a..56ca06094 100644 --- a/field_internal_test.go +++ b/field_internal_test.go @@ -361,3 +361,42 @@ func TestField_PersistAvailableShards(t *testing.T) { } } + +func TestField_PersistAvailableShardsTaxiBug(t *testing.T) { + f := MustOpenField(OptFieldTypeDefault()) + + // bm represents remote available shards. + bm := roaring.NewBitmap() + for i := uint64(0); i < 1204; i += 2 { + bm.Add(i) + } + + 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()) + } + + bm1 := roaring.NewBitmap() + for i := uint64(1); i < 1204; i += 2 { + bm1.Add(i) + } + + if err := f.addRemoteAvailableShards(bm1); err != nil { + t.Fatal(err) + } + + // Reload field and verify that shard data is persisted. + result := bm.Union(bm1) + if err := f.Reopen(); err != nil { + t.Fatal(err) + } else if !reflect.DeepEqual(f.remoteAvailableShards.Slice(), result.Slice()) { + t.Fatalf("unexpected available shards (reopen). expected: %v, but got: %v", bm.Slice(), f.remoteAvailableShards.Slice()) + } + +} From 03e68aaa69e80299b4ca78b92bd32af90ad578f5 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Wed, 17 Oct 2018 12:38:29 -0500 Subject: [PATCH 2/3] Update field_internal_test.go --- field_internal_test.go | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/field_internal_test.go b/field_internal_test.go index a4d1701a4..c2495e9c0 100644 --- a/field_internal_test.go +++ b/field_internal_test.go @@ -362,7 +362,9 @@ func TestField_PersistAvailableShards(t *testing.T) { } -func TestField_PersistAvailableShardsTaxiBug(t *testing.T) { +// Ensure that persisting available shards having a smaller footprint (for example, +// when going from a bitmap to a smaller, RLE representation) succeeds. +func TestField_PersistAvailableShardsFootprint(t *testing.T) { f := MustOpenField(OptFieldTypeDefault()) // bm represents remote available shards. From 9cd0782ef1cf759550ecbdd357fcd1c5290b9b62 Mon Sep 17 00:00:00 2001 From: Travis Turner Date: Wed, 17 Oct 2018 17:33:13 -0500 Subject: [PATCH 3/3] add unprotectedSaveAvailableShards() method --- field.go | 12 ++++++++---- 1 file changed, 8 insertions(+), 4 deletions(-) diff --git a/field.go b/field.go index 462af4f95..63bb08b63 100644 --- a/field.go +++ b/field.go @@ -269,9 +269,13 @@ func (f *Field) loadAvailableShards() error { // saveAvailableShards writes remoteAvailableShards data for the field. func (f *Field) saveAvailableShards() error { - // Open or create file. f.mu.Lock() defer f.mu.Unlock() + return f.unprotectedSaveAvailableShards() +} + +func (f *Field) unprotectedSaveAvailableShards() error { + // Open or create file. path := filepath.Join(f.path, ".available.shards") file, err := os.OpenFile(path, os.O_WRONLY|os.O_CREATE|os.O_TRUNC, 0666) @@ -294,8 +298,8 @@ func (f *Field) saveAvailableShards() error { // // NOTE: This can be overridden on the next sync so all nodes should be updated. func (f *Field) RemoveAvailableShard(v uint64) error { - f.mu.RLock() - defer f.mu.RUnlock() + f.mu.Lock() + defer f.mu.Unlock() b := f.remoteAvailableShards.Clone() if _, err := b.Remove(v); err != nil { @@ -303,7 +307,7 @@ func (f *Field) RemoveAvailableShard(v uint64) error { } f.remoteAvailableShards = b - return f.saveAvailableShards() + return f.unprotectedSaveAvailableShards() } // Type returns the field type.