diff --git a/pkg/storage/indexheader/binary_writer.go b/pkg/storage/indexheader/binary_writer.go index 6dead3bf39b..3a04b09c822 100644 --- a/pkg/storage/indexheader/binary_writer.go +++ b/pkg/storage/indexheader/binary_writer.go @@ -31,6 +31,10 @@ const ( // BinaryFormatV1 represents first version of index-header file. BinaryFormatV1 = 1 + // BinaryFormatV2 represents the second version of the index-header file, + // which contains only the symbols table. + BinaryFormatV2 = 2 + indexTOCLen = 6*8 + crc32.Size BinaryTOCLen = 2*8 + crc32.Size // 16 (2 x uint64) + 4 (CRC32) // HeaderLen represents number of bytes reserved of index header for header. @@ -69,7 +73,7 @@ type BinaryTOC struct { } // WriteBinary build index-header file from the pieces of index in object storage. -func WriteBinary(ctx context.Context, bkt objstore.BucketReader, id ulid.ULID, filename string) (err error) { +func WriteBinary(ctx context.Context, bkt objstore.BucketReader, id ulid.ULID, filename string, writeV2 bool) (err error) { ir, indexVersion, err := newChunkedIndexReader(ctx, bkt, id) if err != nil { return errors.Wrap(err, "new index reader") @@ -79,7 +83,7 @@ func WriteBinary(ctx context.Context, bkt objstore.BucketReader, id ulid.ULID, f // Buffer for copying and encbuffers. // This also will control the size of file writer buffer. buf := make([]byte, 32*1024) - bw, err := newBinaryWriter(tmpFilename, buf) + bw, err := newBinaryWriter(tmpFilename, buf, writeV2) if err != nil { return errors.Wrap(err, "new binary index header writer") } @@ -101,12 +105,14 @@ func WriteBinary(ctx context.Context, bkt objstore.BucketReader, id ulid.ULID, f return errors.Wrap(err, "flush") } - if err := ir.CopyPostingsOffsets(bw.PostingOffsetsWriter(), buf); err != nil { - return err - } + if !writeV2 { + if err := ir.CopyPostingsOffsets(bw.PostingOffsetsWriter(), buf); err != nil { + return err + } - if err := bw.f.Flush(); err != nil { - return errors.Wrap(err, "flush") + if err := bw.f.Flush(); err != nil { + return errors.Wrap(err, "flush") + } } if err := bw.WriteTOC(); err != nil { @@ -243,7 +249,7 @@ type binaryWriter struct { crc32 hash.Hash } -func newBinaryWriter(fn string, buf []byte) (w *binaryWriter, err error) { +func newBinaryWriter(fn string, buf []byte, writeV2 bool) (w *binaryWriter, err error) { df, err := fileutil.OpenDir(filepath.Dir(fn)) if err != nil { return nil, err @@ -274,7 +280,11 @@ func newBinaryWriter(fn string, buf []byte) (w *binaryWriter, err error) { w.buf.Reset() w.buf.PutBE32(MagicIndex) - w.buf.PutByte(BinaryFormatV1) + if writeV2 { + w.buf.PutByte(BinaryFormatV2) + } else { + w.buf.PutByte(BinaryFormatV1) + } return w, w.f.Write(w.buf.Get()) } diff --git a/pkg/storage/indexheader/header.go b/pkg/storage/indexheader/header.go index bbc426b6f3b..9c706c27dc7 100644 --- a/pkg/storage/indexheader/header.go +++ b/pkg/storage/indexheader/header.go @@ -117,21 +117,27 @@ const ( var ( errInvalidIndexHeaderSection = errors.New(fmt.Sprintf("invalid index-header section; must be one of: %s", SectionPostingsOffsetsTable)) + errInvalidIndexHeaderVersion = errors.New("bucket reader must be enabled to write v2 index-header") ) type BucketReaderConfig struct { Enabled bool `yaml:"enabled" category:"experimental"` BucketIndexSections string `yaml:"index_sections" category:"experimental"` + WriteV2IndexHeader bool `yaml:"write_v2_index_header" category:"experimental"` } func (cfg *BucketReaderConfig) RegisterFlagsWithPrefix(f *flag.FlagSet, prefix string) { f.BoolVar(&cfg.Enabled, prefix+"enabled", false, fmt.Sprintf("Enable reading TSDB index-header sections from object storage. When enabled, the configured -%s are not downloaded to local disk.", prefix+"index-sections")) f.StringVar(&cfg.BucketIndexSections, prefix+"index-sections", SectionPostingsOffsetsTable, fmt.Sprintf("Index sections to read from object storage instead of local disk. Valid sections: %s", SectionPostingsOffsetsTable)) + f.BoolVar(&cfg.WriteV2IndexHeader, prefix+"write-v2-index-header", false, "Write only symbols table to on-disk index header.") } func (cfg *BucketReaderConfig) Validate() error { if !slices.Contains([]string{SectionPostingsOffsetsTable}, cfg.BucketIndexSections) { return errInvalidIndexHeaderSection } + if cfg.WriteV2IndexHeader && !cfg.Enabled { + return errInvalidIndexHeaderVersion + } return nil } diff --git a/pkg/storage/indexheader/header_test.go b/pkg/storage/indexheader/header_test.go index 34fb59f311d..f2a43db6ae3 100644 --- a/pkg/storage/indexheader/header_test.go +++ b/pkg/storage/indexheader/header_test.go @@ -105,7 +105,7 @@ func TestReadersComparedToIndexHeader(t *testing.T) { t.Run(testBlock.version, func(t *testing.T) { id := testBlock.id indexName := filepath.Join(tmpDir, id.String(), block.IndexHeaderFilename) - require.NoError(t, WriteBinary(ctx, bkt, id, indexName)) + require.NoError(t, WriteBinary(ctx, bkt, id, indexName, false)) indexFile, err := fileutil.OpenMmapFile(filepath.Join(tmpDir, id.String(), block.IndexFilename)) require.NoError(t, err) @@ -529,7 +529,7 @@ func labelValuesTestCases(t test.TB) (tests map[string][]labelValuesTestCase, bl require.NoError(t, err) indexName := filepath.Join(tmpDir, id.String(), block.IndexHeaderFilename) - require.NoError(t, WriteBinary(ctx, bkt, id, indexName)) + require.NoError(t, WriteBinary(ctx, bkt, id, indexName, false)) indexFile, err := fileutil.OpenMmapFile(filepath.Join(tmpDir, id.String(), block.IndexFilename)) require.NoError(t, err) @@ -602,7 +602,7 @@ func BenchmarkBinaryWrite(t *testing.B) { t.ResetTimer() for i := 0; i < t.N; i++ { - require.NoError(t, WriteBinary(ctx, bkt, m.ULID, fn)) + require.NoError(t, WriteBinary(ctx, bkt, m.ULID, fn, false)) } } diff --git a/pkg/storage/indexheader/lazy_binary_reader.go b/pkg/storage/indexheader/lazy_binary_reader.go index edf76ea1696..cc523b661fb 100644 --- a/pkg/storage/indexheader/lazy_binary_reader.go +++ b/pkg/storage/indexheader/lazy_binary_reader.go @@ -204,7 +204,7 @@ func ensureIndexHeaderOnDisk( level.Debug(logger).Log("msg", "index-header does not exist on disk; will build from bucket", "path", indexHeaderPath) start := time.Now() - if err := WriteBinary(ctx, bkt, blockID, indexHeaderPath); err != nil { + if err := WriteBinary(ctx, bkt, blockID, indexHeaderPath, false); err != nil { level.Error(logger).Log("msg", "failed to create index-header", "err", err) return err } diff --git a/pkg/storage/indexheader/lazy_binary_reader_test.go b/pkg/storage/indexheader/lazy_binary_reader_test.go index 31958fb8a27..f653164e2b7 100644 --- a/pkg/storage/indexheader/lazy_binary_reader_test.go +++ b/pkg/storage/indexheader/lazy_binary_reader_test.go @@ -694,7 +694,7 @@ func BenchmarkLazyBinaryReader_LoadReader(b *testing.B) { require.NoError(b, err) indexName := filepath.Join(bucketDir, idIndexV2.String(), block.IndexHeaderFilename) - require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName)) + require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName, false)) diskReaderBenchFactory := func( cachingBucket *bucketcache.CachingBucket, diff --git a/pkg/storage/indexheader/reader_benchmarks_test.go b/pkg/storage/indexheader/reader_benchmarks_test.go index 4e1bd5dff53..9405f9164ba 100644 --- a/pkg/storage/indexheader/reader_benchmarks_test.go +++ b/pkg/storage/indexheader/reader_benchmarks_test.go @@ -45,7 +45,7 @@ func BenchmarkLookupSymbol(b *testing.B) { require.NoError(b, err) indexName := filepath.Join(bucketDir, idIndexV2.String(), block.IndexHeaderFilename) - require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName)) + require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName, false)) // TODO: are these sensible values for parallelism? for _, parallelism := range []int{1, 2, 4, 8, 20, 100} { @@ -127,7 +127,7 @@ func BenchmarkLabelNames(b *testing.B) { require.NoError(b, err) indexName := filepath.Join(bucketDir, idIndexV2.String(), block.IndexHeaderFilename) - require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName)) + require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName, false)) binaryReader, err := NewStreamBinaryReader(ctx, idIndexV2, objstore.WithNoopInstr(bkt), blockDir, Config{}, 32, log.NewNopLogger(), NewStreamBinaryReaderMetrics(nil)) require.NoError(b, err) @@ -170,7 +170,7 @@ func BenchmarkLabelValuesOffsetsIndexV2(b *testing.B) { require.NoError(b, err) indexName := filepath.Join(dir, blockID.String(), block.IndexHeaderFilename) - require.NoError(b, WriteBinary(ctx, bkt, blockID, indexName)) + require.NoError(b, WriteBinary(ctx, bkt, blockID, indexName, false)) bucketReg := prometheus.NewPedanticRegistry() @@ -355,7 +355,7 @@ func BenchmarkPostingsOffset(b *testing.B) { require.NoError(b, err) indexName := filepath.Join(dir, idIndexV2.String(), block.IndexHeaderFilename) - require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName)) + require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName, false)) b.Run(fmt.Sprintf("%vNames%vValues", nameCount, valueCount), func(b *testing.B) { binaryReader, err := NewStreamBinaryReader(ctx, idIndexV2, objstore.WithNoopInstr(bkt), dir, Config{}, 32, log.NewNopLogger(), NewStreamBinaryReaderMetrics(nil)) @@ -401,7 +401,7 @@ func BenchmarkNewStreamBinaryReader(b *testing.B) { require.NoError(b, err) indexName := filepath.Join(bucketDir, idIndexV2.String(), block.IndexHeaderFilename) - require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName)) + require.NoError(b, WriteBinary(ctx, bkt, idIndexV2, indexName, false)) b.Run(fmt.Sprintf("%vNames%vValues", nameCount, valueCount), func(b *testing.B) { for i := 0; i < b.N; i++ { diff --git a/pkg/storage/indexheader/sparse_header.go b/pkg/storage/indexheader/sparse_header.go index e629216776e..dbdf6d84d8f 100644 --- a/pkg/storage/indexheader/sparse_header.go +++ b/pkg/storage/indexheader/sparse_header.go @@ -281,7 +281,7 @@ func BuildAndWriteSparseHeaderFromTSDBIndex( } allSymbolsCount, sparseSymbolsOffsets, sparsePostingsOffsets, err := buildInMemorySparseHeaderFromIndexHeader( - ctx, indexTOC, filePoolDecbufFactory, sparseSampleFactor, false, l, + ctx, sectionSource{indexTOC, filePoolDecbufFactory}, sectionSource{indexTOC, filePoolDecbufFactory}, sparseSampleFactor, false, l, ) if err != nil { return fmt.Errorf("cannot build sparse index-header values from full index: %w", err) @@ -299,10 +299,15 @@ func BuildAndWriteSparseHeaderFromTSDBIndex( return nil } +type sectionSource struct { + toc *TOCCompat + decbufFactory streamencoding.DecbufFactory +} + func buildInMemorySparseHeaderFromIndexHeader( ctx context.Context, - toc *TOCCompat, - decbufFactory streamencoding.DecbufFactory, + symbols sectionSource, + postingsOffsets sectionSource, sparseSampleFactor int, doChecksum bool, l log.Logger, @@ -321,13 +326,13 @@ func buildInMemorySparseHeaderFromIndexHeader( level.Info(l).Log("msg", "creating sparse index-header from full index-header") allSymbolsCount, sparseSymbolsOffsets, err = streamindex.SparseValuesFromSymbolsTable( - ctx, decbufFactory, int(toc.Symbols), doChecksum, + ctx, symbols.decbufFactory, int(symbols.toc.Symbols), doChecksum, ) if err != nil { return -1, nil, nil, err } - sparsePostingsOffsets, err = streamindex.SparseValuesFromPostingsOffsetsTable(ctx, decbufFactory, int(toc.PostingsOffsetTable), toc.PostingsListEnd, sparseSampleFactor, doChecksum) + sparsePostingsOffsets, err = streamindex.SparseValuesFromPostingsOffsetsTable(ctx, postingsOffsets.decbufFactory, int(postingsOffsets.toc.PostingsOffsetTable), postingsOffsets.toc.PostingsListEnd, sparseSampleFactor, doChecksum) if err != nil { return -1, nil, nil, err } diff --git a/pkg/storage/indexheader/stream_binary_reader.go b/pkg/storage/indexheader/stream_binary_reader.go index e468ed50926..42c79add544 100644 --- a/pkg/storage/indexheader/stream_binary_reader.go +++ b/pkg/storage/indexheader/stream_binary_reader.go @@ -52,6 +52,7 @@ func NewStreamBinaryReaderMetrics(reg prometheus.Registerer) *StreamBinaryReader // 2. Reading only the Symbols table from the v1 index-header file format on disk, // and reading the Postings Offset table directly from the full block index in the bucket. type StreamBinaryReader struct { + indexHeaderVersion int symbolsTOC *TOCCompat postingsOffsetsTOC *TOCCompat @@ -133,7 +134,7 @@ func NewStreamBinaryReader( "path", localIndexHeaderPath, "err", err, ) start := time.Now() - if err = WriteBinary(ctx, bkt, blockID, localIndexHeaderPath); err != nil { + if err = WriteBinary(ctx, bkt, blockID, localIndexHeaderPath, cfg.BucketReader.WriteV2IndexHeader); err != nil { return nil, fmt.Errorf("failed to write index header: %w", err) } level.Info(spanLog).Log( @@ -142,65 +143,40 @@ func NewStreamBinaryReader( ) } - // With full index-header now on disk, initialize the local-disk-backed decbuf factory and read the TOC. + // Initialize the local-disk-backed decbuf factory and read the TOC. // The index-header's TOC is also required to build the sparse index-header if one does not already exist. filePoolDecbufFactory := streamencoding.NewFilePoolDecbufFactory( localIndexHeaderPath, cfg.MaxIdleFileHandles, metrics.filePool, ) - indexHeaderTOC, _, err := TOCFromIndexHeader(ctx, castagnoliTable, filePoolDecbufFactory, l) - if err != nil { - // TOC read checks CRC32; assume a failure here is file corruption and attempt to recreate the index-header. + indexHeaderTOC, indexHeaderVersion, err := TOCFromIndexHeader(ctx, castagnoliTable, filePoolDecbufFactory, l) + // If we can't read the index header, or it's an unsupported version for the current config, we need to rebuild it + needPostingsOffsetsOnDisk := !cfg.BucketReader.Enabled + if err != nil || (indexHeaderVersion == BinaryFormatV2 && needPostingsOffsetsOnDisk) { + // TOC read checks CRC32; assume a failure here is either due to a file corruption. level.Debug(spanLog).Log( "msg", "failed to read table of contents from index-header on disk; will recreate from bucket block index", - "path", localIndexHeaderPath, "err", err, + "path", localIndexHeaderPath, "indexHeaderVersion", indexHeaderVersion, "err", err, ) start := time.Now() - if err = WriteBinary(ctx, bkt, blockID, localIndexHeaderPath); err != nil { + if err = WriteBinary(ctx, bkt, blockID, localIndexHeaderPath, cfg.BucketReader.WriteV2IndexHeader); err != nil { return nil, fmt.Errorf("failed to write index header: %w", err) } level.Info(spanLog).Log( "msg", "created index-header on local disk from bucket block index", "path", localIndexHeaderPath, "elapsed", time.Since(start), ) - indexHeaderTOC, _, err = TOCFromIndexHeader(ctx, castagnoliTable, filePoolDecbufFactory, l) - if err != nil { - // Failure after recreating index-header from bucket; assume this is unrecoverable. - return nil, fmt.Errorf("failed to read table of contents from index-header on disk after recreate from bucket block index: %w", err) - } - } - - // Full index-header is now on disk. - // If we previously failed to load the sparse index-header, build it now from full header. - if !sparseHeaderLoaded { - start := time.Now() - allSymbolsCount, sparseSymbolsOffsets, sparsePostingsOffsets, err = buildInMemorySparseHeaderFromIndexHeader( - ctx, indexHeaderTOC, filePoolDecbufFactory, sparseSampleFactor, cfg.VerifyOnLoad, l, - ) - if err != nil { - // Exhausted all options to load sparse index-header to memory. Not recoverable. - return nil, fmt.Errorf("cannot build sparse index-header values from full index-header: %w", err) - } - - level.Info(spanLog).Log("msg", "built sparse index-header values from full index-header", - "elapsed", time.Since(start), - ) - // Try to write to disk so we do not have to repeat this all again. - sparseHeaderProto := &indexheaderpb.Sparse{ - Symbols: streamindex.SparseSymbolsToProto(allSymbolsCount, sparseSymbolsOffsets), - PostingsOffsetTable: streamindex.SparsePostingsOffsetsTableToProto(sparsePostingsOffsets, sparseSampleFactor), - } - if err = writeSparseHeaderProtoToDisk(localSparseHeaderPath, sparseHeaderProto, l); err != nil { - // Log an error in case there are disk issues, but we can still continue. - level.Error(spanLog).Log( - "msg", "failed to write bucket sparse index-header to disk", "err", err, - ) - } + indexHeaderTOC, indexHeaderVersion, err = TOCFromIndexHeader(ctx, castagnoliTable, filePoolDecbufFactory, l) + } + if err != nil { + // Failure after recreating index-header from bucket; assume this is unrecoverable. + return nil, fmt.Errorf("failed to read table of contents from index-header on disk after recreate from bucket block index: %w", err) } // Everything is now loaded from bucket or disk. streamBinaryReader := &StreamBinaryReader{ sparseSampleFactor: sparseSampleFactor, + indexHeaderVersion: indexHeaderVersion, } // Set up each of the Symbols table and Postings Offsets table readers @@ -236,6 +212,39 @@ func NewStreamBinaryReader( streamBinaryReader.postingsOffsetsTOC = indexHeaderTOC } + // Required index-header section(s) are now on disk. + // If we previously failed to load the sparse index-header, build it now from full header. + // If the bucket reader is enabled, the postings offsets sparse index-header is built from the index header in the bucket. + if !sparseHeaderLoaded { + start := time.Now() + allSymbolsCount, sparseSymbolsOffsets, sparsePostingsOffsets, err = buildInMemorySparseHeaderFromIndexHeader( + ctx, + sectionSource{streamBinaryReader.symbolsTOC, streamBinaryReader.symbolsDecbufFactory}, + sectionSource{streamBinaryReader.postingsOffsetsTOC, streamBinaryReader.postingsOffsetsDecbufFactory}, + sparseSampleFactor, cfg.VerifyOnLoad, l, + ) + if err != nil { + // Exhausted all options to load sparse index-header to memory. Not recoverable. + return nil, fmt.Errorf("cannot build sparse index-header values from full index-header: %w", err) + } + + level.Info(spanLog).Log("msg", "built sparse index-header values from full index-header", + "elapsed", time.Since(start), + ) + + // Try to write to disk so we do not have to repeat this all again. + sparseHeaderProto := &indexheaderpb.Sparse{ + Symbols: streamindex.SparseSymbolsToProto(allSymbolsCount, sparseSymbolsOffsets), + PostingsOffsetTable: streamindex.SparsePostingsOffsetsTableToProto(sparsePostingsOffsets, sparseSampleFactor), + } + if err = writeSparseHeaderProtoToDisk(localSparseHeaderPath, sparseHeaderProto, l); err != nil { + // Log an error in case there are disk issues, but we can still continue. + level.Error(spanLog).Log( + "msg", "failed to write bucket sparse index-header to disk", "err", err, + ) + } + } + // DecbufFactory and TOC for each section are now assigned according to their configured sources. // Finally, initialize the readers for each index-header section. streamBinaryReader.postingsOffsetTable, err = streamindex.NewPostingsOffsetsTableReader( @@ -278,7 +287,7 @@ func (r *StreamBinaryReader) IndexVersion(context.Context) (int, error) { } func (r *StreamBinaryReader) IndexHeaderVersion() int { - return BinaryFormatV1 + return r.indexHeaderVersion } func (r *StreamBinaryReader) PostingsOffset(ctx context.Context, name string, value string) (rng index.Range, returnErr error) { diff --git a/pkg/storage/indexheader/stream_binary_reader_test.go b/pkg/storage/indexheader/stream_binary_reader_test.go index dc7f61371b4..9dca06d0856 100644 --- a/pkg/storage/indexheader/stream_binary_reader_test.go +++ b/pkg/storage/indexheader/stream_binary_reader_test.go @@ -12,6 +12,7 @@ import ( "testing" "github.com/go-kit/log" + "github.com/oklog/ulid/v2" "github.com/prometheus/client_golang/prometheus" "github.com/prometheus/client_golang/prometheus/testutil" "github.com/prometheus/prometheus/model/labels" @@ -100,15 +101,45 @@ func TestStreamBinaryReader_CheckSparseHeadersCorrectnessExtensive(t *testing.T) b := realByteSlice(indexFile.Bytes()) // Write sparse index headers to disk on first build. - _, err = NewStreamBinaryReader(ctx, blockID, bkt, tmpDir, Config{}, 3, log.NewNopLogger(), NewStreamBinaryReaderMetrics(nil)) + r1, err := NewStreamBinaryReader(ctx, blockID, bkt, tmpDir, Config{}, 3, log.NewNopLogger(), NewStreamBinaryReaderMetrics(nil)) require.NoError(t, err) + requireCleanup(t, r1.Close) // Read sparse index headers to disk on second build. r2, err := NewStreamBinaryReader(ctx, blockID, bkt, tmpDir, Config{}, 3, log.NewNopLogger(), NewStreamBinaryReaderMetrics(nil)) require.NoError(t, err) + requireCleanup(t, r2.Close) // Check correctness of sparse index headers. compareIndexToHeader(t, b, r2) compareIndexToHeaderPostings(t, b, r2) + require.False(t, r2.postingsOffsetTable.IsRemote(), + "postings offsets should be read from local disk when the bucket reader is disabled") + + // Build the sparse index-header by reading the postings offset table from object storage, + // with only the symbols table kept on local disk. + bucketDir := filepath.Join(tmpDir, "bucket-reader") + bucketCfg := Config{BucketReader: BucketReaderConfig{ + Enabled: true, + BucketIndexSections: SectionPostingsOffsetsTable, + WriteV2IndexHeader: true, + }} + require.NoError(t, bucketCfg.Validate()) + r3, err := NewStreamBinaryReader(ctx, blockID, bkt, bucketDir, bucketCfg, 3, log.NewNopLogger(), NewStreamBinaryReaderMetrics(nil)) + require.NoError(t, err) + requireCleanup(t, r3.Close) + require.True(t, r3.postingsOffsetTable.IsRemote(), + "postings offsets should be read from object storage when the bucket reader is enabled") + + // Check correctness of sparse index headers built from the bucket. + compareIndexToHeader(t, b, r3) + compareIndexToHeaderPostings(t, b, r3) + + // The v2 index-header holds the symbols table only, and should be smaller than the full one written by r1. + v1Header := readIndexHeaderFromDisk(t, tmpDir, blockID) + v2Header := readIndexHeaderFromDisk(t, bucketDir, blockID) + require.Equal(t, byte(BinaryFormatV1), v1Header[4]) + require.Equal(t, byte(BinaryFormatV2), v2Header[4]) + require.Less(t, len(v2Header), len(v1Header)) }) } } @@ -285,6 +316,133 @@ func TestStreamBinaryReader_UsesSparseHeaderFromObjectStore(t *testing.T) { require.ElementsMatch(t, []string{"a", "b"}, labelNames) } +// TestStreamBinaryReader_IndexHeaderVersionOnDisk tests which index-header format StreamBinaryReader +// ends up with on local disk, for each combination of V2 writing being enabled and of which format +// (if any) an earlier configuration already left on disk. +func TestStreamBinaryReader_IndexHeaderVersionOnDisk(t *testing.T) { + const samplingRate = 3 + + ctx := context.Background() + logger := log.NewNopLogger() + + tmpDir := t.TempDir() + + ubkt, err := filesystem.NewBucket(filepath.Join(tmpDir, "bkt")) + require.NoError(t, err) + bkt := objstore.WithNoopInstr(ubkt) + + t.Cleanup(func() { + require.NoError(t, bkt.Close()) + require.NoError(t, ubkt.Close()) + }) + + blockID, err := block.CreateBlock( + ctx, tmpDir, + generateLabels(generateSymbols("name", 5), generateSymbols("value", 50)), + 100, 0, 1000, labels.FromStrings("ext1", "1"), + ) + require.NoError(t, err) + _, err = block.Upload(ctx, logger, bkt, filepath.Join(tmpDir, blockID.String()), nil) + require.NoError(t, err) + + indexFile, err := fileutil.OpenMmapFile(filepath.Join(tmpDir, blockID.String(), block.IndexFilename)) + require.NoError(t, err) + requireCleanup(t, indexFile.Close) + indexBytes := realByteSlice(indexFile.Bytes()) + + // Writing V2 is only coherent alongside the bucket reader, since the postings offsets it omits + // have to come from somewhere; config validation rejects the other combination. + writeV2Cfg := Config{BucketReader: BucketReaderConfig{ + Enabled: true, + BucketIndexSections: SectionPostingsOffsetsTable, + WriteV2IndexHeader: true, + }} + require.NoError(t, writeV2Cfg.Validate()) + + for _, tc := range []struct { + name string + extantIndexHeaderVersion string + cfg Config + expectVersion int + expectRemote bool + }{ + { + name: "bucket reader enabled, write-v2 disabled", extantIndexHeaderVersion: "", + cfg: Config{BucketReader: BucketReaderConfig{Enabled: true, BucketIndexSections: SectionPostingsOffsetsTable, WriteV2IndexHeader: false}}, + expectVersion: BinaryFormatV1, expectRemote: true, + }, + { + name: "write-v2 disabled, nothing on disk", extantIndexHeaderVersion: "", cfg: Config{}, + expectVersion: BinaryFormatV1, expectRemote: false, + }, + { + name: "write-v2 enabled, nothing on disk", extantIndexHeaderVersion: "", cfg: writeV2Cfg, + expectVersion: BinaryFormatV2, expectRemote: true, + }, + { + name: "write-v2 enabled, v1 on disk", extantIndexHeaderVersion: "v1", cfg: writeV2Cfg, + expectVersion: BinaryFormatV1, expectRemote: true, + }, + { + name: "vwrite-v2 enabled, v2 on disk", extantIndexHeaderVersion: "v2", cfg: writeV2Cfg, + expectVersion: BinaryFormatV2, expectRemote: true, + }, + + { + name: "write-v2 disabled, v1 on disk", extantIndexHeaderVersion: "v1", cfg: Config{}, + expectVersion: BinaryFormatV1, expectRemote: false, + }, + { + name: "write-v2 disabled, v2 on disk", extantIndexHeaderVersion: "v2", cfg: Config{}, + expectVersion: BinaryFormatV1, expectRemote: false, + }, + } { + t.Run(tc.name, func(t *testing.T) { + readerDir := filepath.Join(tmpDir, tc.name) + + if tc.extantIndexHeaderVersion != "" { + seedIndexHeaderOnDisk(t, ctx, bkt, blockID, readerDir, tc.extantIndexHeaderVersion == "v2") + } + + reader, err := NewStreamBinaryReader(ctx, blockID, bkt, readerDir, tc.cfg, samplingRate, logger, NewStreamBinaryReaderMetrics(nil)) + require.NoError(t, err) + requireCleanup(t, reader.Close) + + require.Equal(t, tc.expectVersion, reader.IndexHeaderVersion()) + // Assert the format on disk too, not just what the reader reports, + // so this still fails if the reader and the file ever disagree. + require.Equal(t, byte(tc.expectVersion), readIndexHeaderFromDisk(t, readerDir, blockID)[4]) + + require.Equal(t, tc.expectRemote, reader.postingsOffsetTable.IsRemote()) + + // The reader must resolve symbols, label values, and postings correctly against the block index. + compareIndexToHeader(t, indexBytes, reader) + compareIndexToHeaderPostings(t, indexBytes, reader) + }) + } +} + +// seedIndexHeaderOnDisk writes an index-header of the given format under dir, standing in for one +// left behind by an earlier configuration. +func seedIndexHeaderOnDisk(t *testing.T, ctx context.Context, bkt objstore.InstrumentedBucketReader, blockID ulid.ULID, dir string, writeV2 bool) { + t.Helper() + + blockDir := filepath.Join(dir, blockID.String()) + require.NoError(t, os.MkdirAll(blockDir, os.ModePerm)) + require.NoError(t, WriteBinary(ctx, bkt, blockID, filepath.Join(blockDir, block.IndexHeaderFilename), writeV2)) +} + +// readIndexHeaderFromDisk reads the raw index-header bytes a StreamBinaryReader wrote under dir. +func readIndexHeaderFromDisk(t *testing.T, dir string, blockID ulid.ULID) []byte { + t.Helper() + + raw, err := os.ReadFile(filepath.Join(dir, blockID.String(), block.IndexHeaderFilename)) + require.NoError(t, err) + require.Greater(t, len(raw), HeaderLen, "index-header on disk is too short to contain its header") + + return raw +} + // trackedBucket wraps a BucketReader and tracks details about downloaded files type trackedBucket struct { objstore.InstrumentedBucketReader diff --git a/pkg/storage/indexheader/toc.go b/pkg/storage/indexheader/toc.go index fae5e76dfff..62e899b5fb5 100644 --- a/pkg/storage/indexheader/toc.go +++ b/pkg/storage/indexheader/toc.go @@ -26,14 +26,13 @@ import ( // TOCCompat unifies the Prometheus TSDB index TOC values available from different index types, // containing only the TOC offsets required for index-header reads of the Symbols and Postings Offsets. // -// The StreamBinaryReader can use either file-backed or bucket-backed DecbufFactory to read the index-header. +// TOCCompat can be built from one of: +// - Prometheus TSDB block index in object storage +// - Prometheus TSDB block index on disk +// - a Mimir index-header on disk (either BinaryFormatV1 or BinaryFormatV2) // -// The FilePoolDecbufFactory reads the index-header BinaryFormatV1 from disk. -// This index-header format differs from the full block index, as it only contains the Symbols and Postings Offsets: -// - The section offsets differ from a full Prometheus TSDB index file -// - The file metadata does not contain a full Prometheus block TOC, since not all index sections are present. -// -// The BucketDecbufFactory loads the full Prometheus TSDB index TOC from the block in object storage. +// Because offsets will differ depending on the source of information (Prometheus TSDB index file vs Mimir block index-header), +// the relevant DecbufFactory gets bundled with TOCCompat in sectionSource when resolving the offsets referenced by TOCCompat. type TOCCompat struct { IndexVersion int @@ -47,7 +46,9 @@ type TOCCompat struct { // in which the end of the Postings list is the beginning of the Label Indices table. // Prometheus block index TOC only contains start offsets for sections, not end offsets, // so we use Label Indices Table offset as the end bound of the Postings List. - PostingsListEnd uint64 + PostingsListEnd uint64 + + // If PostingsOffsetTable is 0, the index or index-header this TOCCompat represents should not be used to resolve postings offsets PostingsOffsetTable uint64 } @@ -88,19 +89,19 @@ func TOCFromBucketTSDBIndex( postingsListEnd := tsdbTOC.LabelIndicesTable return &TOCCompat{ - IndexVersion: indexVersion, - Symbols: tsdbTOC.Symbols, - + IndexVersion: indexVersion, + Symbols: tsdbTOC.Symbols, PostingsListEnd: postingsListEnd, PostingsOffsetTable: tsdbTOC.PostingsTable, }, nil } -// TOCFromIndexHeader builds a TOCCompat from the on-disk Mimir BinaryFormatV1. -// This format currently only exists on-disk in the store-gateways. +// TOCFromIndexHeader builds a TOCCompat from the on-disk Mimir BinaryFormatV1 or BinaryFormatV2. +// These formats currently only exists on-disk in the store-gateways. // The BinaryFormatV1 only has two main sections, which are copies of the Symbols and PostingsOffsets tables, // and it has a different layout for the header/metadata and TOC. // This results in different offsets for the relevant sections than a full Prometheus block index in the bucket. +// BinaryFormatV2 does not copy the postings offsets table. func TOCFromIndexHeader( ctx context.Context, castagnoliTable *crc32.Table, @@ -125,43 +126,49 @@ func TOCFromIndexHeader( level.Debug(l).Log("msg", "index header file size", "bytes", indexHeaderSize) indexHeaderVersion = int(decbuf.Byte()) - if indexHeaderVersion != BinaryFormatV1 { - return nil, 0, fmt.Errorf("unknown or unsupported index header format version %d", indexHeaderVersion) - } - - indexVersion := int(decbuf.Byte()) - if indexVersion != index.FormatV2 { - return nil, 0, fmt.Errorf("unknown or unsupported index format version %d", indexVersion) - } - - postingsListEnd := decbuf.Be64() - if err = decbuf.Err(); err != nil { - return nil, 0, fmt.Errorf("cannot read version and index version: %w", err) + switch indexHeaderVersion { + case BinaryFormatV1, BinaryFormatV2: + indexVersion := int(decbuf.Byte()) + if indexVersion != index.FormatV2 { + return nil, indexHeaderVersion, fmt.Errorf("unknown or unsupported index format version %d", indexVersion) + } + + postingsListEnd := decbuf.Be64() + if err = decbuf.Err(); err != nil { + return nil, indexHeaderVersion, fmt.Errorf("cannot read version and index version: %w", err) + } + + indexHeaderTOCOffset := indexHeaderSize - BinaryTOCLen + if decbuf.ResetAt(indexHeaderTOCOffset); decbuf.Err() != nil { + return nil, indexHeaderVersion, decbuf.Err() + } + + if decbuf.CheckCrc32(castagnoliTable); decbuf.Err() != nil { + return nil, indexHeaderVersion, decbuf.Err() + } + decbuf.ResetAt(indexHeaderTOCOffset) + symbols := decbuf.Be64() + postingsOffsetTable := decbuf.Be64() + + if err := decbuf.Err(); err != nil { + return nil, indexHeaderVersion, err + } + + // Index-header version of 2 must have a zero-offset postings offset table, + // all other versions must have a non-zero offset postings offset table. + if (indexHeaderVersion == BinaryFormatV2) != (postingsOffsetTable == 0) { + return nil, indexHeaderVersion, fmt.Errorf("index-header format version %d has offset %d for postings offsets table", indexHeaderVersion, postingsOffsetTable) + } + + return &TOCCompat{ + IndexVersion: indexVersion, + Symbols: symbols, + PostingsListEnd: postingsListEnd, + PostingsOffsetTable: postingsOffsetTable, + }, indexHeaderVersion, nil + default: + return nil, indexHeaderVersion, fmt.Errorf("unknown or unsupported index header format version %d", indexHeaderVersion) } - - indexHeaderTOCOffset := indexHeaderSize - BinaryTOCLen - if decbuf.ResetAt(indexHeaderTOCOffset); decbuf.Err() != nil { - return nil, 0, decbuf.Err() - } - - if decbuf.CheckCrc32(castagnoliTable); decbuf.Err() != nil { - return nil, 0, decbuf.Err() - } - decbuf.ResetAt(indexHeaderTOCOffset) - symbols := decbuf.Be64() - postingsOffsetTable := decbuf.Be64() - - if err := decbuf.Err(); err != nil { - return nil, 0, err - } - - return &TOCCompat{ - IndexVersion: indexVersion, - Symbols: symbols, - - PostingsListEnd: postingsListEnd, - PostingsOffsetTable: postingsOffsetTable, - }, indexHeaderVersion, nil } func fetchRange(ctx context.Context, bkt objstore.BucketReader, objectPath string, offset, length int64) (data []byte, err error) {