Skip to content
Merged
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
149 changes: 79 additions & 70 deletions frac/remote.go
Original file line number Diff line number Diff line change
Expand Up @@ -44,22 +44,16 @@ 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

// Per-section index files and their readers (new split format only).
// Per-section index files (new split format only).
infoFile storage.ImmutableFile
tokenFile storage.ImmutableFile
offsetsFile storage.ImmutableFile
idFile storage.ImmutableFile
lidFile storage.ImmutableFile

tokenReader storage.IndexReader
offsetsReader storage.IndexReader
idReader storage.IndexReader
lidReader storage.IndexReader

indexCache *IndexCache

initMu *sync.RWMutex
Expand Down Expand Up @@ -171,32 +165,25 @@ func (f *Remote) createDataProvider(ctx context.Context) (*sealedDataProvider, e
return nil, err
}

tokenReader := &f.tokenReader
lidReader := &f.lidReader
idReader := &f.idReader

if f.IsLegacy {
tokenReader = &f.legacyReader
lidReader = &f.legacyReader
idReader = &f.legacyReader
}

ir := f.indexReaders()
return &sealedDataProvider{
ctx: ctx,
fractionTypeLabel: "remote",

info: f.info,
config: f.Config,
docsReader: &f.docsReader,
blocksOffsets: f.blocksData.BlocksOffsets,
lidsTable: f.blocksData.LIDsTable,
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)),
info: f.info,
config: f.Config,
docsReader: &f.docsReader,
blocksOffsets: f.blocksData.BlocksOffsets,

lidsTable: f.blocksData.LIDsTable,
lidsLoader: lids.NewLoader(f.info.BinaryDataVer, &ir.LID, cache.NewSession(f.indexCache.LIDs)),

tokenBlockLoader: token.NewBlockLoader(f.BaseFileName, f.Info().BinaryDataVer, &ir.Token, cache.NewSession(f.indexCache.Tokens)),
tokenTableLoader: token.NewTableLoader(f.BaseFileName, f.Info().BinaryDataVer, f.IsLegacy, &ir.Token, cache.NewSession(f.indexCache.TokenTable)),

idsTable: &f.blocksData.IDsTable,
idsProvider: seqids.NewProvider(
idReader,
&ir.ID,
cache.NewSession(f.indexCache.MIDs),
cache.NewSession(f.indexCache.RIDs),
cache.NewSession(f.indexCache.Params),
Expand All @@ -207,6 +194,38 @@ func (f *Remote) createDataProvider(ctx context.Context) (*sealedDataProvider, e
}, nil
}

func (f *Remote) indexReaders() IndexReaders {
if f.IsLegacy {
r := storage.NewIndexReader(
f.readLimiter, f.legacyFile.Name(), f.legacyFile,
cache.NewSession(f.indexCache.LegacyRegistry),
)
return IndexReaders{Token: r, Offsets: r, ID: r, LID: r}
}

return IndexReaders{
Token: storage.NewIndexReader(
f.readLimiter, f.tokenFile.Name(), f.tokenFile,
cache.NewSession(f.indexCache.TokenRegistry),
),

Offsets: storage.NewIndexReader(
f.readLimiter, f.offsetsFile.Name(), f.offsetsFile,
cache.NewSession(f.indexCache.OffsetsRegistry),
),

ID: storage.NewIndexReader(
f.readLimiter, f.idFile.Name(), f.idFile,
cache.NewSession(f.indexCache.IDRegistry),
),

LID: storage.NewIndexReader(
f.readLimiter, f.lidFile.Name(), f.lidFile,
cache.NewSession(f.indexCache.LIDRegistry),
),
}
}

