Skip to content
Open
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
7 changes: 6 additions & 1 deletion cmd/distribution/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -48,7 +48,12 @@ func getReader(path string) (storage.IndexReader, *os.File) {
if err != nil {
panic(err)
}
return storage.NewIndexReader(readLimiter, f.Name(), f, c), f
readerProvider := storage.NewReaderProvider(readLimiter, f.Name(), f, c)
reader, err := readerProvider.GetReader()
if err != nil {
panic(err)
}
return reader, f
}

func readBlock(reader storage.IndexReader, blockIndex uint32) ([]byte, error) {
Expand Down
71 changes: 43 additions & 28 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
legacyReader storage.IndexReader
IsLegacy bool
legacyFile storage.ImmutableFile
legacyReaderProvider *storage.ReaderProvider

// 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

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

indexCache *IndexCache

Expand Down Expand Up @@ -161,6 +161,14 @@ 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 @@ -171,14 +179,21 @@ func (f *Remote) createDataProvider(ctx context.Context) (*sealedDataProvider, e
return nil, err
}

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

if f.IsLegacy {
tokenReader = &f.legacyReader
lidReader = &f.legacyReader
idReader = &f.legacyReader
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)
}

return &sealedDataProvider{
Expand All @@ -190,13 +205,13 @@ 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, 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),

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

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

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

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

