From 8bc110458568bb0e62dcb37e3a07d3c38ac03032 Mon Sep 17 00:00:00 2001 From: Matt Jaffee Date: Tue, 20 Nov 2018 08:56:12 -0600 Subject: [PATCH] fix fragment checksums race condition --- fragment.go | 6 +++--- fragment_internal_test.go | 17 +++++++++++++++++ 2 files changed, 20 insertions(+), 3 deletions(-) diff --git a/fragment.go b/fragment.go index a4e4fd72d..bab1f3ec0 100644 --- a/fragment.go +++ b/fragment.go @@ -1492,9 +1492,6 @@ func (f *fragment) bulkImportStandard(rowIDs, columnIDs []uint64, options *Impor lastRowID = rowID rowSet[rowID] = struct{}{} } - - // Invalidate block checksum. - delete(f.checksums, int(rowID/HashBlockSize)) } f.mu.Lock() @@ -1518,6 +1515,9 @@ func (f *fragment) bulkImportStandard(rowIDs, columnIDs []uint64, options *Impor // Update cache counts for all affected rows. for rowID := range rowSet { + // Invalidate block checksum. + delete(f.checksums, int(rowID/HashBlockSize)) + n := results.CountRange(rowID*ShardWidth, (rowID+1)*ShardWidth) f.cache.BulkAdd(rowID, n) } diff --git a/fragment_internal_test.go b/fragment_internal_test.go index 4e3007957..b36bf3380 100644 --- a/fragment_internal_test.go +++ b/fragment_internal_test.go @@ -25,6 +25,8 @@ import ( "testing" "testing/quick" + "golang.org/x/sync/errgroup" + "github.com/davecgh/go-spew/spew" "github.com/pilosa/pilosa/pql" "github.com/pilosa/pilosa/roaring" @@ -1399,6 +1401,21 @@ func TestFragment_ImportSet(t *testing.T) { } } +func TestFragment_ConcurrentImport(t *testing.T) { + t.Run("bulkImportStandard", func(t *testing.T) { + f := mustOpenFragment("i", "f", viewStandard, 0, "") + defer f.Close() + + eg := errgroup.Group{} + eg.Go(func() error { return f.bulkImportStandard([]uint64{1, 2}, []uint64{1, 2}, &ImportOptions{}) }) + eg.Go(func() error { return f.bulkImportStandard([]uint64{3, 4}, []uint64{3, 4}, &ImportOptions{}) }) + err := eg.Wait() + if err != nil { + t.Fatalf("importing data to fragment: %v", err) + } + }) +} + // Ensure a fragment can import mutually exclusive values. func TestFragment_ImportMutex(t *testing.T) { tests := []struct {