func (f *Remote) Info() *common.Info {
return f.info
}
Expand Down Expand Up @@ -260,9 +279,20 @@ func (f *Remote) loadInfo() error {
if err := f.openInfoLegacy(); err != nil {
return err
}
if f.info, err = loadInfoLegacy(f.legacyReader); err != nil {
logger.Fatal("error loading Info", zap.String("fraction", f.BaseFileName), zap.Error(err))

legacyReader := storage.NewIndexReader(
f.readLimiter, f.legacyFile.Name(),
f.legacyFile, f.indexCache.LegacyRegistry,
)

if f.info, err = loadInfoLegacy(legacyReader); err != nil {
logger.Fatal(
"error loading Info",
zap.String("fraction", f.BaseFileName),
zap.Error(err),
)
}

return nil
}

Expand All @@ -271,8 +301,13 @@ func (f *Remote) loadInfo() error {
}

if f.info, err = loadInfo(f.infoFile); err != nil {
logger.Fatal("error loading Info", zap.String("fraction", f.BaseFileName), zap.Error(err))
logger.Fatal(
"error loading Info",
zap.String("fraction", f.BaseFileName),
zap.Error(err),
)
}

return nil
}

Expand All @@ -293,17 +328,19 @@ func (f *Remote) init() error {
}

if f.IsLegacy {
(&LegacyLoader{}).Load(&f.blocksData, f.info, f.legacyReader)
(&LegacyLoader{}).Load(
&f.blocksData, f.info,
storage.NewIndexReader(
f.readLimiter, f.legacyFile.Name(),
f.legacyFile, f.indexCache.LegacyRegistry,
),
)

f.isInited = true
return nil
}

(&Loader{}).Load(&f.blocksData, f.info, IndexReaders{
Token: f.tokenReader,
Offsets: f.offsetsReader,
ID: f.idReader,
LID: f.lidReader,
})
(&Loader{}).Load(&f.blocksData, f.info, f.indexReaders())

f.isInited = true
return nil
Expand All @@ -316,10 +353,6 @@ func (f *Remote) openInfoLegacy() error {

return f.openRemoteFile(consts.IndexFileSuffix, func(file storage.ImmutableFile) {
f.legacyFile = file
f.legacyReader = storage.NewIndexReader(
f.readLimiter, file.Name(),
file, f.indexCache.LegacyRegistry,
)
})
}

Expand Down Expand Up @@ -348,13 +381,7 @@ func (f *Remote) openIndex() error {
if f.tokenFile == nil {
if err := f.openRemoteFile(
consts.TokenFileSuffix,
func(file storage.ImmutableFile) {
f.tokenFile = file
f.tokenReader = storage.NewIndexReader(
f.readLimiter, file.Name(),
file, f.indexCache.TokenRegistry,
)
},
func(file storage.ImmutableFile) { f.tokenFile = file },
); err != nil {
return err
}
Expand All @@ -363,13 +390,7 @@ func (f *Remote) openIndex() error {
if f.offsetsFile == nil {
if err := f.openRemoteFile(
consts.OffsetsFileSuffix,
func(file storage.ImmutableFile) {
f.offsetsFile = file
f.offsetsReader = storage.NewIndexReader(
f.readLimiter, file.Name(),
file, f.indexCache.OffsetsRegistry,
)
},
func(file storage.ImmutableFile) { f.offsetsFile = file },
); err != nil {
return err
}
Expand All @@ -378,13 +399,7 @@ func (f *Remote) openIndex() error {
if f.idFile == nil {
if err := f.openRemoteFile(
consts.IDFileSuffix,
func(file storage.ImmutableFile) {
f.idFile = file
f.idReader = storage.NewIndexReader(
f.readLimiter, file.Name(),
file, f.indexCache.IDRegistry,
)
},
func(file storage.ImmutableFile) { f.idFile = file },
); err != nil {
return err
}
Expand All @@ -393,13 +408,7 @@ func (f *Remote) openIndex() error {
if f.lidFile == nil {
if err := f.openRemoteFile(
consts.LIDFileSuffix,
func(file storage.ImmutableFile) {
f.lidFile = file
f.lidReader = storage.NewIndexReader(
f.readLimiter, file.Name(),
file, f.indexCache.LIDRegistry,
)
},
func(file storage.ImmutableFile) { f.lidFile = file },
); err != nil {
return err
}
Expand Down
Loading
Loading