return f.openRemoteFile(consts.IndexFileSuffix, func(file storage.ImmutableFile) {
f.legacyFile = file
f.legacyReader = storage.NewIndexReader(
f.legacyReaderProvider = storage.NewReaderProvider(
f.readLimiter, file.Name(),
file, f.indexCache.LegacyRegistry,
)
Expand Down Expand Up @@ -350,7 +365,7 @@ func (f *Remote) openIndex() error {
consts.TokenFileSuffix,
func(file storage.ImmutableFile) {
f.tokenFile = file
f.tokenReader = storage.NewIndexReader(
f.tokenReaderProvider = storage.NewReaderProvider(
f.readLimiter, file.Name(),
file, f.indexCache.TokenRegistry,
)
Expand All @@ -365,7 +380,7 @@ func (f *Remote) openIndex() error {
consts.OffsetsFileSuffix,
func(file storage.ImmutableFile) {
f.offsetsFile = file
f.offsetsReader = storage.NewIndexReader(
f.offsetsReaderProvider = storage.NewReaderProvider(
f.readLimiter, file.Name(),
file, f.indexCache.OffsetsRegistry,
)
Expand All @@ -380,7 +395,7 @@ func (f *Remote) openIndex() error {
consts.IDFileSuffix,
func(file storage.ImmutableFile) {
f.idFile = file
f.idReader = storage.NewIndexReader(
f.idReaderProvider = storage.NewReaderProvider(
f.readLimiter, file.Name(),
file, f.indexCache.IDRegistry,
)
Expand All @@ -395,7 +410,7 @@ func (f *Remote) openIndex() error {
consts.LIDFileSuffix,
func(file storage.ImmutableFile) {
f.lidFile = file
f.lidReader = storage.NewIndexReader(
f.lidReaderProvider = storage.NewReaderProvider(
f.readLimiter, file.Name(),
file, f.indexCache.LIDRegistry,
)
Expand Down
71 changes: 43 additions & 28 deletions frac/sealed.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
legacyReader storage.IndexReader
IsLegacy bool
legacyFile *os.File
legacyReaderProvider *storage.ReaderProvider

// Per-section index files and their readers (new split format only).
infoFile *os.File
Expand All @@ -51,10 +51,10 @@ type Sealed struct {
idFile *os.File
lidFile *os.File

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

blocksData sealed.BlocksData
indexCache *IndexCache
Expand Down Expand Up @@ -176,7 +176,7 @@ func (f *Sealed) openInfoLegacy() {
}

f.legacyFile = file
f.legacyReader = storage.NewIndexReader(
f.legacyReaderProvider = storage.NewReaderProvider(
f.readLimiter, file.Name(),
file, f.indexCache.LegacyRegistry,
)
Expand Down Expand Up @@ -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.tokenReader = storage.NewIndexReader(f.readLimiter, file.Name(), file, f.indexCache.TokenRegistry)
f.tokenReaderProvider = storage.NewReaderProvider(f.readLimiter, file.Name(), file, f.indexCache.TokenRegistry)
}

if f.offsetsFile == nil {
Expand All @@ -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.offsetsReader = storage.NewIndexReader(f.readLimiter, file.Name(), file, f.indexCache.OffsetsRegistry)
f.offsetsReaderProvider = storage.NewReaderProvider(f.readLimiter, file.Name(), file, f.indexCache.OffsetsRegistry)
}

if f.idFile == nil {
Expand All @@ -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.idReader = storage.NewIndexReader(f.readLimiter, file.Name(), file, f.indexCache.IDRegistry)
f.idReaderProvider = storage.NewReaderProvider(f.readLimiter, file.Name(), file, f.indexCache.IDRegistry)
}

if f.lidFile == nil {
Expand All @@ -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.lidReader = storage.NewIndexReader(f.readLimiter, file.Name(), file, f.indexCache.LIDRegistry)
f.lidReaderProvider = storage.NewReaderProvider(f.readLimiter, file.Name(), file, f.indexCache.LIDRegistry)
}
}

Expand Down Expand Up @@ -279,12 +279,20 @@ 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.legacyReader); err != nil {
if f.info, err = loadInfoLegacy(f.mustGetReader(f.legacyReaderProvider)); err != nil {
logger.Fatal("error loading Info", zap.String("fraction", f.BaseFileName), zap.Error(err))
}
return
Expand All @@ -308,16 +316,16 @@ func (f *Sealed) init(full bool) {
}

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

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

f.isInited = true
Expand Down Expand Up @@ -495,14 +503,21 @@ func (f *Sealed) FindLIDs(ctx context.Context, ids []seq.ID) ([]seq.LID, error)
func (f *Sealed) createDataProvider(ctx context.Context) *sealedDataProvider {
f.init(true)

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

if f.IsLegacy {
tokenReader = &f.legacyReader
lidReader = &f.legacyReader
idReader = &f.legacyReader
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)
}

return &sealedDataProvider{
Expand All @@ -514,13 +529,13 @@ 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, 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),

idsTable: &f.blocksData.IDsTable,
idsProvider: seqids.NewProvider(
idReader,
&idReader,
f.indexCache.MIDs,
f.indexCache.RIDs,
f.indexCache.Params,
Expand Down
5 changes: 1 addition & 4 deletions frac/sealed/token/table_loader.go
Original file line number Diff line number Diff line change
Expand Up @@ -141,10 +141,7 @@ func (l *TableLoader) loadBlocksLegacy() ([]TableBlock, error) {
}

func (l *TableLoader) loadBlocks() ([]TableBlock, error) {
blocksCount, err := l.reader.BlocksCount()
if err != nil {
return nil, err
}
blocksCount := l.reader.BlocksCount()

var blocks []TableBlock
for blockIndex := l.tableIndex; blockIndex < uint32(blocksCount); blockIndex++ {
Expand Down
16 changes: 2 additions & 14 deletions frac/sealed_loader.go
Original file line number Diff line number Diff line change
Expand Up @@ -231,13 +231,7 @@ func (l *Loader) loadIDsTable(r storage.IndexReader, info *common.Info) seqids.T
IDsTotal: info.DocsTotal + 1, // Increment by one for [seq.SystemID]
}

blocksCount, err := r.BlocksCount()
if err != nil {
logger.Fatal(
"cannot get block count",
zap.Error(err),
)
}
blocksCount := r.BlocksCount()

for blockIdx := 0; blockIdx < blocksCount; blockIdx += 3 {
header, err := r.GetBlockHeader(uint32(blockIdx))
Expand Down Expand Up @@ -269,13 +263,7 @@ func (l *Loader) loadLIDsTable(r storage.IndexReader) (*lids.Table, error) {
isContinued []bool
)

blocksCount, err := r.BlocksCount()
if err != nil {
logger.Fatal(
"cannot get block count",
zap.Error(err),
)
}
blocksCount := r.BlocksCount()

for blockIdx := 0; blockIdx < blocksCount; blockIdx++ {
header, err := r.GetBlockHeader(uint32(blockIdx))
Expand Down
Loading
Loading