Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
5 changes: 4 additions & 1 deletion .seqbench/continuous.env
Original file line number Diff line number Diff line change
@@ -1,4 +1,7 @@
GOGC=100
GOGC=200
# GOGC=500
# GOMEMLIMIT=4GiB
GO_VERSION=1.26

SEQDB_RESOURCES_SKIP_FSYNC=true

Expand Down
33 changes: 9 additions & 24 deletions cache/cache.go
Original file line number Diff line number Diff line change
Expand Up @@ -253,45 +253,30 @@ func (c *Cache[V]) handlePanic(key uint32, wg *sync.WaitGroup) {
panic(err)
}

func (c *Cache[V]) Get(key uint32, fn func() (V, int)) V {
// attempt to obtain cached value
// or create an entry for a new one
e, wg, success := c.getOrCreate(key)
if success {
return e.value
}

defer c.handlePanic(key, wg)
// long operation
t := time.Now()
value, refMemSize := fn()
latency := time.Since(t).Seconds()

// all good, just update the cache
c.save(e, wg, value, refMemSize, latency)

return value
func (c *Cache[V]) Get(key uint32, l Loader[V]) (V, error) {
value, _, err := c.get(key, l)
return value, err
}

func (c *Cache[V]) GetWithError(key uint32, fn func() (V, int, error)) (V, error) {
func (c *Cache[V]) get(key uint32, l Loader[V]) (V, *entry[V], error) {
e, wg, success := c.getOrCreate(key)
if success {
return e.value, nil
return e.value, e, nil
}

defer c.handlePanic(key, wg)

t := time.Now()
value, refMemSize, err := fn()
value, refMemSize, err := l.Load(key)
latency := time.Since(t).Seconds()

if err != nil {
c.recover(key, wg)
return value, err
return value, nil, err
}

c.save(e, wg, value, refMemSize, latency)

return value, nil
return value, e, nil
}

func (c *Cache[V]) Release() {
Expand Down
18 changes: 9 additions & 9 deletions cache/cache_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ func TestCacheSize(t *testing.T) {
cleaner := NewCleaner(0, nil)
c := NewCache[[]byte](cleaner, nil)

c.Get(1, func() ([]byte, int) { return make([]byte, SIZE), SIZE })
_, _ = c.Get(1, LoaderFunc[[]byte](func(uint32) ([]byte, int, error) { return make([]byte, SIZE), SIZE, nil }))

total := cleaner.getSize()
assert.Equal(t, uint64(SIZE)+c.entrySize, total, "wrong cache size")
Expand All @@ -36,14 +36,14 @@ func TestClean(t *testing.T) {
c3 := NewCache[[]byte](cleaner, nil)

stat := &CleanStat{}
c1.Get(0, func() ([]byte, int) { return make([]byte, Size1), int(Size1) })
_, _ = c1.Get(0, LoaderFunc[[]byte](func(uint32) ([]byte, int, error) { return make([]byte, Size1), int(Size1), nil }))

cleaner.Rotate()
cleaner.Cleanup(stat)

c1.Get(1, func() ([]byte, int) { return make([]byte, Size2), int(Size2) })
c2.Get(1, func() ([]byte, int) { return make([]byte, Size3), int(Size3) })
c3.Get(1, func() ([]byte, int) { return make([]byte, Size4), int(Size4) })
_, _ = c1.Get(1, LoaderFunc[[]byte](func(uint32) ([]byte, int, error) { return make([]byte, Size2), int(Size2), nil }))
_, _ = c2.Get(1, LoaderFunc[[]byte](func(uint32) ([]byte, int, error) { return make([]byte, Size3), int(Size3), nil }))
_, _ = c3.Get(1, LoaderFunc[[]byte](func(uint32) ([]byte, int, error) { return make([]byte, Size4), int(Size4), nil }))

bytesTotal := cleaner.getSize()

Expand Down Expand Up @@ -111,14 +111,14 @@ func TestStress(t *testing.T) {
}
}()
}
val := c.Get(key, func() ([]uint64, int) {
val, _ := c.Get(key, LoaderFunc[[]uint64](func(uint32) ([]uint64, int, error) {
time.Sleep(1 * time.Millisecond)
if err != nil {
panicFired = true
panic(err)
}
return []uint64{uint64(key)}, 32
})
return []uint64{uint64(key)}, 32, nil
}))
if val == nil {
t.Errorf("cache is corrupted")
}
Expand All @@ -136,7 +136,7 @@ func BenchmarkBucketClean(b *testing.B) {
b.StopTimer()

for i := range 1000 {
c.Get(uint32(i), func() (int, int) { return i, 4 })
_, _ = c.Get(uint32(i), LoaderFunc[int](func(uint32) (int, int, error) { return i, 4, nil }))
}

cleaner.markStale(cleaner.getSize())
Expand Down
92 changes: 92 additions & 0 deletions cache/wrapper.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,92 @@
package cache

import (
"weak"
)

type Loader[V any] interface {
Load(key uint32) (value V, size int, err error)
}

type LoaderFunc[V any] func(key uint32) (V, int, error)

func (f LoaderFunc[V]) Load(key uint32) (V, int, error) {
return f(key)
}

type Wrapper[V any] interface {
Get(key uint32, l Loader[V]) (V, error)
}

var (
_ Wrapper[any] = (*Cache[any])(nil)
_ Wrapper[any] = (*Session[any])(nil)
_ Wrapper[any] = (*Scan[any])(nil)
)

// Session is a wrapper for [Cache] which stores local
// copies of [entry] via weak pointers.
type Session[V any] struct {
cache *Cache[V]
local map[uint32]weak.Pointer[entry[V]]
}

func NewSession[V any](c *Cache[V]) *Session[V] {
return &Session[V]{
cache: c,
local: make(map[uint32]weak.Pointer[entry[V]]),
}
}

func (s *Session[V]) Get(key uint32, l Loader[V]) (V, error) {
if e := s.resolve(key); e != nil {
return e.value, nil
}

value, e, err := s.cache.get(key, l)
if err != nil {
return value, err
}

s.local[key] = weak.Make(e)
return value, nil
}

func (s *Session[V]) resolve(key uint32) *entry[V] {
wp, ok := s.local[key]
if !ok {
return nil
}

e := wp.Value()
if e == nil {
delete(s.local, key)
return nil
}

return e
}

type Scan[V any] struct {
key uint32
value V
ok bool
}

func NewScan[V any]() *Scan[V] {
return &Scan[V]{}
}

func (s *Scan[V]) Get(key uint32, l Loader[V]) (V, error) {
if s.ok && s.key == key {
return s.value, nil
}

value, _, err := l.Load(key)
if err != nil {
return value, err
}

s.key, s.value, s.ok = key, value, true
return value, nil
}
7 changes: 1 addition & 6 deletions cmd/distribution/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -48,12 +48,7 @@ func getReader(path string) (storage.IndexReader, *os.File) {
if err != nil {
panic(err)
}
readerProvider := storage.NewReaderProvider(readLimiter, f.Name(), f, c)
reader, err := readerProvider.GetReader()
if err != nil {
panic(err)
}
return reader, f
return storage.NewIndexReader(readLimiter, f.Name(), f, c), f
}

func readBlock(reader storage.IndexReader, blockIndex uint32) ([]byte, error) {
Expand Down
77 changes: 31 additions & 46 deletions frac/remote.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,9 +44,9 @@ type Remote struct {
docsReader storage.DocsReader

// IsLegacy is true for fractions that use the old single .index file format.
IsLegacy bool
legacyFile storage.ImmutableFile
legacyReaderProvider *storage.ReaderProvider
IsLegacy bool
legacyFile storage.ImmutableFile
legacyReader storage.IndexReader

// Per-section index files and their readers (new split format only).
infoFile storage.ImmutableFile
Expand All @@ -55,10 +55,10 @@ type Remote struct {
idFile storage.ImmutableFile
lidFile storage.ImmutableFile

tokenReaderProvider *storage.ReaderProvider
offsetsReaderProvider *storage.ReaderProvider
idReaderProvider *storage.ReaderProvider
lidReaderProvider *storage.ReaderProvider
tokenReader storage.IndexReader
offsetsReader storage.IndexReader
idReader storage.IndexReader
lidReader storage.IndexReader

indexCache *IndexCache

Expand Down Expand Up @@ -161,14 +161,6 @@ func (f *Remote) FindLIDs(ctx context.Context, ids []seq.ID) ([]seq.LID, error)
return dp.FindLIDs(ids)
}

func (f *Remote) mustGetReader(p *storage.ReaderProvider) storage.IndexReader {
r, err := p.GetReader()
if err != nil {
logger.Fatal("error creating IndexReader", zap.String("fraction", f.BaseFileName), zap.Error(err))
}
return r
}

func (f *Remote) createDataProvider(ctx context.Context) (*sealedDataProvider, error) {
if err := f.init(); err != nil {
logger.Error(
Expand All @@ -179,21 +171,14 @@ func (f *Remote) createDataProvider(ctx context.Context) (*sealedDataProvider, e
return nil, err
}

var (
tokenReader storage.IndexReader
lidReader storage.IndexReader
idReader storage.IndexReader
)
tokenReader := &f.tokenReader
lidReader := &f.lidReader
idReader := &f.idReader

if f.IsLegacy {
legacyReader := f.mustGetReader(f.legacyReaderProvider)
tokenReader = legacyReader
lidReader = legacyReader
idReader = legacyReader
} else {
tokenReader = f.mustGetReader(f.tokenReaderProvider)
lidReader = f.mustGetReader(f.lidReaderProvider)
idReader = f.mustGetReader(f.idReaderProvider)
tokenReader = &f.legacyReader
lidReader = &f.legacyReader
idReader = &f.legacyReader
}

return &sealedDataProvider{
Expand All @@ -205,16 +190,16 @@ func (f *Remote) createDataProvider(ctx context.Context) (*sealedDataProvider, e
docsReader: &f.docsReader,
blocksOffsets: f.blocksData.BlocksOffsets,
lidsTable: f.blocksData.LIDsTable,
lidsLoader: lids.NewLoader(f.info.BinaryDataVer, &lidReader, f.indexCache.LIDs),
tokenBlockLoader: token.NewBlockLoader(f.BaseFileName, f.Info().BinaryDataVer, &tokenReader, f.indexCache.Tokens),
tokenTableLoader: token.NewTableLoader(f.BaseFileName, f.Info().BinaryDataVer, f.IsLegacy, &tokenReader, f.indexCache.TokenTable),
lidsLoader: lids.NewLoader(f.info.BinaryDataVer, lidReader, cache.NewSession(f.indexCache.LIDs)),
tokenBlockLoader: token.NewBlockLoader(f.BaseFileName, f.Info().BinaryDataVer, tokenReader, cache.NewSession(f.indexCache.Tokens)),
tokenTableLoader: token.NewTableLoader(f.BaseFileName, f.Info().BinaryDataVer, f.IsLegacy, tokenReader, cache.NewSession(f.indexCache.TokenTable)),

idsTable: &f.blocksData.IDsTable,
idsProvider: seqids.NewProvider(
&idReader,
f.indexCache.MIDs,
f.indexCache.RIDs,
f.indexCache.Params,
idReader,
cache.NewSession(f.indexCache.MIDs),
cache.NewSession(f.indexCache.RIDs),
cache.NewSession(f.indexCache.Params),
&f.blocksData.IDsTable,
f.info.BinaryDataVer,
),
Expand Down Expand Up @@ -275,7 +260,7 @@ func (f *Remote) loadInfo() error {
if err := f.openInfoLegacy(); err != nil {
return err
}
if f.info, err = loadInfoLegacy(f.mustGetReader(f.legacyReaderProvider)); err != nil {
if f.info, err = loadInfoLegacy(f.legacyReader); err != nil {
logger.Fatal("error loading Info", zap.String("fraction", f.BaseFileName), zap.Error(err))
}
return nil
Expand Down Expand Up @@ -308,16 +293,16 @@ func (f *Remote) init() error {
}

if f.IsLegacy {
(&LegacyLoader{}).Load(&f.blocksData, f.info, f.mustGetReader(f.legacyReaderProvider))
(&LegacyLoader{}).Load(&f.blocksData, f.info, f.legacyReader)
f.isInited = true
return nil
}

(&Loader{}).Load(&f.blocksData, f.info, IndexReaders{
Token: f.mustGetReader(f.tokenReaderProvider),
Offsets: f.mustGetReader(f.offsetsReaderProvider),
ID: f.mustGetReader(f.idReaderProvider),
LID: f.mustGetReader(f.lidReaderProvider),
Token: f.tokenReader,
Offsets: f.offsetsReader,
ID: f.idReader,
LID: f.lidReader,
})

f.isInited = true
Expand All @@ -331,7 +316,7 @@ func (f *Remote) openInfoLegacy() error {

return f.openRemoteFile(consts.IndexFileSuffix, func(file storage.ImmutableFile) {
f.legacyFile = file
f.legacyReaderProvider = storage.NewReaderProvider(
f.legacyReader = storage.NewIndexReader(
f.readLimiter, file.Name(),
file, f.indexCache.LegacyRegistry,
)
Expand Down Expand Up @@ -365,7 +350,7 @@ func (f *Remote) openIndex() error {
consts.TokenFileSuffix,
func(file storage.ImmutableFile) {
f.tokenFile = file
f.tokenReaderProvider = storage.NewReaderProvider(
f.tokenReader = storage.NewIndexReader(
f.readLimiter, file.Name(),
file, f.indexCache.TokenRegistry,
)
Expand All @@ -380,7 +365,7 @@ func (f *Remote) openIndex() error {
consts.OffsetsFileSuffix,
func(file storage.ImmutableFile) {
f.offsetsFile = file
f.offsetsReaderProvider = storage.NewReaderProvider(
f.offsetsReader = storage.NewIndexReader(
f.readLimiter, file.Name(),
file, f.indexCache.OffsetsRegistry,
)
Expand All @@ -395,7 +380,7 @@ func (f *Remote) openIndex() error {
consts.IDFileSuffix,
func(file storage.ImmutableFile) {
f.idFile = file
f.idReaderProvider = storage.NewReaderProvider(
f.idReader = storage.NewIndexReader(
f.readLimiter, file.Name(),
file, f.indexCache.IDRegistry,
)
Expand All @@ -410,7 +395,7 @@ func (f *Remote) openIndex() error {
consts.LIDFileSuffix,
func(file storage.ImmutableFile) {
f.lidFile = file
f.lidReaderProvider = storage.NewReaderProvider(
f.lidReader = storage.NewIndexReader(
f.readLimiter, file.Name(),
file, f.indexCache.LIDRegistry,
)
Expand Down
Loading
Loading