diff --git a/fragment.go b/fragment.go index 622341574..4a42a7255 100644 --- a/fragment.go +++ b/fragment.go @@ -177,7 +177,7 @@ func (f *fragment) Open() error { if err := func() error { // Initialize storage in a function so we can close if anything goes wrong. f.Logger.Debugf("open storage for index/field/view/fragment: %s/%s/%s/%d", f.index, f.field, f.view, f.shard) - if err := f.openStorage(); err != nil { + if err := f.openStorage(true); err != nil { return errors.Wrap(err, "opening storage") } @@ -215,12 +215,43 @@ func (f *fragment) reopen() (mustClose bool, err error) { return mustClose, nil } -// openStorage opens the storage bitmap. -func (f *fragment) openStorage() error { +// openStorage opens the storage bitmap. Usually you also want to read in +// the storage, but in the case where we just wrote that file, such as +// unprotectedWriteToFragment, we could also just... not. If we didn't +// have existing storage, we probably need to unmarshal the data. If the +// file we're asked to open is empty, we probably don't. +// +// If we already had mapped storage previously, we want to unmap that, and +// possibly remap it from the file, but we don't need a full unmarshal, just +// an update of mapped pointers. +// +// unmarshalData is somewhat overloaded. it tells us whether or not we +// need to actually create a bitmap from the data (if the data exists to +// do this from). +// +// usually unmarshalData is only set to false when we're in the middle of +// a snapshot, and unprotectedWriteToFragment just wrote the in-memory data +// out. +// +// If we have existing storage data, and we successfully get new data, +// we will unmap the existing storage data. +// +// This function's design is probably a problem -- it is trying to handle +// both cases where there was existing data before, and cases where we +// just wrote the data. +func (f *fragment) openStorage(unmarshalData bool) error { + oldStorageData := f.storageData + // there's a few places where we might encounter an error, but need + // to continue past it through other error checks, before returning it. + var lastError error + // Create a roaring bitmap to serve as storage for the shard. if f.storage == nil { f.storage = roaring.NewFileBitmap() f.storage.Flags = f.flags + // if we didn't actually have storage, we *do* need to + // unmarshal this data in order to have any. + unmarshalData = true } // Open the data file to be mmap'd and used as an ops log. file, mustClose, err := syswrap.OpenFile(f.path, os.O_RDWR|os.O_CREATE|os.O_APPEND, 0666) @@ -237,13 +268,24 @@ func (f *fragment) openStorage() error { return fmt.Errorf("flock: %s", err) } + // data is the data we would unmarshal from, if we're unmarshalling; it might + // be obtained by calling ReadAll on a file. + // + // newStorageData is the data we should map things to. it is set only if + // mmapped; if we didn't mmap (say, we couldn't), we won't want to unmap + // the ioutil byte slice. (Theoretically, we shouldn't be using the mapped + // flag in that case...) + var data []byte + var newStorageData []byte + // If the file is empty then initialize it with an empty bitmap. fi, err := f.file.Stat() if err != nil { return errors.Wrap(err, "statting file before") } else if fi.Size() == 0 { bi := bufio.NewWriter(f.file) - if _, err := f.storage.WriteTo(bi); err != nil { + var err error + if _, err = f.storage.WriteTo(bi); err != nil { return fmt.Errorf("init storage file: %s", err) } bi.Flush() @@ -251,39 +293,94 @@ func (f *fragment) openStorage() error { if err != nil { return errors.Wrap(err, "statting file after") } + // there's nothing here, we're not going to try to unmarshal it. + unmarshalData = false + f.rowCache = &simpleCache{make(map[uint64]*Row)} } else { // Mmap the underlying file so it can be zero copied. - data, err := syswrap.Mmap(int(f.file.Fd()), 0, int(fi.Size()), syscall.PROT_READ, syscall.MAP_SHARED) + data, err = syswrap.Mmap(int(f.file.Fd()), 0, int(fi.Size()), syscall.PROT_READ, syscall.MAP_SHARED) if err == syswrap.ErrMaxMapCountReached { f.Logger.Debugf("maximum number of maps reached, reading file instead") - data, err = ioutil.ReadAll(file) - if err != nil { - return errors.Wrap(err, "failure file readall") + if unmarshalData { + data, err = ioutil.ReadAll(file) + if err != nil { + return errors.Wrap(err, "failure file readall") + } } } else if err != nil { return errors.Wrap(err, "mmap failed") } else { - f.storageData = data - // Advise the kernel that the mmap is accessed randomly. - if err := madvise(f.storageData, syscall.MADV_RANDOM); err != nil { - return fmt.Errorf("madvise: %s", err) - } + newStorageData = data } - - if err := f.storage.UnmarshalBinary(data); err != nil { - return fmt.Errorf("unmarshal storage: file=%s, err=%s", f.file.Name(), err) - } - } - f.opN = f.storage.Info().OpN + if unmarshalData { + f.storageData = newStorageData + // We're about to either re-read the bitmap, or fail to do so + // and unconditionally unmap the existing stuff. Either way, we + // want to unmap the old storage data after we're done here, but + // we can't unmap it yet because it's still live until sometime + // later, but we can't unmap it later, because we could return + // early... this is what defer is for. + if oldStorageData != nil { + defer func() { + unmapErr := syswrap.Munmap(oldStorageData) + if unmapErr != nil { + f.Logger.Printf("unmap of old storage failed: %s", err) + } + }() + } + // so we have a problem here: if this fails, it's unclear whether + // *either* or *both* of old and new storage data might be in use. + // So we call the thing that should unconditionally unmap both of them... + if err := f.storage.UnmarshalBinary(data); err != nil { + f.storage.RemapRoaringStorage(nil) + return fmt.Errorf("unmarshal storage: file=%s, err=%s", f.file.Name(), err) + } + f.rowCache = &simpleCache{make(map[uint64]*Row)} + f.opN = f.storage.OpN() + // slightly incorrect, but we can't measure it so... + } else { + // we're moving to new storage, so instead of using the OpN + // derived from reading that storage, we notify the bitmap that + // OpN is now effectively zero. + f.opN = 0 + f.storage.SetOpN(0) + // if oldStorageData is nil, this just tries to unmap any bits that + // are currently mapped. otherwise, it will point them at this + // storage (if the containers match). + var mappedAny bool + mappedAny, lastError = f.storage.RemapRoaringStorage(newStorageData) + if oldStorageData != nil { + unmapErr := syswrap.Munmap(oldStorageData) + if unmapErr != nil { + f.Logger.Printf("unmap of old storage failed: %s", err) + } + } + if mappedAny { + // Advise the kernel that the mmap is accessed randomly. + if err := madvise(newStorageData, syscall.MADV_RANDOM); err != nil { + lastError = fmt.Errorf("madvise: %s", err) + } + } else { + // if we did map data, but for some reason none of it got used + // as backing store, we can unmap it, and set the slice to nil, + // so we don't keep the now-invalid slice in f.storageData. + if newStorageData != nil { + unmapErr := syswrap.Munmap(newStorageData) + if unmapErr != nil { + lastError = fmt.Errorf("unmapping unused storage data: %s", err) + } + newStorageData = nil + } + } + f.storageData = newStorageData + } // Attach the file to the bitmap to act as a write-ahead log. f.storage.OpWriter = f.file - f.rowCache = &simpleCache{make(map[uint64]*Row)} - - return nil + return lastError } // openCache initializes the cache from row ids persisted to disk. @@ -343,7 +440,7 @@ func (f *fragment) close() error { } // Close underlying storage. - if err := f.closeStorage(); err != nil { + if err := f.closeStorage(true); err != nil { f.Logger.Printf("fragment: error closing storage: err=%s, path=%s", err, f.path) return errors.Wrap(err, "closing storage") } @@ -374,13 +471,19 @@ func (f *fragment) safeClose() error { return nil } -func (f *fragment) closeStorage() error { +// closeStorage attempts to close storage, including unmapping the old +// storage if includeMap is true. This would normally make sense if you're +// expecting to be done using the fragment, or to reload it. But it's also +// okay to just leave stuff mmapped; you don't have to keep the file +// descriptor open. So in some cases, we'll just leave the old mmapping +// in place, rather than regenerating everything from the new file. +func (f *fragment) closeStorage(includeMap bool) error { // Clear the storage bitmap so it doesn't access the closed mmap. //f.storage = roaring.NewBitmap() // Unmap the file. - if f.storageData != nil { + if includeMap && f.storageData != nil { if err := syswrap.Munmap(f.storageData); err != nil { return fmt.Errorf("munmap: %s", err) } @@ -1994,8 +2097,8 @@ func (f *fragment) importValueSmallWrite(columnIDs []uint64, values []int64, bit } return nil }(); err != nil { - _ = f.closeStorage() - _ = f.openStorage() + _ = f.closeStorage(true) + _ = f.openStorage(true) return err } rowSet := make(map[uint64]struct{}, bitDepth+1) @@ -2033,8 +2136,8 @@ func (f *fragment) importValue(columnIDs []uint64, values []int64, bitDepth uint } return nil }(); err != nil { - _ = f.closeStorage() - _ = f.openStorage() + _ = f.closeStorage(true) + _ = f.openStorage(true) return err } @@ -2118,7 +2221,6 @@ func (f *fragment) snapshot() error { // unprotectedWriteToFragment writes the fragment f with bm as the data. It is unprotected, and // f.mu must be locked when calling it. func unprotectedWriteToFragment(f *fragment, bm *roaring.Bitmap) (n int64, err error) { // nolint: interfacer - completeMessage := fmt.Sprintf("fragment: snapshot complete %s/%s/%s/%d", f.index, f.field, f.view, f.shard) start := time.Now() defer track(start, completeMessage, f.stats, f.Logger) @@ -2142,7 +2244,7 @@ func unprotectedWriteToFragment(f *fragment, bm *roaring.Bitmap) (n int64, err e } // Close current storage. - if err := f.closeStorage(); err != nil { + if err := f.closeStorage(false); err != nil { return n, fmt.Errorf("close storage: %s", err) } @@ -2151,8 +2253,12 @@ func unprotectedWriteToFragment(f *fragment, bm *roaring.Bitmap) (n int64, err e return n, fmt.Errorf("rename snapshot: %s", err) } + // if we reloaded from the file, we'd end up with this bitmap + // as our storage. so... let's use this bitmap. as our storage. + f.storage = bm + // Reopen storage. - if err := f.openStorage(); err != nil { + if err := f.openStorage(false); err != nil { return n, fmt.Errorf("open storage: %s", err) } @@ -2341,7 +2447,7 @@ func (f *fragment) readStorageFromArchive(r io.Reader) error { } // Close current storage. - if err := f.closeStorage(); err != nil { + if err := f.closeStorage(true); err != nil { return errors.Wrap(err, "closing") } @@ -2351,7 +2457,7 @@ func (f *fragment) readStorageFromArchive(r io.Reader) error { } // Reopen storage. - if err := f.openStorage(); err != nil { + if err := f.openStorage(true); err != nil { return errors.Wrap(err, "opening") } diff --git a/fragment_internal_test.go b/fragment_internal_test.go index dfa470144..703deba7d 100644 --- a/fragment_internal_test.go +++ b/fragment_internal_test.go @@ -117,7 +117,7 @@ func TestFragment_RowcacheMap(t *testing.T) { _, _ = f.setBit(0, uint64(i*32)) } // force snapshot so we get a mmapped row... - _ = f.snapshot() + _ = f.Snapshot() row := f.row(0) segment := row.Segments()[0] bitmap := segment.data @@ -3227,7 +3227,7 @@ func TestImportClearRestart(t *testing.T) { f2.MaxOpN = maxOpN f2.CacheType = f.CacheType - err = f.closeStorage() + err = f.closeStorage(true) if err != nil { t.Fatalf("closing storage: %v", err) } @@ -3261,7 +3261,7 @@ func TestImportClearRestart(t *testing.T) { f3.MaxOpN = maxOpN f3.CacheType = f.CacheType - err = f2.closeStorage() + err = f2.closeStorage(true) if err != nil { t.Fatalf("f2 closing storage: %v", err) } diff --git a/roaring/roaring.go b/roaring/roaring.go index ab48cf338..6a360ba9b 100644 --- a/roaring/roaring.go +++ b/roaring/roaring.go @@ -1240,6 +1240,82 @@ func (r *roaringIterator) Current() (key uint64, cType byte, n int, length int, return r.currentKey, r.currentType, r.currentN, r.currentLen, r.currentPointer, r.lastErr } +// RemapRoaringStorage tries to update all containers to refer to +// the roaring bitmap in the provided []byte. If any containers are +// marked as mapped, but do not match the provided storage, they will +// be unmapped. The boolean return indicates whether or not any +// containers were mapped to the given storage. +// +// Regardless, after this function runs, no containers have +// mapped storage which does not refer to data; either they got mapped +// to the new storage, or storage was allocated for them. +// +// Data should be in the Pilosa roaring format. +func (b *Bitmap) RemapRoaringStorage(data []byte) (mappedAny bool, returnErr error) { + if b.Containers == nil { + return false, nil + } + var itr *roaringIterator + var err error + var itrKey uint64 + var itrCType byte + var itrN int + var itrPointer *uint16 + var itrErr error + + if data != nil { + itr, err = newRoaringIterator(data) + } + // don't return early: we still have to do the unmapping + if err != nil { + returnErr = err + } + + if itr != nil { + itrKey, itrCType, itrN, _, itrPointer, itrErr = itr.Next() + } + if itrErr != nil { + // iterator errored out, so we won't check it in the loop below + itr = nil + } + + b.Containers.UpdateEvery(func(key uint64, oldC *Container, existed bool) (newC *Container, write bool) { + if itr != nil { + for itrKey < key && itrErr == nil { + itrKey, itrCType, itrN, _, itrPointer, itrErr = itr.Next() + } + if itrErr != nil { + itr = nil + } + // container might be similar enough that we should trust it: + if itrKey == key && itrCType == oldC.typ() && itrN == int(oldC.N()) { + if oldC.frozen() { + // we don't use Clone, because that would copy the + // storage, and we don't need that. + var halfCopy Container + halfCopy = *oldC + halfCopy.flags &^= flagFrozen + newC = &halfCopy + } else { + newC = oldC + } + mappedAny = true + newC.pointer = itrPointer + newC.flags |= flagMapped + return newC, true + } + } + // if the container isn't mapped, we don't need to do anything + if !oldC.Mapped() { + return oldC, false + } + // forcibly unmap it, so the old mapping can be unmapped safely. + newC = oldC.unmapOrClone() + return newC, true + }) + return mappedAny, returnErr +} + // ImportRoaringBits sets-or-clears bits based on a provided Roaring bitmap. // This should be equivalent to unmarshalling the bitmap, then executing // either `b = Union(b, newB)` or `b = Difference(b, newB)`, but with lower diff --git a/view.go b/view.go index 85c9ebbba..620bbcee9 100644 --- a/view.go +++ b/view.go @@ -441,11 +441,11 @@ func upgradeViewBSIv2(v *view, bitDepth uint) (ok bool, _ error) { if tmpPath, err := upgradeRoaringBSIv2(frag, bitDepth); err != nil { return ok, errors.Wrap(err, "upgrading bsi v2") - } else if err := frag.closeStorage(); err != nil { + } else if err := frag.closeStorage(true); err != nil { return ok, errors.Wrap(err, "closing after bsi v2 upgrade") } else if err := os.Rename(tmpPath, frag.path); err != nil { return ok, errors.Wrap(err, "renaming after bsi v2 upgrade") - } else if err := frag.openStorage(); err != nil { + } else if err := frag.openStorage(true); err != nil { return ok, errors.Wrap(err, "re-opening after bsi v2 upgrade") } }