diff --git a/fragment.go b/fragment.go index 78ef0597f..6a801bc50 100644 --- a/fragment.go +++ b/fragment.go @@ -1776,8 +1776,8 @@ func (f *fragment) mergeBlock(id int, data []pairSet) (sets, clears []pairSet, e clears = make([]pairSet, len(data)+1) // Limit upper row/column pair. - maxRowID := uint64(id+1) * HashBlockSize - maxColumnID := uint64(ShardWidth) + maxRowID := (uint64(id+1) * HashBlockSize) - 1 + maxColumnID := uint64(ShardWidth) - 1 // Create buffered iterator for local block. itrs := make([]*bufIterator, 1, len(data)+1) @@ -1857,8 +1857,8 @@ func (f *fragment) mergeBlock(id int, data []pairSet) (sets, clears []pairSet, e sets[i].rowIDs = append(sets[i].rowIDs, min.rowID) sets[i].columnIDs = append(sets[i].columnIDs, min.columnID) } else { - clears[i].rowIDs = append(sets[i].rowIDs, min.rowID) - clears[i].columnIDs = append(sets[i].columnIDs, min.columnID) + clears[i].rowIDs = append(clears[i].rowIDs, min.rowID) + clears[i].columnIDs = append(clears[i].columnIDs, min.columnID) } } } diff --git a/holder_test.go b/holder_test.go index d9d11f3f6..42c010079 100644 --- a/holder_test.go +++ b/holder_test.go @@ -490,6 +490,112 @@ func TestHolderSyncer_SyncHolder(t *testing.T) { } } +// Ensure holder can sync with a remote holder and respects +// the row boundaries of the block. +func TestHolderSyncer_BlockIteratorLimits(t *testing.T) { + c := test.MustNewCluster(t, 3) + c[0].Config.Cluster.ReplicaN = 3 + c[0].Config.AntiEntropy.Interval = 0 + c[1].Config.Cluster.ReplicaN = 3 + c[1].Config.AntiEntropy.Interval = 0 + err := c.Start() + if err != nil { + t.Fatalf("starting cluster: %v", err) + } + defer c.Close() + + _, err = c[0].API.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}) + if err != nil { + t.Fatalf("creating index i: %v", err) + } + _, err = c[0].API.CreateField(context.Background(), "i", "f", pilosa.OptFieldTypeSet(pilosa.DefaultCacheType, pilosa.DefaultCacheSize)) + if err != nil { + t.Fatalf("creating field f: %v", err) + } + + blockEdge := uint64(pilosa.HashBlockSize) + + hldr0 := &test.Holder{Holder: c[0].Server.Holder()} + hldr1 := &test.Holder{Holder: c[1].Server.Holder()} + hldr2 := &test.Holder{Holder: c[2].Server.Holder()} + + // Set data on the local holder. + hldr0.SetBit("i", "f", blockEdge-1, 10) + hldr0.SetBit("i", "f", blockEdge, 20) + + // Set the same data on one of the replicas + // so that we have a quorum. + hldr1.SetBit("i", "f", blockEdge-1, 10) + hldr1.SetBit("i", "f", blockEdge, 20) + + // Leave the third replica empty to force a block merge. + // + + err = c[0].Server.SyncData() + if err != nil { + t.Fatalf("syncing node 0: %v", err) + } + + // Verify data is the same on all nodes. + for i, hldr := range []*test.Holder{hldr0, hldr1, hldr2} { + if a := hldr.Row("i", "f", blockEdge-1).Columns(); !reflect.DeepEqual(a, []uint64{10}) { + t.Errorf("unexpected columns(%d/block 0): %+v", i, a) + } + if a := hldr.Row("i", "f", blockEdge).Columns(); !reflect.DeepEqual(a, []uint64{20}) { + t.Errorf("unexpected columns(%d/block 1): %+v", i, a) + } + } +} + +// Ensure holder correctly handles clears during block sync. +func TestHolderSyncer_Clears(t *testing.T) { + c := test.MustNewCluster(t, 3) + c[0].Config.Cluster.ReplicaN = 3 + c[0].Config.AntiEntropy.Interval = 0 + c[1].Config.Cluster.ReplicaN = 3 + c[1].Config.AntiEntropy.Interval = 0 + err := c.Start() + if err != nil { + t.Fatalf("starting cluster: %v", err) + } + defer c.Close() + + _, err = c[0].API.CreateIndex(context.Background(), "i", pilosa.IndexOptions{}) + if err != nil { + t.Fatalf("creating index i: %v", err) + } + _, err = c[0].API.CreateField(context.Background(), "i", "f", pilosa.OptFieldTypeSet(pilosa.DefaultCacheType, pilosa.DefaultCacheSize)) + if err != nil { + t.Fatalf("creating field f: %v", err) + } + + hldr0 := &test.Holder{Holder: c[0].Server.Holder()} + hldr1 := &test.Holder{Holder: c[1].Server.Holder()} + hldr2 := &test.Holder{Holder: c[2].Server.Holder()} + + // Set data on the local holder that should be cleared + // because it's the only instance of this value. + hldr0.SetBit("i", "f", 0, 30) + + // Set similar data on the replicas, but + // different from what's on local. This should end + // up being set on all replicas + hldr1.SetBit("i", "f", 0, 20) + hldr2.SetBit("i", "f", 0, 20) + + err = c[0].Server.SyncData() + if err != nil { + t.Fatalf("syncing node 0: %v", err) + } + + // Verify data is the same on all nodes. + for i, hldr := range []*test.Holder{hldr0, hldr1, hldr2} { + if a := hldr.Row("i", "f", 0).Columns(); !reflect.DeepEqual(a, []uint64{20}) { + t.Errorf("unexpected columns(%d): %+v", i, a) + } + } +} + // Ensure holder can sync time quantum views with a remote holder. func TestHolderSyncer_TimeQuantum(t *testing.T) { c := test.MustNewCluster(t, 2)