From cfdd89d5887ac69c487fe0854026a2a47cf0ad14 Mon Sep 17 00:00:00 2001 From: Daniil Porokhnin Date: Tue, 11 Aug 2026 14:17:00 +0300 Subject: [PATCH 1/6] feat: `Session` and `Scan` wrappers for `cache` --- cache/cache.go | 20 ++++-- cache/wrapper.go | 105 ++++++++++++++++++++++++++++++ cmd/distribution/main.go | 7 +- frac/remote.go | 77 +++++++++------------- frac/sealed.go | 77 +++++++++------------- frac/sealed/lids/loader.go | 4 +- frac/sealed/seqids/loader.go | 6 +- frac/sealed/seqids/provider.go | 6 +- frac/sealed/token/block_loader.go | 4 +- frac/sealed/token/table_loader.go | 9 ++- frac/sealed_loader.go | 16 ++++- frac/sealed_source.go | 58 ++++++++++++----- storage/index_reader.go | 91 ++++++++++---------------- 13 files changed, 289 insertions(+), 191 deletions(-) create mode 100644 cache/wrapper.go diff --git a/cache/cache.go b/cache/cache.go index 2463bd5b1..f85fe6752 100644 --- a/cache/cache.go +++ b/cache/cache.go @@ -254,11 +254,16 @@ func (c *Cache[V]) handlePanic(key uint32, wg *sync.WaitGroup) { } func (c *Cache[V]) Get(key uint32, fn func() (V, int)) V { + value, _ := c.get(key, fn) + return value +} + +func (c *Cache[V]) get(key uint32, fn func() (V, int)) (V, *entry[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 + return e.value, e } defer c.handlePanic(key, wg) @@ -270,13 +275,18 @@ func (c *Cache[V]) Get(key uint32, fn func() (V, int)) V { // all good, just update the cache c.save(e, wg, value, refMemSize, latency) - return value + return value, e } func (c *Cache[V]) GetWithError(key uint32, fn func() (V, int, error)) (V, error) { + value, _, err := c.getWithError(key, fn) + return value, err +} + +func (c *Cache[V]) getWithError(key uint32, fn func() (V, int, error)) (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) @@ -286,12 +296,12 @@ func (c *Cache[V]) GetWithError(key uint32, fn func() (V, int, error)) (V, error 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() { diff --git a/cache/wrapper.go b/cache/wrapper.go new file mode 100644 index 000000000..10f31cd99 --- /dev/null +++ b/cache/wrapper.go @@ -0,0 +1,105 @@ +package cache + +import ( + "weak" +) + +type Wrapper[V any] interface { + Get(key uint32, fn func() (V, int)) V + GetWithError(key uint32, fn func() (V, int, error)) (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, fn func() (V, int)) V { + if e := s.resolve(key); e != nil { + return e.value + } + + value, e := s.cache.get(key, fn) + s.local[key] = weak.Make(e) + + return value +} + +func (s *Session[V]) GetWithError(key uint32, fn func() (V, int, error)) (V, error) { + if e := s.resolve(key); e != nil { + return e.value, nil + } + + value, e, err := s.cache.getWithError(key, fn) + 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, fn func() (V, int)) V { + if s.ok && s.key == key { + return s.value + } + + value, _ := fn() + s.key, s.value, s.ok = key, value, true + + return value +} + +func (s *Scan[V]) GetWithError(key uint32, fn func() (V, int, error)) (V, error) { + if s.ok && s.key == key { + return s.value, nil + } + + value, _, err := fn() + if err != nil { + return value, err + } + + s.key, s.value, s.ok = key, value, true + return value, nil +} diff --git a/cmd/distribution/main.go b/cmd/distribution/main.go index 371e7bba7..527c484ff 100644 --- a/cmd/distribution/main.go +++ b/cmd/distribution/main.go @@ -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) { diff --git a/frac/remote.go b/frac/remote.go index d118d41cb..f95a74f34 100644 --- a/frac/remote.go +++ b/frac/remote.go @@ -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 @@ -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 @@ -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( @@ -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{ @@ -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, ), @@ -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 @@ -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 @@ -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, ) @@ -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, ) @@ -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, ) @@ -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, ) @@ -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, ) diff --git a/frac/sealed.go b/frac/sealed.go index 28754be91..4ffd79424 100644 --- a/frac/sealed.go +++ b/frac/sealed.go @@ -40,9 +40,9 @@ type Sealed struct { docsReader storage.DocsReader // IsLegacy is true for fractions that use the old single .index file format. - IsLegacy bool - legacyFile *os.File - legacyReaderProvider *storage.ReaderProvider + IsLegacy bool + legacyFile *os.File + legacyReader storage.IndexReader // Per-section index files and their readers (new split format only). infoFile *os.File @@ -51,10 +51,10 @@ type Sealed struct { idFile *os.File lidFile *os.File - tokenReaderProvider *storage.ReaderProvider - offsetsReaderProvider *storage.ReaderProvider - idReaderProvider *storage.ReaderProvider - lidReaderProvider *storage.ReaderProvider + tokenReader storage.IndexReader + offsetsReader storage.IndexReader + idReader storage.IndexReader + lidReader storage.IndexReader blocksData sealed.BlocksData indexCache *IndexCache @@ -176,7 +176,7 @@ func (f *Sealed) openInfoLegacy() { } f.legacyFile = file - f.legacyReaderProvider = storage.NewReaderProvider( + f.legacyReader = storage.NewIndexReader( f.readLimiter, file.Name(), file, f.indexCache.LegacyRegistry, ) @@ -216,7 +216,7 @@ func (f *Sealed) openIndex() { logger.Fatal("can't open token file", zap.String("file", name), zap.Error(err)) } f.tokenFile = file - f.tokenReaderProvider = storage.NewReaderProvider(f.readLimiter, file.Name(), file, f.indexCache.TokenRegistry) + f.tokenReader = storage.NewIndexReader(f.readLimiter, file.Name(), file, f.indexCache.TokenRegistry) } if f.offsetsFile == nil { @@ -226,7 +226,7 @@ func (f *Sealed) openIndex() { logger.Fatal("can't open offsets file", zap.String("file", name), zap.Error(err)) } f.offsetsFile = file - f.offsetsReaderProvider = storage.NewReaderProvider(f.readLimiter, file.Name(), file, f.indexCache.OffsetsRegistry) + f.offsetsReader = storage.NewIndexReader(f.readLimiter, file.Name(), file, f.indexCache.OffsetsRegistry) } if f.idFile == nil { @@ -236,7 +236,7 @@ func (f *Sealed) openIndex() { logger.Fatal("can't open id file", zap.String("file", name), zap.Error(err)) } f.idFile = file - f.idReaderProvider = storage.NewReaderProvider(f.readLimiter, file.Name(), file, f.indexCache.IDRegistry) + f.idReader = storage.NewIndexReader(f.readLimiter, file.Name(), file, f.indexCache.IDRegistry) } if f.lidFile == nil { @@ -246,7 +246,7 @@ func (f *Sealed) openIndex() { logger.Fatal("can't open lid file", zap.String("file", name), zap.Error(err)) } f.lidFile = file - f.lidReaderProvider = storage.NewReaderProvider(f.readLimiter, file.Name(), file, f.indexCache.LIDRegistry) + f.lidReader = storage.NewIndexReader(f.readLimiter, file.Name(), file, f.indexCache.LIDRegistry) } } @@ -279,20 +279,12 @@ func (f *Sealed) openDocs() { f.docsReader = storage.NewDocsReader(f.readLimiter, f.docsFile, f.docsCache) } -func (f *Sealed) 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 *Sealed) loadInfo() { var err error if f.IsLegacy { f.openInfoLegacy() - 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 @@ -316,16 +308,16 @@ func (f *Sealed) init(full bool) { } 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 } (&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 @@ -503,21 +495,14 @@ func (f *Sealed) FindLIDs(ctx context.Context, ids []seq.ID) ([]seq.LID, error) func (f *Sealed) createDataProvider(ctx context.Context) *sealedDataProvider { f.init(true) - 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{ @@ -529,16 +514,16 @@ func (f *Sealed) createDataProvider(ctx context.Context) *sealedDataProvider { 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, ), diff --git a/frac/sealed/lids/loader.go b/frac/sealed/lids/loader.go index cf987a979..2992355a8 100644 --- a/frac/sealed/lids/loader.go +++ b/frac/sealed/lids/loader.go @@ -45,14 +45,14 @@ func (b *UnpackBuffer) Reset(fracVer config.BinaryDataVersion) { // NOT THREAD SAFE. Do not use concurrently. // Use your own Loader instance for each search query type Loader struct { - cache *cache.Cache[*Block] + cache cache.Wrapper[*Block] reader *storage.IndexReader unpackBuf *UnpackBuffer blockBuf []byte fracVer config.BinaryDataVersion } -func NewLoader(fracVer config.BinaryDataVersion, r *storage.IndexReader, c *cache.Cache[*Block]) *Loader { +func NewLoader(fracVer config.BinaryDataVersion, r *storage.IndexReader, c cache.Wrapper[*Block]) *Loader { return &Loader{ cache: c, reader: r, diff --git a/frac/sealed/seqids/loader.go b/frac/sealed/seqids/loader.go index 5f41c3a52..3d92bede8 100644 --- a/frac/sealed/seqids/loader.go +++ b/frac/sealed/seqids/loader.go @@ -28,9 +28,9 @@ func (Table) BlockStartLID(blockIndex uint32) uint32 { type Loader struct { reader *storage.IndexReader table *Table - cacheMIDs *cache.Cache[[]byte] - cacheRIDs *cache.Cache[BlockRIDs] - cacheParams *cache.Cache[BlockParams] + cacheMIDs cache.Wrapper[[]byte] + cacheRIDs cache.Wrapper[BlockRIDs] + cacheParams cache.Wrapper[BlockParams] fracVersion config.BinaryDataVersion } diff --git a/frac/sealed/seqids/provider.go b/frac/sealed/seqids/provider.go index a2de6593f..d1c5f8cff 100644 --- a/frac/sealed/seqids/provider.go +++ b/frac/sealed/seqids/provider.go @@ -18,9 +18,9 @@ type Provider struct { func NewProvider( indexReader *storage.IndexReader, - cacheMIDs *cache.Cache[[]byte], - cacheRIDs *cache.Cache[BlockRIDs], - cacheParams *cache.Cache[BlockParams], + cacheMIDs cache.Wrapper[[]byte], + cacheRIDs cache.Wrapper[BlockRIDs], + cacheParams cache.Wrapper[BlockParams], table *Table, fracVersion config.BinaryDataVersion, ) *Provider { diff --git a/frac/sealed/token/block_loader.go b/frac/sealed/token/block_loader.go index 266315725..81eb46d9e 100644 --- a/frac/sealed/token/block_loader.go +++ b/frac/sealed/token/block_loader.go @@ -191,7 +191,7 @@ func (b *Block) find(from, to int, searcher pattern.Searcher) ([]int, error) { type BlockLoader struct { fracName string fracVer config.BinaryDataVersion - cache *cache.Cache[*Block] + cache cache.Wrapper[*Block] reader *storage.IndexReader unpackBuf *UnpackBuffer blockBuf []byte @@ -201,7 +201,7 @@ func NewBlockLoader( fracName string, fracVer config.BinaryDataVersion, reader *storage.IndexReader, - c *cache.Cache[*Block], + c cache.Wrapper[*Block], ) *BlockLoader { return &BlockLoader{ fracName: fracName, diff --git a/frac/sealed/token/table_loader.go b/frac/sealed/token/table_loader.go index f4bf56c76..82d3840ea 100644 --- a/frac/sealed/token/table_loader.go +++ b/frac/sealed/token/table_loader.go @@ -24,7 +24,7 @@ type TableLoader struct { isLegacy bool reader *storage.IndexReader - cache *cache.Cache[Table] + cache cache.Wrapper[Table] once sync.Once tableIndex uint32 @@ -37,7 +37,7 @@ func NewTableLoader( fracVer config.BinaryDataVersion, isLegacy bool, reader *storage.IndexReader, - c *cache.Cache[Table], + c cache.Wrapper[Table], ) *TableLoader { return &TableLoader{ fracName: fracName, @@ -141,7 +141,10 @@ func (l *TableLoader) loadBlocksLegacy() ([]TableBlock, error) { } func (l *TableLoader) loadBlocks() ([]TableBlock, error) { - blocksCount := l.reader.BlocksCount() + blocksCount, err := l.reader.BlocksCount() + if err != nil { + return nil, err + } var blocks []TableBlock for blockIndex := l.tableIndex; blockIndex < uint32(blocksCount); blockIndex++ { diff --git a/frac/sealed_loader.go b/frac/sealed_loader.go index 17782072a..10d95a2d3 100644 --- a/frac/sealed_loader.go +++ b/frac/sealed_loader.go @@ -231,7 +231,13 @@ func (l *Loader) loadIDsTable(r storage.IndexReader, info *common.Info) seqids.T IDsTotal: info.DocsTotal + 1, // Increment by one for [seq.SystemID] } - blocksCount := r.BlocksCount() + blocksCount, err := r.BlocksCount() + if err != nil { + logger.Fatal( + "cannot get block count", + zap.Error(err), + ) + } for blockIdx := 0; blockIdx < blocksCount; blockIdx += 3 { header, err := r.GetBlockHeader(uint32(blockIdx)) @@ -263,7 +269,13 @@ func (l *Loader) loadLIDsTable(r storage.IndexReader) (*lids.Table, error) { isContinued []bool ) - blocksCount := r.BlocksCount() + blocksCount, err := r.BlocksCount() + if err != nil { + logger.Fatal( + "cannot get block count", + zap.Error(err), + ) + } for blockIdx := 0; blockIdx < blocksCount; blockIdx++ { header, err := r.GetBlockHeader(uint32(blockIdx)) diff --git a/frac/sealed_source.go b/frac/sealed_source.go index 31f7e9f2e..dfab4a567 100644 --- a/frac/sealed_source.go +++ b/frac/sealed_source.go @@ -4,6 +4,7 @@ import ( "iter" "slices" + "github.com/ozontech/seq-db/cache" "github.com/ozontech/seq-db/frac/common" "github.com/ozontech/seq-db/frac/sealed/lids" "github.com/ozontech/seq-db/frac/sealed/seqids" @@ -21,6 +22,10 @@ type DocBlockLocation = util.Pair[[]byte, uint64] type SealedSource struct { f *Sealed + tokenReader storage.IndexReader + idReader storage.IndexReader + lidReader storage.IndexReader + idsProvider *seqids.Provider lidsLoader *lids.Loader @@ -31,24 +36,43 @@ type SealedSource struct { func NewSealedSource(f *Sealed) *SealedSource { f.init(true) - idReader := f.mustGetReader(f.idReaderProvider) - lidReader := f.mustGetReader(f.lidReaderProvider) - tokenReader := f.mustGetReader(f.tokenReaderProvider) - - return &SealedSource{ - f: f, - idsProvider: seqids.NewProvider( - &idReader, - f.indexCache.MIDs, - f.indexCache.RIDs, - f.indexCache.Params, - &f.blocksData.IDsTable, - f.info.BinaryDataVer, - ), - 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), + tokenFile, idFile, lidFile := f.tokenFile, f.idFile, f.lidFile + if f.IsLegacy { + tokenFile, idFile, lidFile = f.legacyFile, f.legacyFile, f.legacyFile + } + + s := &SealedSource{ + f: f, + tokenReader: storage.NewIndexReader(f.readLimiter, tokenFile.Name(), tokenFile, cache.NewScan[[]byte]()), + idReader: storage.NewIndexReader(f.readLimiter, idFile.Name(), idFile, cache.NewScan[[]byte]()), + lidReader: storage.NewIndexReader(f.readLimiter, lidFile.Name(), lidFile, cache.NewScan[[]byte]()), } + + s.idsProvider = seqids.NewProvider( + &s.idReader, + cache.NewScan[[]byte](), + cache.NewScan[seqids.BlockRIDs](), + cache.NewScan[seqids.BlockParams](), + &f.blocksData.IDsTable, + f.info.BinaryDataVer, + ) + + s.lidsLoader = lids.NewLoader( + f.Info().BinaryDataVer, + &s.lidReader, cache.NewScan[*lids.Block](), + ) + + s.tokenBlockLoader = token.NewBlockLoader( + f.BaseFileName, f.Info().BinaryDataVer, + &s.tokenReader, cache.NewScan[*token.Block](), + ) + + s.tokenTableLoader = token.NewTableLoader( + f.BaseFileName, f.Info().BinaryDataVer, + f.IsLegacy, &s.tokenReader, cache.NewScan[token.Table](), + ) + + return s } func (s *SealedSource) Info() *common.Info { diff --git a/storage/index_reader.go b/storage/index_reader.go index 2cbe71fee..a07e08513 100644 --- a/storage/index_reader.go +++ b/storage/index_reader.go @@ -10,50 +10,56 @@ import ( "github.com/ozontech/seq-db/util" ) -type ReaderProvider struct { +type IndexReader struct { limiter *ReadLimiter reader io.ReaderAt readerName string - cache *cache.Cache[[]byte] + cache cache.Wrapper[[]byte] } -func NewReaderProvider( - limiter *ReadLimiter, - readerName string, - reader io.ReaderAt, - registryCache *cache.Cache[[]byte], -) *ReaderProvider { - return &ReaderProvider{ +func NewIndexReader( + limiter *ReadLimiter, readerName string, + reader io.ReaderAt, registryCache cache.Wrapper[[]byte], +) IndexReader { + return IndexReader{ limiter: limiter, - readerName: readerName, reader: reader, + readerName: readerName, cache: registryCache, } } -func (r *ReaderProvider) GetReader() (IndexReader, error) { +func (r *IndexReader) GetBlockHeader(index uint32) (IndexBlockHeader, error) { registry, err := r.cache.GetWithError(1, func() ([]byte, int, error) { data, err := r.readRegistry() return data, cap(data), err }) if err != nil { - return IndexReader{}, err + return nil, err + } + + if (uint64(index)+1)*IndexBlockHeaderSize > uint64(len(registry)) { + return nil, fmt.Errorf( + "too large index block in file %s, with index %d, registry size %d", + r.readerName, index, len(registry), + ) } - return NewIndexReader(r.limiter, r.readerName, r.reader, registry), nil + pos := index * IndexBlockHeaderSize + return registry[pos : pos+IndexBlockHeaderSize], nil } -func (r *ReaderProvider) readRegistry() ([]byte, error) { +func (r *IndexReader) readRegistry() ([]byte, error) { numBuf := make([]byte, 16) n, err := r.limiter.ReadAt(r.reader, numBuf, 0) if err != nil { - return nil, fmt.Errorf("can't read disk registry from file %s: %s", r.readerName, err.Error()) + return nil, fmt.Errorf("can't read disk registry, %s", err.Error()) } if n == 0 { - return nil, fmt.Errorf("can't read disk registry from file %s, n=0", r.readerName) + return nil, fmt.Errorf("can't read disk registry, n=0") } pos := binary.LittleEndian.Uint64(numBuf) @@ -62,55 +68,20 @@ func (r *ReaderProvider) readRegistry() ([]byte, error) { n, err = r.limiter.ReadAt(r.reader, buf, int64(pos)) if err != nil && err != io.EOF { - return nil, fmt.Errorf("can't read disk registry from file %s: %s", r.readerName, err.Error()) + return nil, fmt.Errorf("can't read disk registry, %s", err.Error()) } if uint64(n) != l { - return nil, fmt.Errorf("can't read disk registry, read=%d, requested=%d in file %s", n, l, r.readerName) + return nil, fmt.Errorf("can't read disk registry, read=%d, requested=%d", n, l) } if len(buf)%IndexBlockHeaderSize != 0 { - return nil, fmt.Errorf("wrong registry format in file %s", r.readerName) + return nil, fmt.Errorf("wrong registry format") } return buf, nil } -type IndexReader struct { - limiter *ReadLimiter - - reader io.ReaderAt - readerName string - - registry []byte -} - -func NewIndexReader( - limiter *ReadLimiter, - readerName string, - reader io.ReaderAt, - registry []byte, -) IndexReader { - return IndexReader{ - limiter: limiter, - reader: reader, - readerName: readerName, - registry: registry, - } -} - -func (r *IndexReader) GetBlockHeader(index uint32) (IndexBlockHeader, error) { - if (uint64(index)+1)*IndexBlockHeaderSize > uint64(len(r.registry)) { - return nil, fmt.Errorf( - "too large index block in file %s, with index %d, registry size %d", - r.readerName, index, len(r.registry), - ) - } - - pos := index * IndexBlockHeaderSize - return r.registry[pos : pos+IndexBlockHeaderSize], nil -} - func (r *IndexReader) ReadIndexBlock(blockIndex uint32, dst []byte) ([]byte, uint64, error) { header, err := r.GetBlockHeader(blockIndex) if err != nil { @@ -137,6 +108,14 @@ func (r *IndexReader) ReadIndexBlock(blockIndex uint32, dst []byte) ([]byte, uin return dst, uint64(n), err } -func (r *IndexReader) BlocksCount() int { - return len(r.registry) / IndexBlockHeaderSize +func (r *IndexReader) BlocksCount() (int, error) { + registry, err := r.cache.GetWithError(1, func() ([]byte, int, error) { + data, err := r.readRegistry() + return data, cap(data), err + }) + if err != nil { + return 0, err + } + + return len(registry) / IndexBlockHeaderSize, nil } From f2345376b7103a3c6819a7404d79ae1dd1b03383 Mon Sep 17 00:00:00 2001 From: Daniil Porokhnin Date: Tue, 11 Aug 2026 18:56:52 +0300 Subject: [PATCH 2/6] chore: tune GOGC/GOMEMLIMIT --- .seqbench/continuous.env | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/.seqbench/continuous.env b/.seqbench/continuous.env index 1d5c154e7..5589bd189 100644 --- a/.seqbench/continuous.env +++ b/.seqbench/continuous.env @@ -1,4 +1,5 @@ -GOGC=100 +GOGC=500 +GOMEMLIMIT=4GiB SEQDB_RESOURCES_SKIP_FSYNC=true From ee3f91bf39bedd8be5560ac3c2d6b1f58a30b478 Mon Sep 17 00:00:00 2001 From: Daniil Porokhnin Date: Fri, 14 Aug 2026 12:30:42 +0300 Subject: [PATCH 3/6] refactor: use `Loader` interface in cache --- cache/cache.go | 35 ++------ cache/cache_test.go | 18 ++-- cache/wrapper.go | 43 ++++------ frac/sealed.go | 6 +- frac/sealed/lids/loader.go | 17 ++-- frac/sealed/seqids/loader.go | 131 +++++++++++++++++------------- frac/sealed/seqids/provider.go | 6 +- frac/sealed/token/block_loader.go | 30 ++++--- frac/sealed/token/table_loader.go | 46 ++++++----- skipmaskmanager/loader.go | 18 ++-- storage/docs_reader.go | 21 +++-- storage/index_reader.go | 78 +++++++++--------- 12 files changed, 227 insertions(+), 222 deletions(-) diff --git a/cache/cache.go b/cache/cache.go index f85fe6752..3ff827716 100644 --- a/cache/cache.go +++ b/cache/cache.go @@ -253,45 +253,21 @@ func (c *Cache[V]) handlePanic(key uint32, wg *sync.WaitGroup) { panic(err) } -func (c *Cache[V]) Get(key uint32, fn func() (V, int)) V { - value, _ := c.get(key, fn) - return value -} - -func (c *Cache[V]) get(key uint32, fn func() (V, int)) (V, *entry[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, e - } - - 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, e -} - -func (c *Cache[V]) GetWithError(key uint32, fn func() (V, int, error)) (V, error) { - value, _, err := c.getWithError(key, fn) +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, *entry[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, 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 { @@ -300,7 +276,6 @@ func (c *Cache[V]) getWithError(key uint32, fn func() (V, int, error)) (V, *entr } c.save(e, wg, value, refMemSize, latency) - return value, e, nil } diff --git a/cache/cache_test.go b/cache/cache_test.go index 9c9d6730f..475f14920 100644 --- a/cache/cache_test.go +++ b/cache/cache_test.go @@ -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") @@ -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() @@ -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") } @@ -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()) diff --git a/cache/wrapper.go b/cache/wrapper.go index 10f31cd99..b107fd4ea 100644 --- a/cache/wrapper.go +++ b/cache/wrapper.go @@ -4,9 +4,18 @@ 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, fn func() (V, int)) V - GetWithError(key uint32, fn func() (V, int, error)) (V, error) + Get(key uint32, l Loader[V]) (V, error) } var ( @@ -29,23 +38,12 @@ func NewSession[V any](c *Cache[V]) *Session[V] { } } -func (s *Session[V]) Get(key uint32, fn func() (V, int)) V { - if e := s.resolve(key); e != nil { - return e.value - } - - value, e := s.cache.get(key, fn) - s.local[key] = weak.Make(e) - - return value -} - -func (s *Session[V]) GetWithError(key uint32, fn func() (V, int, error)) (V, error) { +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.getWithError(key, fn) + value, e, err := s.cache.get(key, l) if err != nil { return value, err } @@ -79,23 +77,12 @@ func NewScan[V any]() *Scan[V] { return &Scan[V]{} } -func (s *Scan[V]) Get(key uint32, fn func() (V, int)) V { - if s.ok && s.key == key { - return s.value - } - - value, _ := fn() - s.key, s.value, s.ok = key, value, true - - return value -} - -func (s *Scan[V]) GetWithError(key uint32, fn func() (V, int, error)) (V, error) { +func (s *Scan[V]) Get(key uint32, l Loader[V]) (V, error) { if s.ok && s.key == key { return s.value, nil } - value, _, err := fn() + value, _, err := l.Load(key) if err != nil { return value, err } diff --git a/frac/sealed.go b/frac/sealed.go index 4ffd79424..ff5061386 100644 --- a/frac/sealed.go +++ b/frac/sealed.go @@ -143,9 +143,9 @@ func NewSealedPreloaded( } // Put token table built during sealing into the cache. - indexCache.TokenTable.Get(token.CacheKeyTable, func() (token.Table, int) { - return preloaded.TokenTable, preloaded.TokenTable.Size() - }) + _, _ = indexCache.TokenTable.Get(token.CacheKeyTable, cache.LoaderFunc[token.Table](func(uint32) (token.Table, int, error) { + return preloaded.TokenTable, preloaded.TokenTable.Size(), nil + })) docsCountK := float64(f.info.DocsTotal) / 1000 logger.Info("sealed fraction created from active", diff --git a/frac/sealed/lids/loader.go b/frac/sealed/lids/loader.go index 2992355a8..aea7757d5 100644 --- a/frac/sealed/lids/loader.go +++ b/frac/sealed/lids/loader.go @@ -62,14 +62,15 @@ func NewLoader(fracVer config.BinaryDataVersion, r *storage.IndexReader, c cache } func (l *Loader) GetLIDsBlock(blockIndex uint32) (*Block, error) { - return l.cache.GetWithError(blockIndex, func() (*Block, int, error) { - block, err := l.readLIDsBlock(blockIndex) - if err != nil { - return block, 0, err - } - size := block.GetSizeBytes() - return block, size, nil - }) + return l.cache.Get(blockIndex, l) +} + +func (l *Loader) Load(blockIndex uint32) (*Block, int, error) { + block, err := l.readLIDsBlock(blockIndex) + if err != nil { + return block, 0, err + } + return block, block.GetSizeBytes(), nil } func (l *Loader) readLIDsBlock(blockIndex uint32) (*Block, error) { diff --git a/frac/sealed/seqids/loader.go b/frac/sealed/seqids/loader.go index 3d92bede8..9329a3864 100644 --- a/frac/sealed/seqids/loader.go +++ b/frac/sealed/seqids/loader.go @@ -26,85 +26,104 @@ func (Table) BlockStartLID(blockIndex uint32) uint32 { } type Loader struct { + fracVersion config.BinaryDataVersion reader *storage.IndexReader table *Table - cacheMIDs cache.Wrapper[[]byte] - cacheRIDs cache.Wrapper[BlockRIDs] - cacheParams cache.Wrapper[BlockParams] - fracVersion config.BinaryDataVersion + + mids cache.Wrapper[[]byte] + rids cache.Wrapper[BlockRIDs] + params cache.Wrapper[BlockParams] +} + +type midsLoader Loader + +func (l *midsLoader) midBlockIndex(index uint32) uint32 { + return l.table.StartBlockIndex + index*3 +} + +func (s *midsLoader) Load(index uint32) ([]byte, int, error) { + data, _, err := s.reader.ReadIndexBlock(s.midBlockIndex(index), nil) + return data, cap(data), err } func (l *Loader) GetMIDsBlock(index uint32, unpackCache *unpackCache) (BlockMIDs, error) { - // load binary from index - data, err := l.cacheMIDs.GetWithError(index, func() ([]byte, int, error) { - data, _, err := l.reader.ReadIndexBlock(l.midBlockIndex(index), nil) - return data, cap(data), err - }) - // check errors + data, err := l.mids.Get(index, (*midsLoader)(l)) if err == nil && len(data) == 0 { - err = errors.New("empty block") + return BlockMIDs{}, errors.New("empty block") } + if err != nil { return BlockMIDs{}, err } + // unpack block := BlockMIDs{Values: unpackCache.values[:0]} if err := block.Unpack(data, l.fracVersion, unpackCache); err != nil { return BlockMIDs{}, err } + return block, nil } +// ridsSource implements [cache.Loader] for the RIDs cache layer of [Loader]. +type ridsSource Loader + +func (s *ridsSource) Load(index uint32) (BlockRIDs, int, error) { + l := (*Loader)(s) + + data, _, err := l.reader.ReadIndexBlock(l.ridBlockIndex(index), nil) + if err != nil { + return BlockRIDs{}, 0, err + } + + block := BlockRIDs{ + fracVersion: l.fracVersion, + Values: make([]uint64, 0, consts.IDsPerBlock), + } + + err = block.Unpack(data) + if err != nil { + return BlockRIDs{}, 0, err + } + + if len(block.Values) == 0 { + return BlockRIDs{}, 0, errors.New("empty block") + } + + const ui64 = int(unsafe.Sizeof(uint64(0))) + return block, cap(block.Values) * ui64, err +} + func (l *Loader) GetRIDsBlock(index uint32) (BlockRIDs, error) { - return l.cacheRIDs.GetWithError(index, func() (BlockRIDs, int, error) { - data, _, err := l.reader.ReadIndexBlock(l.ridBlockIndex(index), nil) - if err != nil { - return BlockRIDs{}, 0, err - } - - block := BlockRIDs{ - fracVersion: l.fracVersion, - Values: make([]uint64, 0, consts.IDsPerBlock), - } - - err = block.Unpack(data) - if err != nil { - return BlockRIDs{}, 0, err - } - - if len(block.Values) == 0 { - return BlockRIDs{}, 0, errors.New("empty block") - } - - const ui64 = int(unsafe.Sizeof(uint64(0))) - return block, cap(block.Values) * ui64, err - }) + return l.rids.Get(index, (*ridsSource)(l)) } -func (l *Loader) GetParamsBlock(index uint32) (BlockParams, error) { - return l.cacheParams.GetWithError(index, func() (BlockParams, int, error) { - data, _, err := l.reader.ReadIndexBlock(l.paramsBlockIndex(index), nil) - if err != nil { - return BlockParams{}, 0, err - } - - block := BlockParams{Values: make([]uint64, 0, consts.IDsPerBlock)} - if err := block.Unpack(data); err != nil { - return BlockParams{}, 0, err - } - - if len(block.Values) == 0 { - return BlockParams{}, 0, errors.New("empty block") - } - - const ui64 = int(unsafe.Sizeof(uint64(0))) - return block, cap(block.Values) * ui64, nil - }) +// paramsSource implements [cache.Loader] for the params cache layer of [Loader]. +type paramsSource Loader + +func (s *paramsSource) Load(index uint32) (BlockParams, int, error) { + l := (*Loader)(s) + + data, _, err := l.reader.ReadIndexBlock(l.paramsBlockIndex(index), nil) + if err != nil { + return BlockParams{}, 0, err + } + + block := BlockParams{Values: make([]uint64, 0, consts.IDsPerBlock)} + if err := block.Unpack(data); err != nil { + return BlockParams{}, 0, err + } + + if len(block.Values) == 0 { + return BlockParams{}, 0, errors.New("empty block") + } + + const ui64 = int(unsafe.Sizeof(uint64(0))) + return block, cap(block.Values) * ui64, nil } -// blocks are stored as triplets on disk, (MID + RID + Pos), check docs/format-index-file.go -func (l *Loader) midBlockIndex(index uint32) uint32 { - return l.table.StartBlockIndex + index*3 +func (l *Loader) GetParamsBlock(index uint32) (BlockParams, error) { + return l.params.Get(index, (*paramsSource)(l)) } func (l *Loader) ridBlockIndex(index uint32) uint32 { diff --git a/frac/sealed/seqids/provider.go b/frac/sealed/seqids/provider.go index d1c5f8cff..bc46dd437 100644 --- a/frac/sealed/seqids/provider.go +++ b/frac/sealed/seqids/provider.go @@ -29,9 +29,9 @@ func NewProvider( loader: Loader{ reader: indexReader, table: table, - cacheMIDs: cacheMIDs, - cacheRIDs: cacheRIDs, - cacheParams: cacheParams, + mids: cacheMIDs, + rids: cacheRIDs, + params: cacheParams, fracVersion: fracVersion, }, midCache: NewCache(), diff --git a/frac/sealed/token/block_loader.go b/frac/sealed/token/block_loader.go index 81eb46d9e..2ff149bbe 100644 --- a/frac/sealed/token/block_loader.go +++ b/frac/sealed/token/block_loader.go @@ -212,19 +212,25 @@ func NewBlockLoader( } } +type blockSource BlockLoader + +func (s *blockSource) Load(index uint32) (*Block, int, error) { + block, err := (*BlockLoader)(s).read(index) + if err != nil { + return nil, 0, err + } + return block, block.Size(), nil +} + func (l *BlockLoader) Load(index uint32) *Block { - block := l.cache.Get(index, func() (*Block, int) { - block, err := l.read(index) - if err != nil { - logger.Panic("error reading tokens block", // todo: get rid of panic here - zap.Error(err), - zap.Uint32("index", index), - zap.String("frac", l.fracName), - ) - } - size := block.Size() - return block, size - }) + block, err := l.cache.Get(index, (*blockSource)(l)) + if err != nil { + logger.Panic("error reading tokens block", // todo: get rid of panic here + zap.Error(err), + zap.Uint32("index", index), + zap.String("frac", l.fracName), + ) + } return block } diff --git a/frac/sealed/token/table_loader.go b/frac/sealed/token/table_loader.go index 82d3840ea..a82fb845e 100644 --- a/frac/sealed/token/table_loader.go +++ b/frac/sealed/token/table_loader.go @@ -48,28 +48,34 @@ func NewTableLoader( } } -func (l *TableLoader) Load() Table { - table, err := l.cache.GetWithError(CacheKeyTable, func() (Table, int, error) { - var ( - blocks []TableBlock - err error - ) - - l.advanceToTable() - - if l.isLegacy { - blocks, err = l.loadBlocksLegacy() - } else { - blocks, err = l.loadBlocks() - } +type tableSource TableLoader - if err != nil { - return nil, 0, err - } +func (s *tableSource) Load(uint32) (Table, int, error) { + l := (*TableLoader)(s) - table := TableFromBlocks(blocks) - return table, table.Size(), nil - }) + var ( + blocks []TableBlock + err error + ) + + l.advanceToTable() + + if l.isLegacy { + blocks, err = l.loadBlocksLegacy() + } else { + blocks, err = l.loadBlocks() + } + + if err != nil { + return nil, 0, err + } + + table := TableFromBlocks(blocks) + return table, table.Size(), nil +} + +func (l *TableLoader) Load() Table { + table, err := l.cache.Get(CacheKeyTable, (*tableSource)(l)) if err != nil { logger.Fatal("load token table error", zap.String("frac", l.fracName), diff --git a/skipmaskmanager/loader.go b/skipmaskmanager/loader.go index ca01eb962..2020cb416 100644 --- a/skipmaskmanager/loader.go +++ b/skipmaskmanager/loader.go @@ -42,14 +42,16 @@ func (l *loader) getFile() (*os.File, error) { } func (l *loader) getHeaders() ([]lidsBlockHeader, error) { - return l.headersCache.GetWithError(l.cashKey, func() ([]lidsBlockHeader, int, error) { - headers, err := l.loadHeaders() - if err != nil { - return headers, 0, err - } - size := len(headers) * int(lidsBlockHeaderSizeBytes) - return headers, size, nil - }) + return l.headersCache.Get(l.cashKey, l) +} + +func (l *loader) Load(uint32) ([]lidsBlockHeader, int, error) { + headers, err := l.loadHeaders() + if err != nil { + return headers, 0, err + } + size := len(headers) * int(lidsBlockHeaderSizeBytes) + return headers, size, nil } func (l *loader) loadHeaders() ([]lidsBlockHeader, error) { diff --git a/storage/docs_reader.go b/storage/docs_reader.go index ad5edbd83..dbc014690 100644 --- a/storage/docs_reader.go +++ b/storage/docs_reader.go @@ -41,14 +41,21 @@ func (r *DocsReader) ReadDocs(blockOffset uint64, docOffsets []uint64) ([][]byte return res, nil } +type docsLoader struct { + reader *DocBlocksReader + blockOffset uint64 +} + +func (s docsLoader) Load(uint32) ([]byte, int, error) { + block, _, err := s.reader.ReadDocBlockPayload(int64(s.blockOffset)) + if err != nil { + return nil, 0, fmt.Errorf("can't fetch doc at pos %d: %w", s.blockOffset, err) + } + return block, cap(block), nil +} + func (r *DocsReader) ReadDocsFunc(blockOffset uint64, docOffsets []uint64, cb func([]byte) error) error { - block, err := r.cache.GetWithError(uint32(blockOffset), func() ([]byte, int, error) { - block, _, err := r.reader.ReadDocBlockPayload(int64(blockOffset)) - if err != nil { - return nil, 0, fmt.Errorf("can't fetch doc at pos %d: %w", blockOffset, err) - } - return block, cap(block), nil - }) + block, err := r.cache.Get(uint32(blockOffset), docsLoader{reader: &r.reader, blockOffset: blockOffset}) if err != nil { return err } diff --git a/storage/index_reader.go b/storage/index_reader.go index a07e08513..7a9825673 100644 --- a/storage/index_reader.go +++ b/storage/index_reader.go @@ -10,6 +10,8 @@ import ( "github.com/ozontech/seq-db/util" ) +const registryCacheKey = 1 + type IndexReader struct { limiter *ReadLimiter @@ -31,55 +33,58 @@ func NewIndexReader( } } -func (r *IndexReader) GetBlockHeader(index uint32) (IndexBlockHeader, error) { - registry, err := r.cache.GetWithError(1, func() ([]byte, int, error) { - data, err := r.readRegistry() - return data, cap(data), err - }) - if err != nil { - return nil, err - } +type registryLoader IndexReader - if (uint64(index)+1)*IndexBlockHeaderSize > uint64(len(registry)) { - return nil, fmt.Errorf( - "too large index block in file %s, with index %d, registry size %d", - r.readerName, index, len(registry), - ) - } - - pos := index * IndexBlockHeaderSize - return registry[pos : pos+IndexBlockHeaderSize], nil -} - -func (r *IndexReader) readRegistry() ([]byte, error) { - numBuf := make([]byte, 16) - - n, err := r.limiter.ReadAt(r.reader, numBuf, 0) +func (r *registryLoader) Load(uint32) ([]byte, int, error) { + prefix := make([]byte, 16) + n, err := r.limiter.ReadAt(r.reader, prefix, 0) if err != nil { - return nil, fmt.Errorf("can't read disk registry, %s", err.Error()) + return nil, 0, fmt.Errorf("can't read disk registry, %s", err.Error()) } + if n == 0 { - return nil, fmt.Errorf("can't read disk registry, n=0") + return nil, 0, fmt.Errorf("can't read disk registry, n=0") } - pos := binary.LittleEndian.Uint64(numBuf) - l := binary.LittleEndian.Uint64(numBuf[8:]) - buf := make([]byte, l) + pos := binary.LittleEndian.Uint64(prefix) + size := binary.LittleEndian.Uint64(prefix[8:]) + buf := make([]byte, size) n, err = r.limiter.ReadAt(r.reader, buf, int64(pos)) if err != nil && err != io.EOF { - return nil, fmt.Errorf("can't read disk registry, %s", err.Error()) + return nil, 0, fmt.Errorf("can't read disk registry, %s", err.Error()) } - if uint64(n) != l { - return nil, fmt.Errorf("can't read disk registry, read=%d, requested=%d", n, l) + if uint64(n) != size { + return nil, 0, fmt.Errorf("can't read disk registry, read=%d, requested=%d", n, size) } if len(buf)%IndexBlockHeaderSize != 0 { - return nil, fmt.Errorf("wrong registry format") + return nil, 0, fmt.Errorf("wrong registry format") } - return buf, nil + return buf, cap(buf), nil +} + +func (r *IndexReader) registry() ([]byte, error) { + return r.cache.Get(registryCacheKey, (*registryLoader)(r)) +} + +func (r *IndexReader) GetBlockHeader(index uint32) (IndexBlockHeader, error) { + reg, err := r.registry() + if err != nil { + return nil, err + } + + if (uint64(index)+1)*IndexBlockHeaderSize > uint64(len(reg)) { + return nil, fmt.Errorf( + "too large index block in file %s, with index %d, registry size %d", + r.readerName, index, len(reg), + ) + } + + pos := index * IndexBlockHeaderSize + return reg[pos : pos+IndexBlockHeaderSize], nil } func (r *IndexReader) ReadIndexBlock(blockIndex uint32, dst []byte) ([]byte, uint64, error) { @@ -109,13 +114,10 @@ func (r *IndexReader) ReadIndexBlock(blockIndex uint32, dst []byte) ([]byte, uin } func (r *IndexReader) BlocksCount() (int, error) { - registry, err := r.cache.GetWithError(1, func() ([]byte, int, error) { - data, err := r.readRegistry() - return data, cap(data), err - }) + reg, err := r.registry() if err != nil { return 0, err } - return len(registry) / IndexBlockHeaderSize, nil + return len(reg) / IndexBlockHeaderSize, nil } From 32613475c776e6763601b26bae9837b79498f2d4 Mon Sep 17 00:00:00 2001 From: Daniil Porokhnin Date: Fri, 14 Aug 2026 12:31:37 +0300 Subject: [PATCH 4/6] chore: tune gogc/gomemlimit --- .seqbench/continuous.env | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) diff --git a/.seqbench/continuous.env b/.seqbench/continuous.env index 5589bd189..a8348157c 100644 --- a/.seqbench/continuous.env +++ b/.seqbench/continuous.env @@ -1,5 +1,6 @@ -GOGC=500 -GOMEMLIMIT=4GiB +GOGC=100 +# GOGC=500 +# GOMEMLIMIT=4GiB SEQDB_RESOURCES_SKIP_FSYNC=true From 59eaea1b8e5b2aa800e159fab7804a69e6299bbf Mon Sep 17 00:00:00 2001 From: Daniil Porokhnin Date: Fri, 14 Aug 2026 16:20:19 +0300 Subject: [PATCH 5/6] chore: GOGC=200 --- .seqbench/continuous.env | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.seqbench/continuous.env b/.seqbench/continuous.env index a8348157c..5eee83dd0 100644 --- a/.seqbench/continuous.env +++ b/.seqbench/continuous.env @@ -1,4 +1,4 @@ -GOGC=100 +GOGC=200 # GOGC=500 # GOMEMLIMIT=4GiB From 9262aea95b9c57abafc1237aacd6f8db43c56669 Mon Sep 17 00:00:00 2001 From: Daniil Porokhnin Date: Fri, 14 Aug 2026 17:06:18 +0300 Subject: [PATCH 6/6] chore: GO_VERSION=1.26 --- .seqbench/continuous.env | 1 + 1 file changed, 1 insertion(+) diff --git a/.seqbench/continuous.env b/.seqbench/continuous.env index 5eee83dd0..7de3086a1 100644 --- a/.seqbench/continuous.env +++ b/.seqbench/continuous.env @@ -1,6 +1,7 @@ GOGC=200 # GOGC=500 # GOMEMLIMIT=4GiB +GO_VERSION=1.26 SEQDB_RESOURCES_SKIP_FSYNC=true