From 4fedbdf3e4030fdc65a4153b1727fc7a2134c6e1 Mon Sep 17 00:00:00 2001 From: Daniil Porokhnin Date: Sat, 19 Sep 2026 16:52:50 +0300 Subject: [PATCH 1/2] fix: cache block index instead of block offset --- frac/active.go | 29 +++++++++++++++------------- frac/active_index.go | 15 ++++----------- frac/active_sealing_source.go | 9 +++++---- frac/processor/fetch.go | 5 ++--- frac/remote.go | 21 +++++++++++--------- frac/sealed.go | 21 +++++++++++--------- frac/sealed_index.go | 15 ++++----------- storage/docs_reader.go | 36 +++++++++++++++++++++++++---------- 8 files changed, 81 insertions(+), 70 deletions(-) diff --git a/frac/active.go b/frac/active.go index 737a76800..36b2de8b8 100644 --- a/frac/active.go +++ b/frac/active.go @@ -46,11 +46,11 @@ type Active struct { DocsPositions *DocsPositions IDsToLIDs *ActiveLIDs - docsFile *os.File - docsReader storage.DocsReader - sortReader storage.DocsReader - docsCache *cache.ConcurrentCache[[]byte] - sortCache *cache.ConcurrentCache[[]byte] + docsFile *os.File + docsCache *cache.ConcurrentCache[[]byte] + sortCache *cache.ConcurrentCache[[]byte] + + readLimiter *storage.ReadLimiter walFile *os.File walReader *storage.WalReader @@ -81,11 +81,11 @@ func NewActive( RIDs: NewIDs(), DocBlocks: NewIDs(), - docsFile: docsFile, - docsCache: docsCache, - sortCache: sortCache, - docsReader: storage.NewDocsReader(readLimiter, docsFile, docsCache), - sortReader: storage.NewDocsReader(readLimiter, docsFile, sortCache), + docsFile: docsFile, + docsCache: docsCache, + sortCache: sortCache, + + readLimiter: readLimiter, walFile: walFile, walReader: walReader, @@ -357,10 +357,13 @@ func (f *Active) createDataProvider(ctx context.Context) *activeDataProvider { rids: f.RIDs, tokenList: f.TokenList, - blocksOffsets: f.DocBlocks.GetVals(), - docsPositions: f.DocsPositions, idsToLids: f.IDsToLIDs, - docsReader: &f.docsReader, + docsPositions: f.DocsPositions, + + docsReader: storage.NewDocsReader( + f.readLimiter, f.docsFile, + f.docsCache, f.DocBlocks.GetVals(), + ), skipMaskProvider: f.skipMaskProvider, } diff --git a/frac/active_index.go b/frac/active_index.go index 28fd2ea9b..622856c95 100644 --- a/frac/active_index.go +++ b/frac/active_index.go @@ -25,10 +25,9 @@ type activeDataProvider struct { tokenList *TokenList - blocksOffsets []uint64 docsPositions *DocsPositions idsToLids *ActiveLIDs - docsReader *storage.DocsReader + docsReader storage.DocsReader idsIndex *activeIDsIndex @@ -79,10 +78,9 @@ func (dp *activeDataProvider) Fetch(ids []seq.ID, noSkipMasks bool) ([][]byte, e res := make([][]byte, len(ids)) indexes := []activeFetchIndex{{ - blocksOffsets: dp.blocksOffsets, docsPositions: dp.docsPositions, idsToLids: dp.idsToLids, - docsReader: dp.docsReader, + docsReader: &dp.docsReader, skipMaskProvider: dp.skipMaskProvider, fracName: dp.info.Name(), }} @@ -295,7 +293,6 @@ func inverseLIDs(unmapped []uint32, inv *inverser, minLID, maxLID uint32) []uint } type activeFetchIndex struct { - blocksOffsets []uint64 docsPositions *DocsPositions idsToLids *ActiveLIDs docsReader *storage.DocsReader @@ -303,10 +300,6 @@ type activeFetchIndex struct { fracName string } -func (di *activeFetchIndex) GetBlocksOffsets(num uint32) uint64 { - return di.blocksOffsets[num] -} - func (di *activeFetchIndex) GetDocPos(ids []seq.ID, noSkipMasks bool) ([]seq.DocPos, error) { docsPos := make([]seq.DocPos, len(ids)) for i, id := range ids { @@ -343,6 +336,6 @@ func (di *activeFetchIndex) GetDocPos(ids []seq.ID, noSkipMasks bool) ([]seq.Doc return docsPos, nil } -func (di *activeFetchIndex) ReadDocs(blockOffset uint64, docOffsets []uint64) ([][]byte, error) { - return di.docsReader.ReadDocs(blockOffset, docOffsets) +func (di *activeFetchIndex) ReadDocs(blockIndex uint32, docOffsets []uint64) ([][]byte, error) { + return di.docsReader.ReadDocs(blockIndex, docOffsets) } diff --git a/frac/active_sealing_source.go b/frac/active_sealing_source.go index aa7afebd5..5feac7bd7 100644 --- a/frac/active_sealing_source.go +++ b/frac/active_sealing_source.go @@ -51,7 +51,7 @@ type ActiveSealingSource struct { docPosMap map[seq.ID]seq.DocPos // Original document positions docPosSorted []seq.DocPos // Document positions after sorting - docsReader *storage.DocsReader // Document storage reader + docsReader storage.DocsReader // Document storage reader } func NewActiveSealingSource(active *Active, params common.SealParams) (*ActiveSealingSource, error) { @@ -79,7 +79,9 @@ func NewActiveSealingSource(active *Active, params common.SealParams) (*ActiveSe docPosMap: active.DocsPositions.idToPos, blocksOffsets: active.DocBlocks.vals, - docsReader: &active.sortReader, + docsReader: storage.NewDocsReader( + active.readLimiter, active.docsFile, active.sortCache, active.DocBlocks.vals, + ), } src.prepareInfo() @@ -265,11 +267,10 @@ func (src *ActiveSealingSource) Docs() iter.Seq2[Document, error] { // doc reads a document from storage by its position. func (src *ActiveSealingSource) doc(pos seq.DocPos) ([]byte, error) { blockIndex, docOffset := pos.Unpack() - blockOffset := src.blocksOffsets[blockIndex] var doc []byte err := src.docsReader.ReadDocsFunc( - blockOffset, []uint64{docOffset}, + blockIndex, []uint64{docOffset}, func(b []byte) error { doc = b return nil diff --git a/frac/processor/fetch.go b/frac/processor/fetch.go index 4a3661b8c..cefc45d6a 100644 --- a/frac/processor/fetch.go +++ b/frac/processor/fetch.go @@ -6,9 +6,8 @@ import ( ) type fetchIndex interface { - GetBlocksOffsets(uint32) uint64 GetDocPos([]seq.ID, bool) ([]seq.DocPos, error) - ReadDocs(blockOffset uint64, docOffsets []uint64) ([][]byte, error) + ReadDocs(blockIndex uint32, docOffsets []uint64) ([][]byte, error) } func IndexFetch(ids []seq.ID, noSkipMasks bool, sw *stopwatch.Stopwatch, fetchIndex fetchIndex, res [][]byte) error { @@ -22,7 +21,7 @@ func IndexFetch(ids []seq.ID, noSkipMasks bool, sw *stopwatch.Stopwatch, fetchIn m = sw.Start("read_doc") for i, docOffsets := range offsets { - docs, err := fetchIndex.ReadDocs(fetchIndex.GetBlocksOffsets(blocks[i]), docOffsets) + docs, err := fetchIndex.ReadDocs(blocks[i], docOffsets) if err != nil { return err } diff --git a/frac/remote.go b/frac/remote.go index dad0f9ff7..b0e444427 100644 --- a/frac/remote.go +++ b/frac/remote.go @@ -39,9 +39,8 @@ type Remote struct { info *common.Info - docsFile storage.ImmutableFile - docsCache *cache.ConcurrentCache[[]byte] - docsReader storage.DocsReader + docsFile storage.ImmutableFile + docsCache *cache.ConcurrentCache[[]byte] // IsLegacy is true for fractions that use the old single .index file format. IsLegacy bool @@ -170,10 +169,9 @@ func (f *Remote) createDataProvider(ctx context.Context) (*sealedDataProvider, e ctx: ctx, fractionTypeLabel: "remote", - info: f.info, - config: f.Config, - docsReader: &f.docsReader, - blocksOffsets: f.blocksData.BlocksOffsets, + info: f.info, + config: f.Config, + docsReader: f.docsReader(), lidsTable: f.blocksData.LIDsTable, lidsLoader: lids.NewLoader(f.info.BinaryDataVer, &ir.LID, cache.NewSession(f.indexCache.LIDs)), @@ -226,6 +224,13 @@ func (f *Remote) indexReaders() IndexReaders { } } +func (f *Remote) docsReader() storage.DocsReader { + return storage.NewDocsReader( + f.readLimiter, f.docsFile, + f.docsCache, f.blocksData.BlocksOffsets, + ) +} + func (f *Remote) Info() *common.Info { return f.info } @@ -454,7 +459,6 @@ func (f *Remote) openDocs() error { if unsortedExists { f.docsFile = s3.NewReader(f.ctx, f.s3cli, unsortedName) - f.docsReader = storage.NewDocsReader(f.readLimiter, f.docsFile, f.docsCache) return nil } @@ -468,7 +472,6 @@ func (f *Remote) openDocs() error { if sortedExists { f.docsFile = s3.NewReader(f.ctx, f.s3cli, sortedName) - f.docsReader = storage.NewDocsReader(f.readLimiter, f.docsFile, f.docsCache) return nil } diff --git a/frac/sealed.go b/frac/sealed.go index 32b5f01d4..65db0fb06 100644 --- a/frac/sealed.go +++ b/frac/sealed.go @@ -35,9 +35,8 @@ type Sealed struct { info *common.Info - docsFile *os.File - docsCache *cache.ConcurrentCache[[]byte] - docsReader storage.DocsReader + docsFile *os.File + docsCache *cache.ConcurrentCache[[]byte] // IsLegacy is true for fractions that use the old single .index file format. IsLegacy bool @@ -258,8 +257,6 @@ func (f *Sealed) openDocs() { ) } } - - f.docsReader = storage.NewDocsReader(f.readLimiter, f.docsFile, f.docsCache) } func (f *Sealed) loadInfo() { @@ -499,10 +496,9 @@ func (f *Sealed) createDataProvider(ctx context.Context) *sealedDataProvider { ctx: ctx, fractionTypeLabel: "sealed", - info: f.info, - config: f.Config, - docsReader: &f.docsReader, - blocksOffsets: f.blocksData.BlocksOffsets, + info: f.info, + config: f.Config, + docsReader: f.docsReader(), lidsTable: f.blocksData.LIDsTable, lidsLoader: lids.NewLoader(f.info.BinaryDataVer, &ir.LID, cache.NewSession(f.indexCache.LIDs)), @@ -568,6 +564,13 @@ func (f *Sealed) indexReaders() IndexReaders { } } +func (f *Sealed) docsReader() storage.DocsReader { + return storage.NewDocsReader( + f.readLimiter, f.docsFile, + f.docsCache, f.blocksData.BlocksOffsets, + ) +} + // computeIndexOnDisk returns the total on-disk size of index files for a local fraction. func (f *Sealed) computeIndexSize() { suffixes := []string{ diff --git a/frac/sealed_index.go b/frac/sealed_index.go index e520d651e..55a4d5546 100644 --- a/frac/sealed_index.go +++ b/frac/sealed_index.go @@ -43,8 +43,7 @@ type sealedDataProvider struct { tokenBlockLoader *token.BlockLoader tokenTableLoader *token.TableLoader - blocksOffsets []uint64 - docsReader *storage.DocsReader + docsReader storage.DocsReader // fractionTypeLabel can be either 'sealed' or 'remote'. // This value is used in metrics to distinguish between operations over local and remote fractions. @@ -65,8 +64,7 @@ func (dp *sealedDataProvider) getFetchIndex() *sealedFetchIndex { return &sealedFetchIndex{ fracName: dp.info.Name(), idsIndex: dp.getIDsIndex(), - docsReader: dp.docsReader, - blocksOffsets: dp.blocksOffsets, + docsReader: &dp.docsReader, skipMaskProvider: dp.skipMaskProvider, } } @@ -362,14 +360,9 @@ type sealedFetchIndex struct { fracName string idsIndex *sealedIDsIndex docsReader *storage.DocsReader - blocksOffsets []uint64 skipMaskProvider skipMaskProvider } -func (fi *sealedFetchIndex) GetBlocksOffsets(num uint32) uint64 { - return fi.blocksOffsets[num] -} - func (fi *sealedFetchIndex) GetDocPos(ids []seq.ID, noSkipMasks bool) ([]seq.DocPos, error) { allLids := fi.findLIDs(ids) @@ -406,8 +399,8 @@ func (fi *sealedFetchIndex) GetDocPos(ids []seq.ID, noSkipMasks bool) ([]seq.Doc return fi.getDocPosByLIDs(allLids), nil } -func (fi *sealedFetchIndex) ReadDocs(blockOffset uint64, docOffsets []uint64) ([][]byte, error) { - return fi.docsReader.ReadDocs(blockOffset, docOffsets) +func (fi *sealedFetchIndex) ReadDocs(blockIndex uint32, docOffsets []uint64) ([][]byte, error) { + return fi.docsReader.ReadDocs(blockIndex, docOffsets) } // findLIDs returns a slice of LIDs. If seq.ID is not found, LID has the value 0 at the corresponding position diff --git a/storage/docs_reader.go b/storage/docs_reader.go index 939c24e75..a4302d098 100644 --- a/storage/docs_reader.go +++ b/storage/docs_reader.go @@ -9,21 +9,28 @@ import ( ) type DocsReader struct { - reader DocBlocksReader - cache *cache.ConcurrentCache[[]byte] + reader DocBlocksReader + cache *cache.ConcurrentCache[[]byte] + blockOffsets []uint64 } -func NewDocsReader(limiter *ReadLimiter, reader io.ReaderAt, docsCache *cache.ConcurrentCache[[]byte]) DocsReader { +func NewDocsReader( + limiter *ReadLimiter, + reader io.ReaderAt, + docsCache *cache.ConcurrentCache[[]byte], + blockOffsets []uint64, +) DocsReader { return DocsReader{ - reader: NewDocBlocksReader(limiter, reader), - cache: docsCache, + reader: NewDocBlocksReader(limiter, reader), + cache: docsCache, + blockOffsets: blockOffsets, } } -func (r *DocsReader) ReadDocs(blockOffset uint64, docOffsets []uint64) ([][]byte, error) { +func (r *DocsReader) ReadDocs(blockIndex uint32, docOffsets []uint64) ([][]byte, error) { bufSize := 0 res := make([][]byte, 0, len(docOffsets)) - err := r.ReadDocsFunc(blockOffset, docOffsets, func(doc []byte) error { + err := r.ReadDocsFunc(blockIndex, docOffsets, func(doc []byte) error { bufSize += len(doc) res = append(res, doc) return nil @@ -41,16 +48,25 @@ func (r *DocsReader) ReadDocs(blockOffset uint64, docOffsets []uint64) ([][]byte return res, nil } -func (r *DocsReader) Load(blockOffset uint32) ([]byte, int, error) { +func (r *DocsReader) Load(blockIndex uint32) ([]byte, int, error) { + if uint64(blockIndex) >= uint64(len(r.blockOffsets)) { + return nil, 0, fmt.Errorf( + "doc block index %d is out of range [0, %d)", + blockIndex, len(r.blockOffsets), + ) + } + + blockOffset := r.blockOffsets[blockIndex] 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 } -func (r *DocsReader) ReadDocsFunc(blockOffset uint64, docOffsets []uint64, cb func([]byte) error) error { - block, err := r.cache.Get(uint32(blockOffset), r) +func (r *DocsReader) ReadDocsFunc(blockIndex uint32, docOffsets []uint64, cb func([]byte) error) error { + block, err := r.cache.Get(blockIndex, r) if err != nil { return err } From 89a6e55c0b3c0b05313e2a27b153064ca4690df5 Mon Sep 17 00:00:00 2001 From: Daniil Porokhnin Date: Sat, 19 Sep 2026 17:08:37 +0300 Subject: [PATCH 2/2] chore: add build target for release build --- Makefile | 9 ++++++++- 1 file changed, 8 insertions(+), 1 deletion(-) diff --git a/Makefile b/Makefile index 990fbb034..1bae88475 100644 --- a/Makefile +++ b/Makefile @@ -22,6 +22,13 @@ build-image: -t ${IMAGE}:${VERSION} \ . +.PHONY: build +build: + CGO_ENABLED=1 \ + go build -tags 'netgo osusergo' \ + -o ${LOCAL_BIN}/${OS}-${ARCH}/ \ + ./cmd/... + .PHONY: build-debug build-debug: CGO_ENABLED=0 \ @@ -30,7 +37,7 @@ build-debug: ./cmd/... .PHONY: run -run: build-debug +run: build SEQDB_STORAGE_DATA_DIR=$(shell mktemp -d) \ ${LOCAL_BIN}/${OS}-${ARCH}/seq-db \ --mode=single \