diff --git a/cmd/index_analyzer/main.go b/cmd/index_analyzer/main.go index 5c8d08ef2..e0d1403a5 100644 --- a/cmd/index_analyzer/main.go +++ b/cmd/index_analyzer/main.go @@ -143,7 +143,6 @@ func openFrac( rl *storage.ReadLimiter, ) (indexwriter.Source, func()) { base := basePath(path) - legacy := strings.HasSuffix(path, consts.IndexFileSuffix) sealed := frac.NewSealed( base, @@ -153,7 +152,6 @@ func openFrac( nil, &frac.Config{}, noopSkipMaskProvider{}, - legacy, ) return frac.NewSealedSource(sealed), sealed.Release } diff --git a/consts/consts.go b/consts/consts.go index ff86ee5fe..3c27e644f 100644 --- a/consts/consts.go +++ b/consts/consts.go @@ -85,7 +85,8 @@ const ( // We can remove it in the future releases. IndexDelFileSuffix = ".index.del" - RemoteFractionSuffix = ".remote" + RemoteFractionSuffix = ".remote" + RemoteFractionTmpSuffix = "._remote" FracCacheFileSuffix = ".frac-cache" CompactionPlan = ".compaction-plan" diff --git a/frac/common/info.go b/frac/common/info.go index 46edae9a7..7f76c75d7 100644 --- a/frac/common/info.go +++ b/frac/common/info.go @@ -30,6 +30,7 @@ type Info struct { DocsRaw uint64 `json:"docs_raw"` // How much raw docs data is appended. MetaOnDisk uint64 `json:"meta_on_disk"` // How much compressed metadata is stored on disk. IndexOnDisk uint64 `json:"index_on_disk"` // How much compressed index data is stored on disk. + InfoOnDisk uint64 `json:"-"` ConstRegularBlockSize int `json:"const_regular_block_size"` ConstIDsPerBlock int `json:"const_ids_per_block"` diff --git a/frac/fraction_test.go b/frac/fraction_test.go index f632d360d..e160e5f3e 100644 --- a/frac/fraction_test.go +++ b/frac/fraction_test.go @@ -2881,7 +2881,6 @@ func (s *SealedLoadedFractionTestSuite) newSealedLoaded(bulks ...[]string) *frac nil, s.config, testSkipMaskProvider{}, - false, ) s.fraction = sealed @@ -2929,9 +2928,8 @@ func (s *RemoteFractionTestSuite) SetupTest() { ) s.Require().NoError(err, "s3 client setup failed") - offloaded, err := sealed.Offload(context.Background(), s3.NewUploader(s3cli)) + err = sealed.Offload(context.Background(), s3.NewUploader(s3cli)) s.Require().NoError(err, "offload failed") - s.Require().True(offloaded, "didn't offload frac") remoteFrac := frac.NewRemote( context.Background(), @@ -2943,7 +2941,6 @@ func (s *RemoteFractionTestSuite) SetupTest() { s.config, s3cli, testSkipMaskProvider{}, - false, ) s.fraction = remoteFrac diff --git a/frac/remote.go b/frac/remote.go index b0e444427..b48613b95 100644 --- a/frac/remote.go +++ b/frac/remote.go @@ -2,13 +2,16 @@ package frac import ( "context" + "errors" "fmt" + "os" "path/filepath" "sync" "go.uber.org/zap" "github.com/ozontech/seq-db/cache" + "github.com/ozontech/seq-db/config" "github.com/ozontech/seq-db/consts" "github.com/ozontech/seq-db/frac/common" "github.com/ozontech/seq-db/frac/processor" @@ -42,17 +45,14 @@ type Remote struct { docsFile storage.ImmutableFile docsCache *cache.ConcurrentCache[[]byte] - // IsLegacy is true for fractions that use the old single .index file format. - IsLegacy bool - legacyFile storage.ImmutableFile - // Per-section index files (new split format only). - infoFile storage.ImmutableFile tokenFile storage.ImmutableFile offsetsFile storage.ImmutableFile idFile storage.ImmutableFile lidFile storage.ImmutableFile + legacyFile storage.ImmutableFile + indexCache *IndexCache initMu *sync.RWMutex @@ -72,10 +72,9 @@ func NewRemote( indexCache *IndexCache, docsCache *cache.ConcurrentCache[[]byte], info *common.Info, - config *Config, + cfg *Config, s3cli *s3.Client, skipMaskProvider skipMaskProvider, - isLegacy bool, ) *Remote { f := &Remote{ ctx: ctx, @@ -88,12 +87,10 @@ func NewRemote( info: info, BaseFileName: baseFile, - Config: config, + Config: cfg, s3cli: s3cli, skipMaskProvider: skipMaskProvider, - - IsLegacy: isLegacy, } // Fast path if fraction-info cache exists AND it has valid index size. @@ -104,12 +101,11 @@ func NewRemote( return f } - // FIXME(dkharms): For now almost any availability issues with S3 will cause seq-db to panic during initialisation phase. - // I wrote a small proposal on how we can reduce impact of such events. - // https://github.com/ozontech/seq-db/issues/92 - if err := f.loadInfo(); err != nil { - logger.Error( + // FIXME(dkharms): For now almost any availability issues with S3 will cause seq-db to panic + // during initialisation phase. I wrote a small proposal on how we can reduce impact of such + // events. https://github.com/ozontech/seq-db/issues/92 + logger.Fatal( "cannot open info file: any subsequent operation will fail", zap.String("fraction", filepath.Base(f.BaseFileName)), zap.Error(err), @@ -177,7 +173,7 @@ func (f *Remote) createDataProvider(ctx context.Context) (*sealedDataProvider, e 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)), + tokenTableLoader: token.NewTableLoader(f.BaseFileName, f.Info().BinaryDataVer, f.IsSingleIndex(), &ir.Token, cache.NewSession(f.indexCache.TokenTable)), idsTable: &f.blocksData.IDsTable, idsProvider: seqids.NewProvider( @@ -193,7 +189,7 @@ func (f *Remote) createDataProvider(ctx context.Context) (*sealedDataProvider, e } func (f *Remote) indexReaders() IndexReaders { - if f.IsLegacy { + if f.IsSingleIndex() { r := storage.NewIndexReader( f.readLimiter, f.legacyFile.Name(), f.legacyFile, cache.NewSession(f.indexCache.LegacyRegistry), @@ -277,43 +273,95 @@ func (f *Remote) String() string { return fracToString(f, "remote") } +func (f *Remote) IsSingleIndex() bool { + return f.info.BinaryDataVer < config.BinaryDataV3 +} + +// loadInfo loads the remote fraction information from available sources in priority order: +// 1. Local non-empty *.remote file (offload stores info inside .remote itself). +// 2. Remote .info file on S3 (legacy but still supported). +// 3. Legacy *.index file on S3 (oldest scenario). func (f *Remote) loadInfo() error { - var err error + err := f.tryLoadInfoLocal() + if err == nil { + return nil + } - if f.IsLegacy { - if err := f.openInfoLegacy(); err != nil { - return err - } + logger.Warn( + "cannot open local info file for remote fraction, falling back to S3", + zap.String("fraction", f.BaseFileName), + zap.Error(err), + ) - legacyReader := storage.NewIndexReader( - f.readLimiter, f.legacyFile.Name(), - f.legacyFile, f.indexCache.LegacyRegistry, - ) + err = f.tryLoadInfoRemote() + if err == nil { + return nil + } - if f.info, err = loadInfoLegacy(legacyReader); err != nil { - logger.Fatal( - "error loading Info", - zap.String("fraction", f.BaseFileName), - zap.Error(err), - ) - } + logger.Warn( + "cannot open remote info file, falling back to legacy index", + zap.String("fraction", f.BaseFileName), + zap.Error(err), + ) - return nil + return f.loadInfoLegacy() +} + +// tryLoadInfoLocal attempts to load fraction information from a local non-empty +// .remote file. This is the most preferred and modern approach, where all data +// is already present on disk. An empty .remote is a legacy marker and means the +// fraction was offloaded before info was stored inside .remote. +func (f *Remote) tryLoadInfoLocal() (err error) { + var ( + file *os.File + stat os.FileInfo + ) + + if file, err = os.Open(f.BaseFileName + consts.RemoteFractionSuffix); err != nil { + return err } - if err := f.openInfo(); err != nil { + defer file.Close() + + if stat, err = file.Stat(); err != nil { return err } - if f.info, err = loadInfo(f.infoFile); err != nil { - logger.Fatal( - "error loading Info", - zap.String("fraction", f.BaseFileName), - zap.Error(err), - ) + if stat.Size() == 0 { + return errors.New("it's a legacy empty *.remote file") } - return nil + f.info, err = loadInfo(file) + return err +} + +// tryLoadInfoRemote attempts to load fraction information from a remote .info file +// located on S3. This is an intermediate fallback: it is used when the local +// .remote is empty, but an .info file still exists on S3 (maintained for +// backward compatibility). +func (f *Remote) tryLoadInfoRemote() error { + infoFile, err := f.openRemoteFile(consts.InfoFileSuffix, true) + if err == nil { + f.info, err = loadInfo(infoFile) + } + return err +} + +// loadInfoLegacy loads fraction information from the legacy index stored on S3. +// This is the oldest fallback, used when only an empty *.remote file exists locally +// and a single *.index file resides on S3 containing all necessary data. +func (f *Remote) loadInfoLegacy() (err error) { + if err := f.openIndexLegacyRemote(); err != nil { + return err + } + + reader := storage.NewIndexReader( + f.readLimiter, f.legacyFile.Name(), f.legacyFile, + cache.NewSession(f.indexCache.LegacyRegistry), + ) + + f.info, err = loadInfoLegacy(reader) + return err } func (f *Remote) init() error { @@ -332,12 +380,12 @@ func (f *Remote) init() error { return nil } - if f.IsLegacy { + if f.IsSingleIndex() { (&LegacyLoader{}).Load( &f.blocksData, f.info, storage.NewIndexReader( - f.readLimiter, f.legacyFile.Name(), - f.legacyFile, f.indexCache.LegacyRegistry, + f.readLimiter, f.legacyFile.Name(), f.legacyFile, + cache.NewSession(f.indexCache.LegacyRegistry), ), ) @@ -351,70 +399,40 @@ func (f *Remote) init() error { return nil } -func (f *Remote) openInfoLegacy() error { - if f.legacyFile != nil { - return nil +func (f *Remote) openIndexLegacyRemote() (err error) { + if f.legacyFile == nil { + f.legacyFile, err = f.openRemoteFile(consts.IndexFileSuffix, true) } - - return f.openRemoteFile(consts.IndexFileSuffix, func(file storage.ImmutableFile) { - f.legacyFile = file - }) -} - -func (f *Remote) openInfo() error { - if f.infoFile != nil { - return nil - } - - return f.openRemoteFile( - consts.InfoFileSuffix, - func(file storage.ImmutableFile) { - f.infoFile = file - }, - ) + return err } func (f *Remote) openIndex() error { - if f.IsLegacy { - return f.openInfoLegacy() + if f.IsSingleIndex() { + return f.openIndexLegacyRemote() } - if err := f.openInfo(); err != nil { - return err - } + var err error if f.tokenFile == nil { - if err := f.openRemoteFile( - consts.TokenFileSuffix, - func(file storage.ImmutableFile) { f.tokenFile = file }, - ); err != nil { + if f.tokenFile, err = f.openRemoteFile(consts.TokenFileSuffix, true); err != nil { return err } } if f.offsetsFile == nil { - if err := f.openRemoteFile( - consts.OffsetsFileSuffix, - func(file storage.ImmutableFile) { f.offsetsFile = file }, - ); err != nil { + if f.offsetsFile, err = f.openRemoteFile(consts.OffsetsFileSuffix, true); err != nil { return err } } if f.idFile == nil { - if err := f.openRemoteFile( - consts.IDFileSuffix, - func(file storage.ImmutableFile) { f.idFile = file }, - ); err != nil { + if f.idFile, err = f.openRemoteFile(consts.IDFileSuffix, true); err != nil { return err } } if f.lidFile == nil { - if err := f.openRemoteFile( - consts.LIDFileSuffix, - func(file storage.ImmutableFile) { f.lidFile = file }, - ); err != nil { + if f.lidFile, err = f.openRemoteFile(consts.LIDFileSuffix, true); err != nil { return err } } @@ -422,23 +440,25 @@ func (f *Remote) openIndex() error { return nil } -func (f *Remote) openRemoteFile(suffix string, assign func(storage.ImmutableFile)) error { +// openRemoteFile returns (nil, nil) if the file is missing and mustExist is false. +func (f *Remote) openRemoteFile(suffix string, mustExist bool) (storage.ImmutableFile, error) { name := filepath.Base(f.BaseFileName) + suffix - ok, err := f.s3cli.Exists(f.ctx, name) if err != nil { - return fmt.Errorf( + return nil, fmt.Errorf( "cannot check existence of %q file: %w", suffix, err, ) } if !ok { - return fmt.Errorf("missing %q file", suffix) + if mustExist { + return nil, fmt.Errorf("missing %q file", suffix) + } + return nil, nil } - assign(s3.NewReader(f.ctx, f.s3cli, name)) - return nil + return s3.NewReader(f.ctx, f.s3cli, name), nil } func (f *Remote) openDocs() error { @@ -446,36 +466,23 @@ func (f *Remote) openDocs() error { return nil } - sortedName := filepath.Base(f.BaseFileName) + consts.SdocsFileSuffix - unsortedName := filepath.Base(f.BaseFileName) + consts.DocsFileSuffix - - unsortedExists, err := f.s3cli.Exists(f.ctx, unsortedName) - if err != nil { - return fmt.Errorf( - "cannot check existence of %q file: %w", - consts.DocsFileSuffix, err, - ) - } - - if unsortedExists { - f.docsFile = s3.NewReader(f.ctx, f.s3cli, unsortedName) - return nil - } - - sortedExists, err := f.s3cli.Exists(f.ctx, sortedName) + docsFile, err := f.openRemoteFile(consts.DocsFileSuffix, false) if err != nil { - return fmt.Errorf( - "cannot check existence of %q file: %w", - consts.SdocsFileSuffix, err, - ) + return err } - if sortedExists { - f.docsFile = s3.NewReader(f.ctx, f.s3cli, sortedName) - return nil + if docsFile == nil { + docsFile, err = f.openRemoteFile(consts.SdocsFileSuffix, false) + if err != nil { + return err + } + if docsFile == nil { + return fmt.Errorf("missing %q and %q files", consts.DocsFileSuffix, consts.SdocsFileSuffix) + } } - return fmt.Errorf("missing %q and %q files", consts.DocsFileSuffix, consts.SdocsFileSuffix) + f.docsFile = docsFile + return nil } func (f *Remote) computeIndexSize() { @@ -487,21 +494,21 @@ func (f *Remote) computeIndexSize() { return } + f.info.IndexOnDisk = f.info.InfoOnDisk files := []storage.ImmutableFile{ - f.infoFile, f.tokenFile, f.offsetsFile, f.idFile, f.lidFile, } - if f.IsLegacy { + if f.IsSingleIndex() { + f.info.IndexOnDisk = 0 files = []storage.ImmutableFile{ f.legacyFile, } } - f.info.IndexOnDisk = 0 for _, file := range files { st, err := file.Stat() if err != nil { diff --git a/frac/sealed.go b/frac/sealed.go index 65db0fb06..6828e9aa0 100644 --- a/frac/sealed.go +++ b/frac/sealed.go @@ -6,13 +6,13 @@ import ( "fmt" "io" "os" - "path/filepath" "sync" "go.uber.org/zap" "golang.org/x/sync/errgroup" "github.com/ozontech/seq-db/cache" + "github.com/ozontech/seq-db/config" "github.com/ozontech/seq-db/consts" "github.com/ozontech/seq-db/frac/common" "github.com/ozontech/seq-db/frac/processor" @@ -38,17 +38,14 @@ type Sealed struct { docsFile *os.File docsCache *cache.ConcurrentCache[[]byte] - // IsLegacy is true for fractions that use the old single .index file format. - IsLegacy bool - legacyFile *os.File - // Per-section index files and their readers (new split format only). - infoFile *os.File tokenFile *os.File offsetsFile *os.File idFile *os.File lidFile *os.File + legacyFile *os.File + blocksData sealed.BlocksData indexCache *IndexCache @@ -77,9 +74,8 @@ func NewSealed( indexCache *IndexCache, docsCache *cache.ConcurrentCache[[]byte], info *common.Info, - config *Config, + cfg *Config, skipMaskProvider skipMaskProvider, - isLegacy bool, ) *Sealed { f := &Sealed{ initMu: &sync.RWMutex{}, @@ -88,10 +84,9 @@ func NewSealed( docsCache: docsCache, indexCache: indexCache, - IsLegacy: isLegacy, info: info, BaseFileName: baseFile, - Config: config, + Config: cfg, PartialSuicideMode: Off, @@ -115,7 +110,7 @@ func NewSealedPreloaded( rl *storage.ReadLimiter, indexCache *IndexCache, docsCache *cache.ConcurrentCache[[]byte], - config *Config, + cfg *Config, skipMaskProvider skipMaskProvider, ) *Sealed { f := &Sealed{ @@ -130,7 +125,7 @@ func NewSealedPreloaded( info: preloaded.Info, BaseFileName: baseFile, - Config: config, + Config: cfg, skipMaskProvider: skipMaskProvider, } @@ -156,68 +151,42 @@ func NewSealedPreloaded( return f } -func (f *Sealed) openInfoLegacy() { - if f.legacyFile != nil { - return - } - - f.openFile( - consts.IndexFileSuffix, - func(file *os.File) { f.legacyFile = file }, - ) +func (f *Sealed) IsSingleIndex() bool { + return f.info.BinaryDataVer < config.BinaryDataV3 } -func (f *Sealed) openInfo() { - if f.infoFile != nil { - return +func (f *Sealed) openIndexLegacy() { + if f.legacyFile == nil { + f.legacyFile = f.openFile(consts.IndexFileSuffix) } - - f.openFile( - consts.InfoFileSuffix, - func(file *os.File) { f.infoFile = file }, - ) } func (f *Sealed) openIndex() { - if f.IsLegacy { + if f.IsSingleIndex() { // We have exactly one `.index` file for legacy sealed fractions. // So opening only this file is sufficient. - f.openInfoLegacy() + f.openIndexLegacy() return } - f.openInfo() - if f.tokenFile == nil { - f.openFile( - consts.TokenFileSuffix, - func(file *os.File) { f.tokenFile = file }, - ) + f.tokenFile = f.openFile(consts.TokenFileSuffix) } if f.offsetsFile == nil { - f.openFile( - consts.OffsetsFileSuffix, - func(file *os.File) { f.offsetsFile = file }, - ) + f.offsetsFile = f.openFile(consts.OffsetsFileSuffix) } if f.idFile == nil { - f.openFile( - consts.IDFileSuffix, - func(file *os.File) { f.idFile = file }, - ) + f.idFile = f.openFile(consts.IDFileSuffix) } if f.lidFile == nil { - f.openFile( - consts.LIDFileSuffix, - func(file *os.File) { f.lidFile = file }, - ) + f.lidFile = f.openFile(consts.LIDFileSuffix) } } -func (f *Sealed) openFile(suffix string, assign func(*os.File)) { +func (f *Sealed) openFile(suffix string) *os.File { name := f.BaseFileName + suffix file, err := os.Open(name) @@ -229,7 +198,7 @@ func (f *Sealed) openFile(suffix string, assign func(*os.File)) { ) } - assign(file) + return file } func (f *Sealed) openDocs() { @@ -260,34 +229,40 @@ func (f *Sealed) openDocs() { } func (f *Sealed) loadInfo() { - var err error - - if f.IsLegacy { - f.openInfoLegacy() - - legacyReader := storage.NewIndexReader( - f.readLimiter, f.legacyFile.Name(), - f.legacyFile, f.indexCache.LegacyRegistry, + if err := f.tryLoadInfo(); err != nil { + logger.Warn( + "cannot open single info file, falling back to legacy index", + zap.String("fraction", f.BaseFileName), + zap.Error(err), ) + f.tryLoadLegacyInfo() + } +} - if f.info, err = loadInfoLegacy(legacyReader); err != nil { - logger.Fatal( - "error loading Info", - zap.String("fraction", f.BaseFileName), - zap.Error(err), - ) - } +func (f *Sealed) tryLoadInfo() error { + infoFile, err := os.Open(f.BaseFileName + consts.InfoFileSuffix) + if err != nil { + return err + } + defer infoFile.Close() - return + if f.info, err = loadInfo(infoFile); err != nil { + logger.Fatal("error loading info", zap.String("fraction", f.BaseFileName), zap.Error(err)) } + return nil +} - f.openInfo() - if f.info, err = loadInfo(f.infoFile); err != nil { - logger.Fatal( - "error loading Info", - zap.String("fraction", f.BaseFileName), - zap.Error(err), - ) +func (f *Sealed) tryLoadLegacyInfo() { + f.openIndexLegacy() + + reader := storage.NewIndexReader( + f.readLimiter, f.legacyFile.Name(), f.legacyFile, + cache.NewSession(f.indexCache.LegacyRegistry), + ) + + var err error + if f.info, err = loadInfoLegacy(reader); err != nil { + logger.Fatal("error loading legacy info", zap.String("fraction", f.BaseFileName), zap.Error(err)) } } @@ -302,12 +277,12 @@ func (f *Sealed) init(full bool) { return } - if f.IsLegacy { + if f.IsSingleIndex() { (&LegacyLoader{}).Load( &f.blocksData, f.info, storage.NewIndexReader( - f.readLimiter, f.legacyFile.Name(), - f.legacyFile, f.indexCache.LegacyRegistry, + f.readLimiter, f.legacyFile.Name(), f.legacyFile, + cache.NewSession(f.indexCache.LegacyRegistry), ), ) @@ -320,36 +295,73 @@ func (f *Sealed) init(full bool) { } // Offload saves all index files and docs to remote storage. -func (f *Sealed) Offload(ctx context.Context, u storage.Uploader) (bool, error) { +func (f *Sealed) Offload(ctx context.Context, u storage.Uploader) error { f.init(false) + if f.IsSingleIndex() { + return f.offloadLegacy(ctx, u) + } + + infoScr := f.BaseFileName + consts.InfoFileSuffix + infoDstTmp := f.BaseFileName + consts.RemoteFractionTmpSuffix + // persist frac.info -> frac._remote (durable link, survives crash). + if err := util.DurableHardLink(infoScr, infoDstTmp); err != nil { + if !errors.Is(err, os.ErrExist) { + return err + } + // frac._remote already exists - it is retry of a previous attempt, not a failure. + } + g, gctx := errgroup.WithContext(ctx) g.Go(func() error { return u.Upload(gctx, f.docsFile) }) + g.Go(func() error { return u.Upload(gctx, f.tokenFile) }) + g.Go(func() error { return u.Upload(gctx, f.offsetsFile) }) + g.Go(func() error { return u.Upload(gctx, f.idFile) }) + g.Go(func() error { return u.Upload(gctx, f.lidFile) }) - if f.IsLegacy { - g.Go(func() error { return u.Upload(gctx, f.legacyFile) }) - } else { - g.Go(func() error { return u.Upload(gctx, f.infoFile) }) - g.Go(func() error { return u.Upload(gctx, f.tokenFile) }) - g.Go(func() error { return u.Upload(gctx, f.offsetsFile) }) - g.Go(func() error { return u.Upload(gctx, f.idFile) }) - g.Go(func() error { return u.Upload(gctx, f.lidFile) }) + g.Go(func() error { + infoFile, err := os.Open(f.BaseFileName + consts.InfoFileSuffix) + if err != nil { + return err + } + defer infoFile.Close() + return u.Upload(gctx, infoFile) + }) + + if err := g.Wait(); err != nil { + // TODO: Clean S3 zombies + remove frac._remote + return err } + infoDst := f.BaseFileName + consts.RemoteFractionSuffix + if err := util.DurableRenameFile(infoDstTmp, infoDst); err != nil { // rename frac._remote -> frac.remote + return err + } + + return nil +} + +func (f *Sealed) offloadLegacy(ctx context.Context, u storage.Uploader) error { + tmp := f.BaseFileName + consts.RemoteFractionTmpSuffix + if err := util.DurableTouchFile(tmp); err != nil { // create empty frac._remote + return err + } + + g, gctx := errgroup.WithContext(ctx) + g.Go(func() error { return u.Upload(gctx, f.docsFile) }) + g.Go(func() error { return u.Upload(gctx, f.legacyFile) }) + if err := g.Wait(); err != nil { - // TODO: Clean S3 zombies - return true, err + // TODO: Clean S3 zombies + remove frac._remote + return err } - remoteFracName := f.BaseFileName + consts.RemoteFractionSuffix - file, err := os.Create(remoteFracName) - if err != nil { - return true, err + dst := f.BaseFileName + consts.RemoteFractionSuffix + if err := util.DurableRenameFile(tmp, dst); err != nil { // rename frac._remote -> frac.remote + return err } - defer file.Close() - util.MustSyncPath(filepath.Dir(remoteFracName)) - return true, nil + return nil } func (f *Sealed) Release() { @@ -357,14 +369,13 @@ func (f *Sealed) Release() { indexFiles := []*os.File{ f.docsFile, - f.infoFile, f.tokenFile, f.offsetsFile, f.idFile, f.lidFile, } - if f.IsLegacy { + if f.IsSingleIndex() { indexFiles = []*os.File{ f.docsFile, f.legacyFile, @@ -426,7 +437,7 @@ func (f *Sealed) Suicide() { consts.LIDFileSuffix, } - if f.IsLegacy { + if f.IsSingleIndex() { indexSuffixes = []string{ consts.IndexFileSuffix, } @@ -504,7 +515,7 @@ func (f *Sealed) createDataProvider(ctx context.Context) *sealedDataProvider { 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)), + tokenTableLoader: token.NewTableLoader(f.BaseFileName, f.Info().BinaryDataVer, f.IsSingleIndex(), &ir.Token, cache.NewSession(f.indexCache.TokenTable)), idsTable: &f.blocksData.IDsTable, idsProvider: seqids.NewProvider( @@ -533,7 +544,7 @@ func (f *Sealed) IsIntersecting(from, to seq.MID) bool { } func (f *Sealed) indexReaders() IndexReaders { - if f.IsLegacy { + if f.IsSingleIndex() { r := storage.NewIndexReader( f.readLimiter, f.legacyFile.Name(), f.legacyFile, cache.NewSession(f.indexCache.LegacyRegistry), @@ -573,21 +584,21 @@ func (f *Sealed) docsReader() storage.DocsReader { // computeIndexOnDisk returns the total on-disk size of index files for a local fraction. func (f *Sealed) computeIndexSize() { + f.info.IndexOnDisk = f.info.InfoOnDisk suffixes := []string{ - consts.InfoFileSuffix, consts.TokenFileSuffix, consts.OffsetsFileSuffix, consts.IDFileSuffix, consts.LIDFileSuffix, } - if f.IsLegacy { + if f.IsSingleIndex() { + f.info.IndexOnDisk = 0 suffixes = []string{ consts.IndexFileSuffix, } } - f.info.IndexOnDisk = 0 for _, suffix := range suffixes { st, err := os.Stat(f.info.Path + suffix) if err != nil { diff --git a/frac/sealed/block_info.go b/frac/sealed/block_info.go index 4aa3025fe..4cac15518 100644 --- a/frac/sealed/block_info.go +++ b/frac/sealed/block_info.go @@ -38,5 +38,7 @@ func (b *BlockInfo) Unpack(data []byte) error { return errors.New("stats unmarshaling error") } b.Info.MetaOnDisk = 0 // todo: make this correction on sealing and remove this next time + b.Info.InfoOnDisk = uint64(len(data)) + return nil } diff --git a/frac/sealed_source.go b/frac/sealed_source.go index 1f58563aa..3e0403db4 100644 --- a/frac/sealed_source.go +++ b/frac/sealed_source.go @@ -63,7 +63,7 @@ func NewSealedSource(f *Sealed) *SealedSource { s.tokenTableLoader = token.NewTableLoader( f.BaseFileName, f.Info().BinaryDataVer, - f.IsLegacy, &s.tokenReader, cache.NewScan[token.Table](), + f.IsSingleIndex(), &s.tokenReader, cache.NewScan[token.Table](), ) return s diff --git a/fracmanager/frac_manifest.go b/fracmanager/frac_manifest.go index a0458f2f1..3df92b27d 100644 --- a/fracmanager/frac_manifest.go +++ b/fracmanager/frac_manifest.go @@ -27,7 +27,10 @@ type fracManifest struct { hasWal bool // presence of WAL with meta hasIndex bool // presence of index file hasSdocs bool // presence of sorted documents - hasRemote bool // presence of remote fraction + hasRemote bool // presence of remote fraction (legacy) + + // Presence of ._remote (case when offloading was interrupted) + hasRemoteTmp bool // Split index file flags hasInfo bool @@ -46,7 +49,15 @@ type fracManifest struct { // hasAllIndexFiles reports whether all 5 split index files are present. func (m *fracManifest) hasAllIndexFiles() bool { - return m.hasInfo && m.hasToken && m.hasOffsets && m.hasID && m.hasLID + return (m.hasInfo && m.hasToken && m.hasOffsets && m.hasID && m.hasLID) || m.hasIndex +} + +func (m *fracManifest) hasDocsFile() bool { + return m.hasSdocs || m.hasDocs +} + +func (m *fracManifest) hasRemoteFile() bool { + return m.hasRemote } // AddExtension adds information about a file with the specified extension @@ -85,6 +96,9 @@ func (m *fracManifest) AddExtension(ext string) error { case consts.CompactionPlan: m.hasCompactionPlan = true + case consts.RemoteFractionTmpSuffix: + m.hasRemoteTmp = true + case consts.IndexTmpFileSuffix, consts.InfoTmpFileSuffix, consts.TokenTmpFileSuffix, consts.OffsetsTmpFileSuffix, consts.IDTmpFileSuffix, consts.LIDTmpFileSuffix, @@ -114,10 +128,10 @@ const ( // Stage determines the current stage of the fraction based on file presence // Key method for making fraction management decisions func (m *fracManifest) Stage() fracStage { - if m.hasRemote { + if m.hasRemoteFile() { return fracStageRemote } - if (m.hasAllIndexFiles() || m.hasIndex) && (m.hasSdocs || m.hasDocs) { + if m.hasAllIndexFiles() && m.hasDocsFile() { return fracStageSealed } if m.hasWal && m.hasDocs { @@ -201,6 +215,14 @@ func removeIndexTmp(m *fracManifest) { } } +func removeRemoteTmp(m *fracManifest) { + if m.hasRemoteTmp { + // TODO: Clean S3 zombies before + util.RemoveFile(m.basePath + consts.RemoteFractionTmpSuffix) + m.hasRemoteTmp = false + } +} + func removeSdocsTmp(m *fracManifest) { util.RemoveFile(m.basePath + consts.SdocsTmpFileSuffix) } @@ -391,6 +413,7 @@ func cleanupTemporary(m *fracManifest) { removeSdocsDel(m) removeDocsDel(m) removeIndexTmp(m) + removeRemoteTmp(m) removeDocsTmp(m) removeSdocsTmp(m) } @@ -403,6 +426,8 @@ func removeAllFiles(basePath string) { consts.SdocsFileSuffix, consts.SdocsDelFileSuffix, consts.SdocsTmpFileSuffix, consts.IndexFileSuffix, consts.IndexDelFileSuffix, consts.IndexTmpFileSuffix, + consts.RemoteFractionTmpSuffix, consts.RemoteFractionSuffix, + consts.InfoFileSuffix, consts.InfoTmpFileSuffix, consts.TokenFileSuffix, consts.TokenTmpFileSuffix, consts.OffsetsFileSuffix, consts.OffsetsTmpFileSuffix, @@ -451,6 +476,7 @@ func (f *fracManifest) MarshalLogObject(enc zapcore.ObjectEncoder) error { enc.AddBool("hasIndex", f.hasIndex) enc.AddBool("hasSdocs", f.hasSdocs) enc.AddBool("hasRemote", f.hasRemote) + enc.AddBool("hasRemoteTmp", f.hasRemoteTmp) enc.AddBool("hasInfo", f.hasInfo) enc.AddBool("hasToken", f.hasToken) diff --git a/fracmanager/fraction_provider.go b/fracmanager/fraction_provider.go index 995f67a41..04b495495 100644 --- a/fracmanager/fraction_provider.go +++ b/fracmanager/fraction_provider.go @@ -74,7 +74,7 @@ func (fp *fractionProvider) NewActive(name string) *frac.Active { ) } -func (fp *fractionProvider) NewSealed(name string, cachedInfo *common.Info, isLegacy bool) *frac.Sealed { +func (fp *fractionProvider) NewSealed(name string, cachedInfo *common.Info) *frac.Sealed { return frac.NewSealed( name, fp.readLimiter, @@ -83,7 +83,6 @@ func (fp *fractionProvider) NewSealed(name string, cachedInfo *common.Info, isLe cachedInfo, // Preloaded meta information &fp.config.Fraction, fp.skipMaskProvider, - isLegacy, ) } @@ -99,7 +98,7 @@ func (fp *fractionProvider) NewSealedPreloaded(name string, preloadedData *seale ) } -func (fp *fractionProvider) NewRemote(ctx context.Context, name string, cachedInfo *common.Info, isLegacy bool) *frac.Remote { +func (fp *fractionProvider) NewRemote(ctx context.Context, name string, cachedInfo *common.Info) *frac.Remote { return frac.NewRemote( ctx, name, @@ -110,7 +109,6 @@ func (fp *fractionProvider) NewRemote(ctx context.Context, name string, cachedIn &fp.config.Fraction, fp.s3cli, fp.skipMaskProvider, - isLegacy, ) } @@ -172,15 +170,11 @@ func (fp *fractionProvider) Seal(a *frac.Active) (*frac.Sealed, error) { // Offload uploads fraction to S3 storage and returns a remote fraction // IMPORTANT: context controls timeouts and operation cancellation func (fp *fractionProvider) Offload(ctx context.Context, f *frac.Sealed) (*frac.Remote, error) { - mustBeOffloaded, err := f.Offload(ctx, s3.NewUploader(fp.s3cli)) + err := f.Offload(ctx, s3.NewUploader(fp.s3cli)) if err != nil { return nil, err } - if !mustBeOffloaded { - return nil, nil - } - info := f.Info() - return fp.NewRemote(ctx, info.Path, info, f.IsLegacy), nil + return fp.NewRemote(ctx, info.Path, info), nil } diff --git a/fracmanager/lifecycle_manager.go b/fracmanager/lifecycle_manager.go index 199025ae7..2be719aca 100644 --- a/fracmanager/lifecycle_manager.go +++ b/fracmanager/lifecycle_manager.go @@ -181,10 +181,8 @@ func (lc *lifecycleManager) tryOffload(ctx context.Context, sealed *frac.Sealed) return nil, err } - if remote != nil { - offloadingTotal.WithLabelValues("success").Inc() - offloadingDurationSeconds.Observe(float64(offloadingDuration)) - } + offloadingTotal.WithLabelValues("success").Inc() + offloadingDurationSeconds.Observe(float64(offloadingDuration)) return remote, nil } diff --git a/fracmanager/loader.go b/fracmanager/loader.go index 2de97d2c4..0fbbf30ef 100644 --- a/fracmanager/loader.go +++ b/fracmanager/loader.go @@ -9,7 +9,6 @@ import ( "go.uber.org/zap" "golang.org/x/sync/errgroup" - "github.com/ozontech/seq-db/consts" "github.com/ozontech/seq-db/frac" "github.com/ozontech/seq-db/logger" ) @@ -138,23 +137,9 @@ func (l *Loader) discover(ctx context.Context) ([]*frac.Active, []*frac.Sealed, case fracStageActive: actives = append(actives, l.provider.NewActive(manifest.basePath)) case fracStageSealed: - locals = append(locals, l.loadSealed(manifest, loadedInfoCache)) + locals = append(locals, l.loadSealed(manifest.basePath, loadedInfoCache)) case fracStageRemote: - // TODO(dkharms): Drop this check once we store `Info` for remote fractions locally. - - indexName := filepath.Base(manifest.basePath) + consts.IndexFileSuffix - hasIndex, err := l.provider.s3cli.Exists(ctx, indexName) - if err != nil { - logger.Error( - "will skip fraction: cannot check existence of .index file", - zap.String("fraction", filepath.Base(manifest.basePath)), - zap.Error(err), - ) - continue - } - - manifest.hasIndex = hasIndex - remotes = append(remotes, l.loadRemote(ctx, manifest, loadedInfoCache)) + remotes = append(remotes, l.loadRemote(ctx, manifest.basePath, loadedInfoCache)) default: logger.Error("unexpected fraction stage", zap.Any("manifest", manifest)) } @@ -169,21 +154,21 @@ func (l *Loader) discover(ctx context.Context) ([]*frac.Active, []*frac.Sealed, } // loadSealed loads a sealed fraction using cache -func (l *Loader) loadSealed(manifest *fracManifest, loadedInfoCache *fracInfoCache) *frac.Sealed { - info, found := loadedInfoCache.Get(filepath.Base(manifest.basePath)) +func (l *Loader) loadSealed(basePath string, loadedInfoCache *fracInfoCache) *frac.Sealed { + info, found := loadedInfoCache.Get(filepath.Base(basePath)) l.updateStats(found) - f := l.provider.NewSealed(manifest.basePath, info, manifest.hasIndex) + f := l.provider.NewSealed(basePath, info) l.infoCache.Add(f.Info()) return f } // loadRemote loads a remote fraction -func (l *Loader) loadRemote(ctx context.Context, manifest *fracManifest, loadedInfoCache *fracInfoCache) *frac.Remote { - info, found := loadedInfoCache.Get(filepath.Base(manifest.basePath)) +func (l *Loader) loadRemote(ctx context.Context, basePath string, loadedInfoCache *fracInfoCache) *frac.Remote { + info, found := loadedInfoCache.Get(filepath.Base(basePath)) l.updateStats(found) - f := l.provider.NewRemote(ctx, manifest.basePath, info, manifest.hasIndex) + f := l.provider.NewRemote(ctx, basePath, info) l.infoCache.Add(f.Info()) return f } diff --git a/fracmanager/loader_test.go b/fracmanager/loader_test.go index 6ba8abb61..22a147d31 100644 --- a/fracmanager/loader_test.go +++ b/fracmanager/loader_test.go @@ -3,13 +3,16 @@ package fracmanager import ( "context" "math/rand" + "os" "path/filepath" "sync" "testing" "time" "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "github.com/ozontech/seq-db/config" "github.com/ozontech/seq-db/consts" "github.com/ozontech/seq-db/frac" "github.com/ozontech/seq-db/frac/common" @@ -273,3 +276,199 @@ func TestDiscover(t *testing.T) { assert.Empty(t, expectedSealed, "we don't expect any more sealed fractions") assert.Empty(t, expectedRemote, "we don't expect any more remote fractions") } + +// createEmptyRemoteFile creates an empty .remote marker file on disk. +// An empty .remote means the fraction was offloaded in the legacy format +// (before info started being stored inside .remote itself). +func createEmptyRemoteFile(t testing.TB, basePath string) { + t.Helper() + + err := os.WriteFile(basePath+consts.RemoteFractionSuffix, nil, 0o644) + require.NoError(t, err) +} + +// TestDiscover_RemoteFileWithInfo verifies that a fraction with non-empty .remote +// (created by offload, which stores info inside .remote itself) is detected +// as remote with the new split format (no S3 request needed). +// No .frac-cache +func TestDiscover_RemoteFileWithInfo(t *testing.T) { + fp, loader, tearDown := setupLoaderTest(t, nil) + defer tearDown() + + // Create a sealed fraction and offload it — this creates non-empty .remote on disk + // (with serialized info) and uploads all files to S3. + a := fp.CreateActive() + appendDocsToActive(t, a, 10) + s, err := fp.Seal(a) + require.NoError(t, err) + + r, err := fp.Offload(t.Context(), s) + require.NoError(t, err) + require.NotNil(t, r) + s.Suicide() + + // Now discover from FS. + actives, locals, remotes, err := loader.discover(t.Context()) + require.NoError(t, err) + + assert.Empty(t, actives, "no active fractions expected") + assert.Empty(t, locals, "no local fractions expected") + require.Len(t, remotes, 1, "one remote fraction expected") + + remote := remotes[0] + assert.Equal(t, r.Info().Name(), remote.Info().Name(), "remote fraction name should match") + assert.False(t, remote.IsSingleIndex(), "remote fraction with info in .remote should be non-legacy") + + stat, err := os.Stat(remote.BaseFileName + consts.RemoteFractionSuffix) + require.NoError(t, err, "file .remote must exists") + assert.Greater(t, stat.Size(), int64(0), ".remote must contain serialized info") +} + +// TestDiscover_EmptyRemote_NewIndex verifies that a fraction with empty .remote +// and no .index in S3 (but split files exist) is detected as non-legacy remote. +// No .frac-cache +func TestDiscover_EmptyRemote_NewIndex(t *testing.T) { + fp, loader, tearDown := setupLoaderTest(t, nil) + defer tearDown() + + // Create a sealed fraction and offload it — this creates real files in S3. + a := fp.CreateActive() + appendDocsToActive(t, a, 10) + s, err := fp.Seal(a) + require.NoError(t, err) + + r, err := fp.Offload(t.Context(), s) + require.NoError(t, err) + require.NotNil(t, r) + s.Suicide() + + basePath := r.BaseFileName + + // Overwrite .remote with an empty marker to simulate legacy offload format. + createEmptyRemoteFile(t, basePath) + + // Discover from FS. + actives, locals, remotes, err := loader.discover(t.Context()) + require.NoError(t, err) + + assert.Empty(t, actives, "no active fractions expected") + assert.Empty(t, locals, "no local fractions expected") + require.Len(t, remotes, 1, "one remote fraction expected") + + remote := remotes[0] + assert.Equal(t, r.Info().Name(), remote.Info().Name(), "remote fraction name should match") + assert.False(t, remote.IsSingleIndex(), "remote fraction without .index in S3 should be non-legacy") +} + +// TestDiscover_EmptyRemote_CacheLegacy verifies that a fraction with empty .remote +// and cached Info with BinaryDataVer < V3 is detected as legacy remote. +func TestDiscover_EmptyRemote_CacheLegacy(t *testing.T) { + fp, loader, tearDown := setupLoaderTest(t, nil) + defer tearDown() + + // Create a sealed fraction and offload it. + a := fp.CreateActive() + appendDocsToActive(t, a, 10) + s, err := fp.Seal(a) + require.NoError(t, err) + + r, err := fp.Offload(t.Context(), s) + require.NoError(t, err) + require.NotNil(t, r) + s.Suicide() + + basePath := r.BaseFileName + baseName := r.Info().Name() + + // Overwrite .remote with an empty marker to simulate legacy offload format. + createEmptyRemoteFile(t, basePath) + + // Add cached Info with BinaryDataVer < V3 (simulating legacy) + // and IndexOnDisk > 0 so NewRemote fast path (info.IndexOnDisk > 0) works. + cachedInfo := &common.Info{ + Path: basePath, + DocsTotal: r.Info().DocsTotal, + BinaryDataVer: config.BinaryDataV2, // < V3 — legacy + IndexOnDisk: 4096, // > 0 — enables fast path in NewRemote + } + loader.infoCache.Add(cachedInfo) + err = loader.infoCache.SyncWithDisk() + require.NoError(t, err) + + // Discover from FS. + actives, locals, remotes, err := loader.discover(t.Context()) + require.NoError(t, err) + + assert.Empty(t, actives, "no active fractions expected") + assert.Empty(t, locals, "no local fractions expected") + require.Len(t, remotes, 1, "one remote fraction expected") + + remote := remotes[0] + assert.Equal(t, baseName, remote.Info().Name(), "remote fraction name should match") + assert.True(t, remote.IsSingleIndex(), "remote fraction with cached BinaryDataVer= V3 is detected as non-legacy remote. +// Since split format is known from cache, no S3 request should be made. +func TestDiscover_EmptyRemote_CacheNew(t *testing.T) { + fp, loader, tearDown := setupLoaderTest(t, nil) + defer tearDown() + + // Create a sealed fraction and offload it. + a := fp.CreateActive() + appendDocsToActive(t, a, 10) + s, err := fp.Seal(a) + require.NoError(t, err) + + // Add cached Info and sync to disk so loadedInfoCache inside discover() picks it up. + loader.infoCache.Add(s.Info()) + err = loader.infoCache.SyncWithDisk() + require.NoError(t, err) + + // Offload and remove localy + r, err := fp.Offload(t.Context(), s) + require.NoError(t, err) + require.NotNil(t, r) + s.Suicide() + + // Overwrite .remote with an empty marker to simulate legacy offload format. + basePath := r.BaseFileName + createEmptyRemoteFile(t, basePath) + + // Discover from FS. + actives, locals, remotes, err := loader.discover(t.Context()) + require.NoError(t, err) + + assert.Empty(t, actives, "no active fractions expected") + assert.Empty(t, locals, "no local fractions expected") + require.Len(t, remotes, 1, "one remote fraction expected") + + remote := remotes[0] + assert.Equal(t, r.Info().Name(), remote.Info().Name(), "remote fraction name should match") + assert.False(t, remote.IsSingleIndex(), "remote fraction with cached BinaryDataVer>=V3 should be non-legacy") +} + +// TestLoadRemote_Legacy verifies loading a legacy remote fraction using cached Info +// with IndexOnDisk > 0 (fast path — no S3 request for info loading). +func TestLoadRemote_Legacy(t *testing.T) { + fp, loader, tearDown := setupLoaderTest(t, nil) + defer tearDown() + + baseName := "seq-db-TESTLEGACYLOAD" + basePath := filepath.Join(fp.config.DataDir, baseName) + + // Create cached Info with IndexOnDisk > 0 (legacy, fast path). + cachedInfo := &common.Info{ + Path: basePath, + BinaryDataVer: config.BinaryDataV2, + DocsTotal: 100, + IndexOnDisk: 4096, + } + loadedInfoCache := NewFracInfoCacheFromDisk(loader.infoCache.fullPath) + loadedInfoCache.Add(cachedInfo) + + remote := loader.loadRemote(t.Context(), basePath, loadedInfoCache) + require.NotNil(t, remote) + assert.True(t, remote.IsSingleIndex(), "should be legacy") +} diff --git a/util/fs.go b/util/fs.go index 91d80678f..f5243ef1a 100644 --- a/util/fs.go +++ b/util/fs.go @@ -161,3 +161,39 @@ func CopyFile(src, dst string) error { _, err = io.Copy(out, in) return err } + +func DurableHardLink(src, dst string) error { + if err := os.Link(src, dst); err != nil { + return err + } + if err := SyncPath(filepath.Dir(dst)); err != nil { + return err + } + return nil +} + +func DurableRenameFile(src, dst string) error { + if err := os.Rename(src, dst); err != nil { + return err + } + if err := SyncPath(filepath.Dir(dst)); err != nil { + return err + } + return nil +} + +func FileExists(filename string) bool { + _, err := os.Stat(filename) + return !os.IsNotExist(err) +} + +func DurableTouchFile(name string) error { + file, err := os.Create(name) + if err != nil { + return err + } + if err := file.Close(); err != nil { + return err + } + return SyncPath(filepath.Dir(name)) +}