From 70c7c84b86c7d8c463fdc15a5075563f899791e5 Mon Sep 17 00:00:00 2001 From: Daniil Lyubaev Date: Mon, 5 Oct 2026 16:21:33 +0300 Subject: [PATCH 1/7] chore: decoding old fraction formats tests --- cmd/fraction/decoding_suites_test.go | 187 ++++++++++++++++ .../{decoding_test.go => sealing_test.go} | 208 +++++------------- cmd/fraction/testdata/legacy/seal-fraction.sh | 134 +++++++++++ cmd/fraction/testdata/legacy/sealer/main.go | 139 ++++++++++++ 4 files changed, 510 insertions(+), 158 deletions(-) create mode 100644 cmd/fraction/decoding_suites_test.go rename cmd/fraction/{decoding_test.go => sealing_test.go} (51%) create mode 100644 cmd/fraction/testdata/legacy/seal-fraction.sh create mode 100644 cmd/fraction/testdata/legacy/sealer/main.go diff --git a/cmd/fraction/decoding_suites_test.go b/cmd/fraction/decoding_suites_test.go new file mode 100644 index 00000000..06deeef5 --- /dev/null +++ b/cmd/fraction/decoding_suites_test.go @@ -0,0 +1,187 @@ +package main + +// The suites run the same decode checks against fractions from different +// sources: sealed on the fly by the current code (default) or sealed by +// the era's own code for every legacy format version ("legacy", +// see testdata/legacy/seal-fraction.sh). + +import ( + "fmt" + "os/exec" + "path/filepath" + "strings" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "github.com/stretchr/testify/suite" + + "github.com/ozontech/seq-db/config" +) + +// FractionDecoderTestSuite holds the decode checks. The source of the +// fraction is decided by the embedded suite: the current-version suite +// seals it with the local code, the legacy one seals it per version with +// the era's code via testdata/legacy/seal-fraction.sh. +type FractionDecoderTestSuite struct { + suite.Suite + + fracName string + expectedVer config.BinaryDataVersion +} + +func (s *FractionDecoderTestSuite) sealDocs(docsPerDocBlock int) { + s.fracName = filepath.Join(s.T().TempDir(), "frac_000001") + err := sealFraction( + s.fracName, + "", + strings.NewReader(strings.Join(docLines(s.T(), testDocs), "\n")), + docsPerDocBlock, + ) + s.Require().NoError(err) + s.expectedVer = config.CurrentFracVersion +} + +func (s *FractionDecoderTestSuite) sealDocsEra(v config.BinaryDataVersion) { + s.fracName = filepath.Join(s.T().TempDir(), fmt.Sprintf("frac_v%d", v)) + + cmd := exec.Command("bash", "testdata/legacy/seal-fraction.sh", fmt.Sprintf("v%d", v), s.fracName) + cmd.Stdin = strings.NewReader(strings.Join(docLines(s.T(), testDocs), "\n")) + out, err := cmd.CombinedOutput() + s.Require().NoError(err, "seal-fraction.sh failed:\n%s", out) + + s.expectedVer = v +} + +func (s *FractionDecoderTestSuite) decodeAll() *collectedContent { + content := &collectedContent{t: s.T()} + err := decodeFraction(s.fracName, "", content) + s.Require().NoError(err) + return content +} + +func (s *FractionDecoderTestSuite) TestDecode() { + s.Run("info", s.checkInfo) + s.Run("docs keep raw json and skip the system doc", s.checkDocs) + s.Run("ids have contiguous lids", s.checkIDs) + s.Run("tokens are sequential and sorted with correct freq and postings", s.checkTokens) + s.Run("tokens without postings when requested alone", s.checkTokensWithoutPostings) + s.Run("only selects sections", s.checkOnlyFiltersSections) +} + +func (s *FractionDecoderTestSuite) checkInfo() { + content := s.decodeAll() + info := content.Info + s.Require().NotNil(info) + + // the decoder must recognize the format version, not just happen to + // read the files + s.Equal(s.expectedVer, info.BinaryDataVer) + + s.Equal(uint32(3), info.DocsTotal) + s.Positive(info.DocsOnDisk) + s.Less(uint64(info.From), uint64(info.To)) +} + +func (s *FractionDecoderTestSuite) checkDocs() { + content := s.decodeAll() + + s.Require().Len(content.Docs, 3) + s.ElementsMatch(docLines(s.T(), testDocs), docTexts(content.Docs)) +} + +func (s *FractionDecoderTestSuite) checkIDs() { + content := s.decodeAll() + + s.Require().Len(content.IDs, 4) + for i, id := range content.IDs { + s.Equal(uint32(i), id.LID) + } +} + +func (s *FractionDecoderTestSuite) checkTokens() { + content := s.decodeAll() + + s.Require().Len(content.Tokens, len(testTokens)) + for i, tok := range content.Tokens { + expected := testTokens[i] + expected.TID = uint32(i + 1) + expected.Kind = kindToken + expected.LIDs = tok.LIDs + s.Equal(expected, tok) + + s.Len(tok.LIDs, int(tok.Freq)) + for _, lid := range tok.LIDs { + s.NotZero(lid) + s.LessOrEqual(lid, uint32(3)) + } + } +} + +func (s *FractionDecoderTestSuite) checkTokensWithoutPostings() { + content := &collectedContent{t: s.T()} + err := decodeFraction(s.fracName, kindToken, content) + s.Require().NoError(err) + + s.Require().Len(content.Tokens, len(testTokens)) + for i, tok := range content.Tokens { + s.Equal(testTokens[i].Token, tok.Token) + s.Equal(testTokens[i].Freq, tok.Freq) + s.Empty(tok.LIDs) + } +} + +func (s *FractionDecoderTestSuite) checkOnlyFiltersSections() { + content := &collectedContent{t: s.T()} + err := decodeFraction(s.fracName, kindInfo+","+kindOffsets, content) + s.Require().NoError(err) + + s.NotNil(content.Info) + s.NotNil(content.Offsets) + s.Empty(content.Docs) + s.Empty(content.Tokens) + s.Empty(content.IDs) +} + +type CurrentVersionSuite struct { + FractionDecoderTestSuite +} + +func (s *CurrentVersionSuite) SetupTest() { + s.sealDocs(defaultDocsPerDocBlock) +} + +// LegacyVersionsSuite runs the same checks against a fraction sealed by +// the era's own code. The suite is instantiated per version by +// TestFractionDecoderLegacy; a future v7 in config/frac_version.go +// automatically adds v6 there. Versions older than v2 are not +// discoverable and are skipped. +type LegacyVersionsSuite struct { + FractionDecoderTestSuite + + Version config.BinaryDataVersion +} + +func (s *LegacyVersionsSuite) SetupTest() { + s.sealDocsEra(s.Version) +} + +func TestFractionDecoderCurrent(t *testing.T) { + suite.Run(t, new(CurrentVersionSuite)) +} + +func TestFractionDecoderLegacy(t *testing.T) { + for ver := config.BinaryDataV2; ver < config.CurrentFracVersion; ver++ { + t.Run(fmt.Sprintf("v%d", ver), func(t *testing.T) { + suite.Run(t, &LegacyVersionsSuite{Version: config.BinaryDataVersion(ver)}) + }) + } +} + +func TestUnknownOnlySection(t *testing.T) { + fracName := filepath.Join(t.TempDir(), "frac_000001") + + err := decodeFraction(fracName, "amogus", &collectedContent{}) + require.Error(t, err) + assert.Contains(t, err.Error(), "unknown section to decode: amogus") +} diff --git a/cmd/fraction/decoding_test.go b/cmd/fraction/sealing_test.go similarity index 51% rename from cmd/fraction/decoding_test.go rename to cmd/fraction/sealing_test.go index 6dc62206..ae33f2ed 100644 --- a/cmd/fraction/decoding_test.go +++ b/cmd/fraction/sealing_test.go @@ -15,10 +15,11 @@ import ( "github.com/ozontech/seq-db/frac/common" ) -// sealAndDecode seals docs (one JSON document per element, empty elements are -// skipped just like on stdin) into a fresh fraction in a temp dir using the -// mapping at mappingPath (empty string means the built-in default) and decodes -// the requested sections back into a collected form for assertions. +// sealAndDecode seals docs (one JSON document per element, empty elements +// are skipped just like on stdin) into a fresh fraction in a temp dir +// using the mapping at mappingPath (empty string means the built-in +// default) and decodes the requested sections back into a collected form +// for assertions. func sealAndDecode( t *testing.T, mappingPath string, @@ -29,8 +30,9 @@ func sealAndDecode( t.Helper() fracName := filepath.Join(t.TempDir(), "frac_000001") + r := strings.NewReader(strings.Join(docs, "\n")) - err := sealFraction(fracName, mappingPath, strings.NewReader(strings.Join(docs, "\n")), docsPerDocBlock) + err := sealFraction(fracName, mappingPath, r, docsPerDocBlock) require.NoError(t, err) for _, suffix := range []string{ @@ -135,153 +137,52 @@ var testTokens = []tokenRecord{ {Field: "service", Token: "billing", Freq: 1}, } -func TestSealAndDecode(t *testing.T) { - tests := []struct { - name string - docs []string - docsPerDocBlock int - only string - check func(t *testing.T, content *collectedContent) - }{ - { - name: "info", - docsPerDocBlock: defaultDocsPerDocBlock, - only: "info", - check: func(t *testing.T, content *collectedContent) { - require.NotNil(t, content.Info) - assert.Equal(t, uint32(3), content.Info.DocsTotal) - assert.Positive(t, content.Info.DocsOnDisk) - assert.Less(t, uint64(content.Info.From), uint64(content.Info.To)) - }, - }, - { - name: "docs keep raw json and skip the system doc", - docsPerDocBlock: defaultDocsPerDocBlock, - only: kindDoc, - check: func(t *testing.T, content *collectedContent) { - require.Len(t, content.Docs, 3) - assert.ElementsMatch(t, docLines(t, testDocs), docTexts(content.Docs)) - }, - }, - { - name: "tokens are sequential and sorted with correct freq and postings", - docsPerDocBlock: defaultDocsPerDocBlock, - only: kindID + "," + kindToken, - check: func(t *testing.T, content *collectedContent) { - require.Len(t, content.Tokens, len(testTokens)) - for i, tok := range content.Tokens { - expected := testTokens[i] - expected.TID = uint32(i + 1) - expected.Kind = kindToken - expected.LIDs = tok.LIDs - assert.Equal(t, expected, tok) - } - for _, tok := range content.Tokens { - assert.Len(t, tok.LIDs, int(tok.Freq)) - for _, lid := range tok.LIDs { - assert.NotZero(t, lid) - assert.LessOrEqual(t, lid, uint32(3)) - } - } - }, - }, - { - name: "tokens without postings when requested alone", - docsPerDocBlock: defaultDocsPerDocBlock, - only: kindToken, - check: func(t *testing.T, content *collectedContent) { - require.Len(t, content.Tokens, len(testTokens)) - for i, tok := range content.Tokens { - assert.Equal(t, testTokens[i].Token, tok.Token) - assert.Equal(t, testTokens[i].Freq, tok.Freq) - assert.Empty(t, tok.LIDs) - } - }, - }, - { - name: "offsets of a single doc block", - docsPerDocBlock: defaultDocsPerDocBlock, - only: kindOffsets, - check: func(t *testing.T, content *collectedContent) { - require.NotNil(t, content.Offsets) - assert.Equal(t, []uint64{0}, content.Offsets.Values) - }, - }, - { - name: "only selects sections", - docsPerDocBlock: defaultDocsPerDocBlock, - only: kindInfo + "," + kindOffsets, - check: func(t *testing.T, content *collectedContent) { - require.NotNil(t, content.Info) - require.NotNil(t, content.Offsets) - assert.Empty(t, content.Docs) - assert.Empty(t, content.Tokens) - assert.Empty(t, content.IDs) - }, - }, - { - name: "empty lines are skipped", - docsPerDocBlock: defaultDocsPerDocBlock, - docs: []string{ - `{"timestamp":"2024-01-01T10:00:00.000Z","level":"info","message":"first"}`, - "", - `{"timestamp":"2024-01-01T10:00:01.000Z","level":"info","message":"second"}`, - }, - only: kindInfo + "," + kindDoc, - check: func(t *testing.T, content *collectedContent) { - require.NotNil(t, content.Info) - assert.Equal(t, uint32(2), content.Info.DocsTotal) - assert.Len(t, content.Docs, 2) - }, - }, - { - name: "docs are split into several blocks", - docs: generateDocs(5), - docsPerDocBlock: 2, - only: kindInfo + "," + kindDoc + "," + kindID + "," + kindToken + "," + kindOffsets, - check: func(t *testing.T, content *collectedContent) { - // 5 docs with 2 docs per block: [2][2][1] - require.Len(t, content.Offsets.Values, 3) - assert.Equal(t, uint32(5), content.Info.DocsTotal) - - // all docs are decoded, ids have contiguous lids across block boundaries - require.Len(t, content.Docs, 5) - require.Len(t, content.IDs, 6) // +1 for the system id - for i, id := range content.IDs { - assert.Equal(t, uint32(i), id.LID) - } - - // every posting lid points to an existing doc - for _, tok := range content.Tokens { - assert.Len(t, tok.LIDs, int(tok.Freq)) - for _, lid := range tok.LIDs { - assert.NotZero(t, lid) - assert.LessOrEqual(t, lid, uint32(5)) - } - } - - // every doc position points into a real doc block, and every - // block is used: doc indexers assign block indexes concurrently, - // so only the counts per block are guaranteed, not the order - blocksPerLID := map[uint32]int{} - for _, id := range content.IDs[1:] { // system id is not positioned - assert.Less(t, id.BlockIndex, uint32(len(content.Offsets.Values)), "block index out of range") - blocksPerLID[id.BlockIndex]++ - } - assert.Equal(t, map[uint32]int{0: 2, 1: 2, 2: 1}, blocksPerLID) - }, - }, +func TestSealSkipsEmptyLines(t *testing.T) { + content := sealAndDecode(t, "", defaultDocsPerDocBlock, kindInfo+","+kindDoc, + `{"timestamp":"2024-01-01T10:00:00.000Z","level":"info","message":"first"}`, + "", + `{"timestamp":"2024-01-01T10:00:01.000Z","level":"info","message":"second"}`, + ) + + require.NotNil(t, content.Info) + assert.Equal(t, uint32(2), content.Info.DocsTotal) + assert.Len(t, content.Docs, 2) +} + +func TestSealSplitsIntoDocBlocks(t *testing.T) { + content := sealAndDecode(t, "", 2, kindInfo+","+kindDoc+","+kindID+","+kindToken+","+kindOffsets, + generateDocs(5)..., + ) + + // 5 docs with 2 docs per block: [2][2][1] + require.Len(t, content.Offsets.Values, 3) + assert.Equal(t, uint32(5), content.Info.DocsTotal) + + // all docs are decoded, ids have contiguous lids across block boundaries + require.Len(t, content.Docs, 5) + require.Len(t, content.IDs, 6) // +1 for the system id + for i, id := range content.IDs { + assert.Equal(t, uint32(i), id.LID) } - for _, tt := range tests { - t.Run(tt.name, func(t *testing.T) { - docs := tt.docs - if docs == nil { - docs = docLines(t, testDocs) - } - tt.check(t, sealAndDecode(t, "", tt.docsPerDocBlock, tt.only, docs...)) - }) + // every posting lid points to an existing doc + for _, tok := range content.Tokens { + assert.Len(t, tok.LIDs, int(tok.Freq)) + for _, lid := range tok.LIDs { + assert.NotZero(t, lid) + assert.LessOrEqual(t, lid, uint32(5)) + } } + + // every doc position points into a real doc block, and every block is + // used: doc indexers assign block indexes concurrently, so only the + // counts per block are guaranteed, not the order + blocksPerLID := map[uint32]int{} + for _, id := range content.IDs[1:] { // system id is not positioned + assert.Less(t, id.BlockIndex, uint32(len(content.Offsets.Values)), "block index out of range") + blocksPerLID[id.BlockIndex]++ + } + assert.Equal(t, map[uint32]int{0: 2, 1: 2, 2: 1}, blocksPerLID) } func TestSealFailures(t *testing.T) { @@ -308,15 +209,6 @@ func TestSealFailures(t *testing.T) { }) } } - -func TestUnknownOnlySection(t *testing.T) { - fracName := filepath.Join(t.TempDir(), "frac_000001") - - err := decodeFraction(fracName, "amogus", &collectedContent{}) - require.Error(t, err) - assert.Contains(t, err.Error(), "unknown section to decode: amogus") -} - func TestSealCustomMapping(t *testing.T) { mapping := filepath.Join(t.TempDir(), "mapping.yaml") err := os.WriteFile(mapping, []byte("mapping-list:\n - type: \"path\"\n name: \"request_uri\"\n"), 0o600) diff --git a/cmd/fraction/testdata/legacy/seal-fraction.sh b/cmd/fraction/testdata/legacy/seal-fraction.sh new file mode 100644 index 00000000..240e3501 --- /dev/null +++ b/cmd/fraction/testdata/legacy/seal-fraction.sh @@ -0,0 +1,134 @@ +#!/usr/bin/env bash +# Seals a fraction in a given format version. Used by tests and CI to run +# the cmd/fraction tests against fractions produced by older seq-db code. +# +# Usage: +# seal-fraction.sh < docs.jsonl +# +# is a fraction version: v2, v3, v4, v5, ... or "current". +# v2..v5 — a git worktree is created at the last commit of that +# version's era (found dynamically: the parent of the first +# commit introducing the next version in config/frac_version.go) +# and a standalone sealer is run there; old code has no +# `fraction seal`. The sealer is the single shared file +# (testdata/legacy/sealer/main.go), patched for the +# era: the few lines that differ between the eras are +# replaced with sed. +# v6+ — a git worktree is created at the era commit the same way, +# but the era's own `go run ./cmd/fraction seal` is used +# (it exists since v6 and seals in the code's current +# format). Needs no per-version support: a new vN+1 works +# with zero script changes. +# current — the local code's own `go run ./cmd/fraction seal` is used. +# +# is the output fraction base name (without suffixes). +# Documents (one JSON per line) are read from stdin. +# +# Prints the fraction base name on stdout. + +set -euo pipefail + +VERSION="${1:?usage: seal-fraction.sh < docs.jsonl}" +FRAC="${2:?usage: seal-fraction.sh < docs.jsonl}" + +SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" +REPO="$(git -C "$SCRIPT_DIR" rev-parse --show-toplevel)" + +# Last commit of version's era: the parent of the first commit that +# introduced the next version (BinaryDataV) in config/frac_version.go. +# When the next version is not committed yet (e.g. a locally added V7), +# the era has not ended: fall back to the last commit of main. Versions +# whose next one appeared before the file existed (v0, v1) are not +# discoverable and are not supported. +era_commit() { + local next=$((VERSION_NUM + 1)) + for commit in $(git -C "$REPO" rev-list --reverse main -- config/frac_version.go); do + if git -C "$REPO" show "${commit}:config/frac_version.go" 2>/dev/null | grep -q "BinaryDataV${next}\b"; then + git -C "$REPO" rev-parse --short "${commit}^" + return + fi + done + if [[ "$VERSION_NUM" -ge "$(max_known_version_num)" ]]; then + git -C "$REPO" rev-parse --short main + return + fi + echo "cannot find the commit introducing BinaryDataV${next}: version $VERSION is too old (supported: v2+)" >&2 + return 1 +} + +# The highest BinaryDataVN declared in the current config/frac_version.go. +max_known_version_num() { + git -C "$REPO" show "main:config/frac_version.go" | grep -oE 'BinaryDataV[0-9]+' | grep -oE '[0-9]+' | sort -n | tail -1 +} + +# Patches the shared sealer for a given pre-v6 format era: replaces the +# lines that differ between the eras. The file as committed targets the +# newest sealer era (v5); older eras get a sed-reduced version. v6+ eras +# do not reach this function: they seal with their own cmd/fraction. +patch_sealer() { + local ver="$1" file="$2" + if [[ "$ver" == "v5" || "$ver" == "v4" ]]; then + # v4 sealer is identical to the v5 one + cat "$file" + return + fi + if [[ "$ver" == "v3" ]]; then + # v3: sealing lived in frac/sealed/sealing, and SealParams + # had no LIDBlockSize/TokenBlockSize + sed -e 's|"github.com/ozontech/seq-db/sealing"|"github.com/ozontech/seq-db/frac/sealed/sealing"|' \ + -e '/LIDBlockSize:/d' -e '/TokenBlockSize:/d' \ + "$file" + return + fi + # v2: on top of the v3 differences, the skip mask provider interface + # had no GetIDsBitmapByFrac (and roaring is not in the go.mod), and + # GetIDsIteratorByFrac returned no cleanup func + sed -e 's|"github.com/ozontech/seq-db/sealing"|"github.com/ozontech/seq-db/frac/sealed/sealing"|' \ + -e '/LIDBlockSize:/d' -e '/TokenBlockSize:/d' -e '/LIDsBitmapThreshold:/d' \ + -e '/"github.com\/RoaringBitmap\/roaring\/v2"/d' \ + -e '/GetIDsBitmapByFrac(_ string, _, _ uint32)/,/^}/d' \ + -e 's|return node.NewStatic(nil, reverse), false, func() error { return nil }, nil|return node.NewStatic(nil, reverse), false, nil|' \ + -e 's|(_ string, _, _ uint32, reverse bool) (node.Node, bool, func() error, error)|(_ string, _, _ uint32, reverse bool) (node.Node, bool, error)|' \ + "$file" +} + +if [[ "$VERSION" == "current" ]]; then + (cd "$REPO" && go run ./cmd/fraction seal "$FRAC") +else + [[ "$VERSION" =~ ^v[0-9]+$ ]] || { + echo "unknown version: $VERSION (expected vN or current)" >&2 + exit 2 + } + VERSION_NUM="${VERSION#v}" + [[ "$VERSION_NUM" -ge 2 ]] || { + echo "version $VERSION is too old (supported: v2+)" >&2 + exit 2 + } + + commit="$(era_commit)" + wt="${TMPDIR:-/tmp}/fraction-seal-$VERSION" + + if [[ -d "$wt/.git" || -f "$wt/.git" ]]; then + git -C "$wt" checkout -q --detach "$(git -C "$REPO" rev-parse "$commit")" + else + git -C "$REPO" worktree add --detach "$wt" "$commit" >/dev/null + fi + + if [[ "$VERSION_NUM" -ge 6 ]]; then + # v6+ eras have their own cmd/fraction seal + (cd "$wt" && go run ./cmd/fraction seal "$FRAC") + else + sealer="$SCRIPT_DIR/sealer/main.go" + [[ -f "$sealer" ]] || { + echo "standalone sealer is missing: $sealer" >&2 + exit 2 + } + + mkdir -p "$wt/sealer" + patch_sealer "$VERSION" "$sealer" > "$wt/sealer/main.go" + + (cd "$wt" && CGO_ENABLED=0 go run ./sealer "$FRAC") + fi +fi + +echo "$FRAC" diff --git a/cmd/fraction/testdata/legacy/sealer/main.go b/cmd/fraction/testdata/legacy/sealer/main.go new file mode 100644 index 00000000..01c4eab5 --- /dev/null +++ b/cmd/fraction/testdata/legacy/sealer/main.go @@ -0,0 +1,139 @@ +package main + +// Standalone sealer for the legacy fraction formats (V3..V5): seals JSON +// documents from stdin, one per line. Built against the era commit of a +// format: testdata/legacy/seal-fraction.sh copies this file into a +// checked-out worktree of that commit, patching the few lines that +// differ between the eras (imports, SealParams fields). Do not build it +// in the main module: the era code it compiles against is older than +// the current one. + +import ( + "bufio" + "fmt" + "os" + "path/filepath" + "sync" + "time" + + "github.com/RoaringBitmap/roaring/v2" + "github.com/alecthomas/units" + + "github.com/ozontech/seq-db/cache" + "github.com/ozontech/seq-db/frac" + "github.com/ozontech/seq-db/frac/common" + "github.com/ozontech/seq-db/indexer" + "github.com/ozontech/seq-db/node" + "github.com/ozontech/seq-db/sealing" + "github.com/ozontech/seq-db/seq" + "github.com/ozontech/seq-db/storage" + "github.com/ozontech/seq-db/tokenizer" +) + +type stubSkipMaskProvider struct{} + +func (stubSkipMaskProvider) GetIDsIteratorByFrac(_ string, _, _ uint32, reverse bool) (node.Node, bool, func() error, error) { + return node.NewStatic(nil, reverse), false, func() error { return nil }, nil +} + +func (stubSkipMaskProvider) GetIDsBitmapByFrac(_ string, _, _ uint32) (*roaring.Bitmap, error) { + return nil, nil +} + +func (stubSkipMaskProvider) RemoveFrac(_ string) {} + +func main() { + if len(os.Args) != 2 { + fmt.Fprintln(os.Stderr, "usage: sealer < docs.jsonl") + os.Exit(2) + } + fracName := os.Args[1] + + scanner := bufio.NewScanner(os.Stdin) + scanner.Buffer(make([]byte, 0, 1024*1024), 16*1024*1024) + var readNext func() ([]byte, error) + readNext = func() ([]byte, error) { + if !scanner.Scan() { + if err := scanner.Err(); err != nil { + return nil, err + } + return nil, nil + } + line := scanner.Bytes() + if len(line) == 0 { + return readNext() + } + return line, nil + } + + mapping := seq.Mapping{ + "level": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + "service": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + "message": seq.NewSingleType(seq.TokenizerTypeText, "", 0), + } + tokenizers := map[seq.TokenizerType]tokenizer.Tokenizer{ + seq.TokenizerTypeKeyword: tokenizer.NewKeywordTokenizer(20, false, true), + seq.TokenizerTypeText: tokenizer.NewTextTokenizer(20, false, true, 100), + } + + activeIndexer, stopIndexer := frac.NewActiveIndexer(4, 10) + defer stopIndexer() + + active := frac.NewActive( + fracName, + activeIndexer, + storage.NewReadLimiter(1, nil), + cache.NewCache[[]byte](nil, nil), + cache.NewCache[[]byte](nil, nil), + &frac.Config{}, + stubSkipMaskProvider{}, + ) + + proc := indexer.NewProcessor(mapping, tokenizers, 0, 0, 0) + compressor := indexer.GetDocsMetasCompressor(3, 3) + + _, binaryDocs, binaryMeta, err := proc.ProcessBulk(time.Now(), nil, nil, readNext) + if err != nil { + fmt.Fprintf(os.Stderr, "process bulk: %s\n", err) + os.Exit(2) + } + + compressor.CompressDocsAndMetas(binaryDocs, binaryMeta) + docsBlock, metasBlock := compressor.DocsMetas() + + var wg sync.WaitGroup + wg.Add(1) + if err := active.Append(docsBlock, metasBlock, &wg); err != nil { + fmt.Fprintf(os.Stderr, "append: %s\n", err) + os.Exit(2) + } + wg.Wait() + + sealParams := common.SealParams{ + IDsZstdLevel: 1, + LIDsZstdLevel: 1, + TokenListZstdLevel: 1, + DocsPositionsZstdLevel: 1, + TokenTableZstdLevel: 1, + DocBlocksZstdLevel: 1, + LIDBlockSize: 256, + TokenBlockSize: 128, + DocBlockSize: 128 * int(units.KiB), + } + + src, err := frac.NewActiveSealingSource(active, sealParams) + if err != nil { + fmt.Fprintf(os.Stderr, "sealing source: %s\n", err) + os.Exit(2) + } + + if _, err := sealing.Seal(src, sealParams); err != nil { + fmt.Fprintf(os.Stderr, "seal: %s\n", err) + os.Exit(2) + } + + active.Release() + + files, _ := filepath.Glob(fracName + "*") + fmt.Fprintf(os.Stderr, "sealed: %v\n", files) +} From 0c67606eef001670b533ea3e9aacc27d888f11c9 Mon Sep 17 00:00:00 2001 From: Daniil Lyubaev Date: Mon, 5 Oct 2026 16:45:35 +0300 Subject: [PATCH 2/7] fix: add mapping --- cmd/fraction/testdata/legacy/seal-fraction.sh | 39 +++++++++++-- cmd/fraction/testdata/legacy/sealer/main.go | 58 +++++++++++++++---- 2 files changed, 81 insertions(+), 16 deletions(-) diff --git a/cmd/fraction/testdata/legacy/seal-fraction.sh b/cmd/fraction/testdata/legacy/seal-fraction.sh index 240e3501..bd4efb26 100644 --- a/cmd/fraction/testdata/legacy/seal-fraction.sh +++ b/cmd/fraction/testdata/legacy/seal-fraction.sh @@ -5,6 +5,11 @@ # Usage: # seal-fraction.sh < docs.jsonl # +# Optional flag: +# --mapping= mapping YAML for the indexed fields; the file must +# be readable by the era's seq.ReadMapping. Without +# the flag a built-in default mapping is used. +# # is a fraction version: v2, v3, v4, v5, ... or "current". # v2..v5 — a git worktree is created at the last commit of that # version's era (found dynamically: the parent of the first @@ -28,12 +33,24 @@ set -euo pipefail -VERSION="${1:?usage: seal-fraction.sh < docs.jsonl}" -FRAC="${2:?usage: seal-fraction.sh < docs.jsonl}" +VERSION="" +FRAC="" +for arg in "$@"; do + case "$arg" in + --mapping=*) MAPPING="${arg#--mapping=}" ;; + *) if [[ -z "$VERSION" ]]; then VERSION="$arg"; else FRAC="$arg"; fi ;; + esac +done +VERSION="${VERSION:?usage: seal-fraction.sh [--mapping=file] < docs.jsonl}" +FRAC="${FRAC:?usage: seal-fraction.sh [--mapping=file] < docs.jsonl}" +MAPPING="${MAPPING:-}" SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" REPO="$(git -C "$SCRIPT_DIR" rev-parse --show-toplevel)" +# the mapping file is opened by sealers running from other cwd's +[[ -z "$MAPPING" || "$MAPPING" == /* ]] || MAPPING="$(pwd)/$MAPPING" + # Last commit of version's era: the parent of the first commit that # introduced the next version (BinaryDataV) in config/frac_version.go. # When the next version is not committed yet (e.g. a locally added V7), @@ -93,7 +110,11 @@ patch_sealer() { } if [[ "$VERSION" == "current" ]]; then - (cd "$REPO" && go run ./cmd/fraction seal "$FRAC") + if [[ -n "$MAPPING" ]]; then + (cd "$REPO" && go run ./cmd/fraction seal --mapping="$MAPPING" "$FRAC") + else + (cd "$REPO" && go run ./cmd/fraction seal "$FRAC") + fi else [[ "$VERSION" =~ ^v[0-9]+$ ]] || { echo "unknown version: $VERSION (expected vN or current)" >&2 @@ -116,7 +137,11 @@ else if [[ "$VERSION_NUM" -ge 6 ]]; then # v6+ eras have their own cmd/fraction seal - (cd "$wt" && go run ./cmd/fraction seal "$FRAC") + if [[ -n "$MAPPING" ]]; then + (cd "$wt" && go run ./cmd/fraction seal --mapping="$MAPPING" "$FRAC") + else + (cd "$wt" && go run ./cmd/fraction seal "$FRAC") + fi else sealer="$SCRIPT_DIR/sealer/main.go" [[ -f "$sealer" ]] || { @@ -127,7 +152,11 @@ else mkdir -p "$wt/sealer" patch_sealer "$VERSION" "$sealer" > "$wt/sealer/main.go" - (cd "$wt" && CGO_ENABLED=0 go run ./sealer "$FRAC") + if [[ -n "$MAPPING" ]]; then + (cd "$wt" && CGO_ENABLED=0 go run ./sealer "$FRAC" --mapping="$MAPPING") + else + (cd "$wt" && CGO_ENABLED=0 go run ./sealer "$FRAC") + fi fi fi diff --git a/cmd/fraction/testdata/legacy/sealer/main.go b/cmd/fraction/testdata/legacy/sealer/main.go index 01c4eab5..3cabc0ce 100644 --- a/cmd/fraction/testdata/legacy/sealer/main.go +++ b/cmd/fraction/testdata/legacy/sealer/main.go @@ -1,18 +1,19 @@ package main -// Standalone sealer for the legacy fraction formats (V3..V5): seals JSON +// Standalone sealer for the legacy fraction formats (V2..V5): seals JSON // documents from stdin, one per line. Built against the era commit of a // format: testdata/legacy/seal-fraction.sh copies this file into a // checked-out worktree of that commit, patching the few lines that -// differ between the eras (imports, SealParams fields). Do not build it -// in the main module: the era code it compiles against is older than -// the current one. +// differ between the eras (imports, SealParams fields, flags). Do not +// build it in the main module: the era code it compiles against is +// older than the current one. import ( "bufio" "fmt" "os" "path/filepath" + "strings" "sync" "time" @@ -43,11 +44,28 @@ func (stubSkipMaskProvider) GetIDsBitmapByFrac(_ string, _, _ uint32) (*roaring. func (stubSkipMaskProvider) RemoveFrac(_ string) {} func main() { - if len(os.Args) != 2 { - fmt.Fprintln(os.Stderr, "usage: sealer < docs.jsonl") + var mappingPath, fracName string + args := os.Args[1:] + for i := 0; i < len(args); i++ { + arg := args[i] + switch { + case arg == "--mapping": + if i+1 >= len(args) { + fmt.Fprintln(os.Stderr, "--mapping requires a value") + os.Exit(2) + } + i++ + mappingPath = args[i] + case strings.HasPrefix(arg, "--mapping="): + mappingPath = strings.TrimPrefix(arg, "--mapping=") + default: + fracName = arg + } + } + if fracName == "" { + fmt.Fprintln(os.Stderr, "usage: sealer [--mapping=mapping.yaml] < docs.jsonl") os.Exit(2) } - fracName := os.Args[1] scanner := bufio.NewScanner(os.Stdin) scanner.Buffer(make([]byte, 0, 1024*1024), 16*1024*1024) @@ -66,14 +84,24 @@ func main() { return line, nil } - mapping := seq.Mapping{ - "level": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - "service": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - "message": seq.NewSingleType(seq.TokenizerTypeText, "", 0), + mapping := defaultMapping() + if mappingPath != "" { + data, err := os.ReadFile(mappingPath) + if err != nil { + fmt.Fprintf(os.Stderr, "cannot read mapping: %s\n", err) + os.Exit(2) + } + mapping, err = seq.ReadMapping(data) + if err != nil { + fmt.Fprintf(os.Stderr, "cannot parse mapping: %s\n", err) + os.Exit(2) + } } tokenizers := map[seq.TokenizerType]tokenizer.Tokenizer{ seq.TokenizerTypeKeyword: tokenizer.NewKeywordTokenizer(20, false, true), seq.TokenizerTypeText: tokenizer.NewTextTokenizer(20, false, true, 100), + seq.TokenizerTypePath: tokenizer.NewPathTokenizer(512, false, true), + seq.TokenizerTypeExists: tokenizer.NewExistsTokenizer(), } activeIndexer, stopIndexer := frac.NewActiveIndexer(4, 10) @@ -137,3 +165,11 @@ func main() { files, _ := filepath.Glob(fracName + "*") fmt.Fprintf(os.Stderr, "sealed: %v\n", files) } + +func defaultMapping() seq.Mapping { + return seq.Mapping{ + "level": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + "service": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + "message": seq.NewSingleType(seq.TokenizerTypeText, "", 0), + } +} From aed7a0f60c9b1a4c343c518626f553a0f9e6cb32 Mon Sep 17 00:00:00 2001 From: Daniil Lyubaev Date: Mon, 5 Oct 2026 17:11:57 +0300 Subject: [PATCH 3/7] fix: ci with history --- .github/workflows/ci.yml | 2 ++ 1 file changed, 2 insertions(+) diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 7c772178..1b341a9b 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -31,6 +31,8 @@ jobs: - name: Checkout code uses: actions/checkout@v6 + with: + fetch-depth: 0 - name: Install Go uses: actions/setup-go@v6 From ec7c99dbbbee1350c7bc8c9081e9a3a0432ca002 Mon Sep 17 00:00:00 2001 From: Daniil Lyubaev Date: Mon, 5 Oct 2026 17:23:34 +0300 Subject: [PATCH 4/7] fix: history ref --- cmd/fraction/testdata/legacy/seal-fraction.sh | 30 +++++++++++++++---- 1 file changed, 24 insertions(+), 6 deletions(-) diff --git a/cmd/fraction/testdata/legacy/seal-fraction.sh b/cmd/fraction/testdata/legacy/seal-fraction.sh index bd4efb26..135b119e 100644 --- a/cmd/fraction/testdata/legacy/seal-fraction.sh +++ b/cmd/fraction/testdata/legacy/seal-fraction.sh @@ -48,25 +48,41 @@ MAPPING="${MAPPING:-}" SCRIPT_DIR="$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)" REPO="$(git -C "$SCRIPT_DIR" rev-parse --show-toplevel)" +# Era commit discovery needs history; a PR checkout has no local "main" +# branch (detached HEAD), so prefer origin/main, then main, then HEAD. +history_ref() { + for ref in origin/main main HEAD; do + if git -C "$REPO" rev-parse -q --verify "$ref" >/dev/null; then + echo "$ref" + return + fi + done + echo "no git history ref found (tried origin/main, main, HEAD)" >&2 + return 1 +} + # the mapping file is opened by sealers running from other cwd's [[ -z "$MAPPING" || "$MAPPING" == /* ]] || MAPPING="$(pwd)/$MAPPING" # Last commit of version's era: the parent of the first commit that # introduced the next version (BinaryDataV) in config/frac_version.go. # When the next version is not committed yet (e.g. a locally added V7), -# the era has not ended: fall back to the last commit of main. Versions -# whose next one appeared before the file existed (v0, v1) are not -# discoverable and are not supported. +# the era has not ended: fall back to the history ref tip. Versions whose +# next one appeared before the file existed (v0, v1) are not discoverable +# and are not supported. Requires real history: in a shallow CI clone the +# caller must fetch it first (fetch-depth: 0 or git fetch --unshallow). era_commit() { local next=$((VERSION_NUM + 1)) - for commit in $(git -C "$REPO" rev-list --reverse main -- config/frac_version.go); do + local history + history="$(history_ref)" + for commit in $(git -C "$REPO" rev-list --reverse "$history" -- config/frac_version.go); do if git -C "$REPO" show "${commit}:config/frac_version.go" 2>/dev/null | grep -q "BinaryDataV${next}\b"; then git -C "$REPO" rev-parse --short "${commit}^" return fi done if [[ "$VERSION_NUM" -ge "$(max_known_version_num)" ]]; then - git -C "$REPO" rev-parse --short main + git -C "$REPO" rev-parse --short "$history" return fi echo "cannot find the commit introducing BinaryDataV${next}: version $VERSION is too old (supported: v2+)" >&2 @@ -75,7 +91,9 @@ era_commit() { # The highest BinaryDataVN declared in the current config/frac_version.go. max_known_version_num() { - git -C "$REPO" show "main:config/frac_version.go" | grep -oE 'BinaryDataV[0-9]+' | grep -oE '[0-9]+' | sort -n | tail -1 + local history + history="$(history_ref)" + git -C "$REPO" show "${history}:config/frac_version.go" | grep -oE 'BinaryDataV[0-9]+' | grep -oE '[0-9]+' | sort -n | tail -1 } # Patches the shared sealer for a given pre-v6 format era: replaces the From 17fe85e6f756bd88caa7ef857b4b6c896bc7c28f Mon Sep 17 00:00:00 2001 From: Daniil Lyubaev Date: Mon, 5 Oct 2026 17:49:02 +0300 Subject: [PATCH 5/7] fix: comment --- cmd/fraction/decoding_suites_test.go | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/cmd/fraction/decoding_suites_test.go b/cmd/fraction/decoding_suites_test.go index 06deeef5..d4ed04d2 100644 --- a/cmd/fraction/decoding_suites_test.go +++ b/cmd/fraction/decoding_suites_test.go @@ -2,8 +2,8 @@ package main // The suites run the same decode checks against fractions from different // sources: sealed on the fly by the current code (default) or sealed by -// the era's own code for every legacy format version ("legacy", -// see testdata/legacy/seal-fraction.sh). +// the era's own code for every legacy format version +// (see testdata/legacy/seal-fraction.sh). import ( "fmt" From daba14bd054b897a740aa2f1fb048ba440f1a25a Mon Sep 17 00:00:00 2001 From: Daniil Lyubaev Date: Tue, 6 Oct 2026 11:31:15 +0300 Subject: [PATCH 6/7] fix: use commits for legacy sealers --- cmd/fraction/testdata/legacy/seal-fraction.sh | 97 ++++------ cmd/fraction/testdata/legacy/sealer/main.go | 175 ------------------ 2 files changed, 38 insertions(+), 234 deletions(-) delete mode 100644 cmd/fraction/testdata/legacy/sealer/main.go diff --git a/cmd/fraction/testdata/legacy/seal-fraction.sh b/cmd/fraction/testdata/legacy/seal-fraction.sh index 135b119e..af9ee96e 100644 --- a/cmd/fraction/testdata/legacy/seal-fraction.sh +++ b/cmd/fraction/testdata/legacy/seal-fraction.sh @@ -11,19 +11,17 @@ # the flag a built-in default mapping is used. # # is a fraction version: v2, v3, v4, v5, ... or "current". -# v2..v5 — a git worktree is created at the last commit of that -# version's era (found dynamically: the parent of the first -# commit introducing the next version in config/frac_version.go) -# and a standalone sealer is run there; old code has no -# `fraction seal`. The sealer is the single shared file -# (testdata/legacy/sealer/main.go), patched for the -# era: the few lines that differ between the eras are -# replaced with sed. -# v6+ — a git worktree is created at the era commit the same way, -# but the era's own `go run ./cmd/fraction seal` is used -# (it exists since v6 and seals in the code's current -# format). Needs no per-version support: a new vN+1 works -# with zero script changes. +# v2..v5 — a git worktree is created at that version's sealer commit +# (an era commit with cmd/sealer added on top, see the +# sealer_commit table below) and the committed standalone +# sealer is run there; old code has no `fraction seal`. +# v6+ — a git worktree is created at the era commit found +# dynamically (the parent of the first commit introducing +# the next version in config/frac_version.go) and the era's +# own `go run ./cmd/fraction seal` is used (it exists since +# v6 and seals in the code's current format). Needs no +# per-version support: a new vN+1 works with zero script +# changes. # current — the local code's own `go run ./cmd/fraction seal` is used. # # is the output fraction base name (without suffixes). @@ -67,10 +65,9 @@ history_ref() { # Last commit of version's era: the parent of the first commit that # introduced the next version (BinaryDataV) in config/frac_version.go. # When the next version is not committed yet (e.g. a locally added V7), -# the era has not ended: fall back to the history ref tip. Versions whose -# next one appeared before the file existed (v0, v1) are not discoverable -# and are not supported. Requires real history: in a shallow CI clone the -# caller must fetch it first (fetch-depth: 0 or git fetch --unshallow). +# the era has not ended: fall back to the history ref tip. Requires real +# history: in a shallow CI clone the caller must fetch it first +# (fetch-depth: 0 or git fetch --unshallow). era_commit() { local next=$((VERSION_NUM + 1)) local history @@ -85,7 +82,7 @@ era_commit() { git -C "$REPO" rev-parse --short "$history" return fi - echo "cannot find the commit introducing BinaryDataV${next}: version $VERSION is too old (supported: v2+)" >&2 + echo "cannot find the commit introducing BinaryDataV${next}: version $VERSION is too old (supported: v6+)" >&2 return 1 } @@ -96,35 +93,15 @@ max_known_version_num() { git -C "$REPO" show "${history}:config/frac_version.go" | grep -oE 'BinaryDataV[0-9]+' | grep -oE '[0-9]+' | sort -n | tail -1 } -# Patches the shared sealer for a given pre-v6 format era: replaces the -# lines that differ between the eras. The file as committed targets the -# newest sealer era (v5); older eras get a sed-reduced version. v6+ eras -# do not reach this function: they seal with their own cmd/fraction. -patch_sealer() { - local ver="$1" file="$2" - if [[ "$ver" == "v5" || "$ver" == "v4" ]]; then - # v4 sealer is identical to the v5 one - cat "$file" - return - fi - if [[ "$ver" == "v3" ]]; then - # v3: sealing lived in frac/sealed/sealing, and SealParams - # had no LIDBlockSize/TokenBlockSize - sed -e 's|"github.com/ozontech/seq-db/sealing"|"github.com/ozontech/seq-db/frac/sealed/sealing"|' \ - -e '/LIDBlockSize:/d' -e '/TokenBlockSize:/d' \ - "$file" - return - fi - # v2: on top of the v3 differences, the skip mask provider interface - # had no GetIDsBitmapByFrac (and roaring is not in the go.mod), and - # GetIDsIteratorByFrac returned no cleanup func - sed -e 's|"github.com/ozontech/seq-db/sealing"|"github.com/ozontech/seq-db/frac/sealed/sealing"|' \ - -e '/LIDBlockSize:/d' -e '/TokenBlockSize:/d' -e '/LIDsBitmapThreshold:/d' \ - -e '/"github.com\/RoaringBitmap\/roaring\/v2"/d' \ - -e '/GetIDsBitmapByFrac(_ string, _, _ uint32)/,/^}/d' \ - -e 's|return node.NewStatic(nil, reverse), false, func() error { return nil }, nil|return node.NewStatic(nil, reverse), false, nil|' \ - -e 's|(_ string, _, _ uint32, reverse bool) (node.Node, bool, func() error, error)|(_ string, _, _ uint32, reverse bool) (node.Node, bool, error)|' \ - "$file" +# Sealer commits for the pre-v6 formats: an era commit with the +# standalone sealer (cmd/sealer) committed on top. +sealer_commit() { + case "$VERSION" in + v2) echo "98d29f2af902ee37ac24da18255a0a69a0da8266" ;; + v3) echo "1e214d0a282f29f577f38dcc920a8b8c6d95eab1" ;; + v4) echo "113f8663b80ce44149c12f1333237313b535e1d1" ;; + v5) echo "02e864f6a26b193fbfdcc973a5de3b8b9dd15f60" ;; + esac } if [[ "$VERSION" == "current" ]]; then @@ -144,7 +121,17 @@ else exit 2 } - commit="$(era_commit)" + if [[ "$VERSION_NUM" -ge 6 ]]; then + # v6+ eras seal with their own cmd/fraction at the era commit + commit="$(era_commit)" + else + commit="$(sealer_commit)" + [[ -n "$commit" ]] || { + echo "no sealer commit for version $VERSION" >&2 + exit 2 + } + fi + wt="${TMPDIR:-/tmp}/fraction-seal-$VERSION" if [[ -d "$wt/.git" || -f "$wt/.git" ]]; then @@ -161,19 +148,11 @@ else (cd "$wt" && go run ./cmd/fraction seal "$FRAC") fi else - sealer="$SCRIPT_DIR/sealer/main.go" - [[ -f "$sealer" ]] || { - echo "standalone sealer is missing: $sealer" >&2 - exit 2 - } - - mkdir -p "$wt/sealer" - patch_sealer "$VERSION" "$sealer" > "$wt/sealer/main.go" - + # the sealer is committed at the era commit (cmd/sealer) if [[ -n "$MAPPING" ]]; then - (cd "$wt" && CGO_ENABLED=0 go run ./sealer "$FRAC" --mapping="$MAPPING") + (cd "$wt" && CGO_ENABLED=0 go run ./cmd/sealer "$FRAC" --mapping="$MAPPING") else - (cd "$wt" && CGO_ENABLED=0 go run ./sealer "$FRAC") + (cd "$wt" && CGO_ENABLED=0 go run ./cmd/sealer "$FRAC") fi fi fi diff --git a/cmd/fraction/testdata/legacy/sealer/main.go b/cmd/fraction/testdata/legacy/sealer/main.go deleted file mode 100644 index 3cabc0ce..00000000 --- a/cmd/fraction/testdata/legacy/sealer/main.go +++ /dev/null @@ -1,175 +0,0 @@ -package main - -// Standalone sealer for the legacy fraction formats (V2..V5): seals JSON -// documents from stdin, one per line. Built against the era commit of a -// format: testdata/legacy/seal-fraction.sh copies this file into a -// checked-out worktree of that commit, patching the few lines that -// differ between the eras (imports, SealParams fields, flags). Do not -// build it in the main module: the era code it compiles against is -// older than the current one. - -import ( - "bufio" - "fmt" - "os" - "path/filepath" - "strings" - "sync" - "time" - - "github.com/RoaringBitmap/roaring/v2" - "github.com/alecthomas/units" - - "github.com/ozontech/seq-db/cache" - "github.com/ozontech/seq-db/frac" - "github.com/ozontech/seq-db/frac/common" - "github.com/ozontech/seq-db/indexer" - "github.com/ozontech/seq-db/node" - "github.com/ozontech/seq-db/sealing" - "github.com/ozontech/seq-db/seq" - "github.com/ozontech/seq-db/storage" - "github.com/ozontech/seq-db/tokenizer" -) - -type stubSkipMaskProvider struct{} - -func (stubSkipMaskProvider) GetIDsIteratorByFrac(_ string, _, _ uint32, reverse bool) (node.Node, bool, func() error, error) { - return node.NewStatic(nil, reverse), false, func() error { return nil }, nil -} - -func (stubSkipMaskProvider) GetIDsBitmapByFrac(_ string, _, _ uint32) (*roaring.Bitmap, error) { - return nil, nil -} - -func (stubSkipMaskProvider) RemoveFrac(_ string) {} - -func main() { - var mappingPath, fracName string - args := os.Args[1:] - for i := 0; i < len(args); i++ { - arg := args[i] - switch { - case arg == "--mapping": - if i+1 >= len(args) { - fmt.Fprintln(os.Stderr, "--mapping requires a value") - os.Exit(2) - } - i++ - mappingPath = args[i] - case strings.HasPrefix(arg, "--mapping="): - mappingPath = strings.TrimPrefix(arg, "--mapping=") - default: - fracName = arg - } - } - if fracName == "" { - fmt.Fprintln(os.Stderr, "usage: sealer [--mapping=mapping.yaml] < docs.jsonl") - os.Exit(2) - } - - scanner := bufio.NewScanner(os.Stdin) - scanner.Buffer(make([]byte, 0, 1024*1024), 16*1024*1024) - var readNext func() ([]byte, error) - readNext = func() ([]byte, error) { - if !scanner.Scan() { - if err := scanner.Err(); err != nil { - return nil, err - } - return nil, nil - } - line := scanner.Bytes() - if len(line) == 0 { - return readNext() - } - return line, nil - } - - mapping := defaultMapping() - if mappingPath != "" { - data, err := os.ReadFile(mappingPath) - if err != nil { - fmt.Fprintf(os.Stderr, "cannot read mapping: %s\n", err) - os.Exit(2) - } - mapping, err = seq.ReadMapping(data) - if err != nil { - fmt.Fprintf(os.Stderr, "cannot parse mapping: %s\n", err) - os.Exit(2) - } - } - tokenizers := map[seq.TokenizerType]tokenizer.Tokenizer{ - seq.TokenizerTypeKeyword: tokenizer.NewKeywordTokenizer(20, false, true), - seq.TokenizerTypeText: tokenizer.NewTextTokenizer(20, false, true, 100), - seq.TokenizerTypePath: tokenizer.NewPathTokenizer(512, false, true), - seq.TokenizerTypeExists: tokenizer.NewExistsTokenizer(), - } - - activeIndexer, stopIndexer := frac.NewActiveIndexer(4, 10) - defer stopIndexer() - - active := frac.NewActive( - fracName, - activeIndexer, - storage.NewReadLimiter(1, nil), - cache.NewCache[[]byte](nil, nil), - cache.NewCache[[]byte](nil, nil), - &frac.Config{}, - stubSkipMaskProvider{}, - ) - - proc := indexer.NewProcessor(mapping, tokenizers, 0, 0, 0) - compressor := indexer.GetDocsMetasCompressor(3, 3) - - _, binaryDocs, binaryMeta, err := proc.ProcessBulk(time.Now(), nil, nil, readNext) - if err != nil { - fmt.Fprintf(os.Stderr, "process bulk: %s\n", err) - os.Exit(2) - } - - compressor.CompressDocsAndMetas(binaryDocs, binaryMeta) - docsBlock, metasBlock := compressor.DocsMetas() - - var wg sync.WaitGroup - wg.Add(1) - if err := active.Append(docsBlock, metasBlock, &wg); err != nil { - fmt.Fprintf(os.Stderr, "append: %s\n", err) - os.Exit(2) - } - wg.Wait() - - sealParams := common.SealParams{ - IDsZstdLevel: 1, - LIDsZstdLevel: 1, - TokenListZstdLevel: 1, - DocsPositionsZstdLevel: 1, - TokenTableZstdLevel: 1, - DocBlocksZstdLevel: 1, - LIDBlockSize: 256, - TokenBlockSize: 128, - DocBlockSize: 128 * int(units.KiB), - } - - src, err := frac.NewActiveSealingSource(active, sealParams) - if err != nil { - fmt.Fprintf(os.Stderr, "sealing source: %s\n", err) - os.Exit(2) - } - - if _, err := sealing.Seal(src, sealParams); err != nil { - fmt.Fprintf(os.Stderr, "seal: %s\n", err) - os.Exit(2) - } - - active.Release() - - files, _ := filepath.Glob(fracName + "*") - fmt.Fprintf(os.Stderr, "sealed: %v\n", files) -} - -func defaultMapping() seq.Mapping { - return seq.Mapping{ - "level": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - "service": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - "message": seq.NewSingleType(seq.TokenizerTypeText, "", 0), - } -} From b46bf30a7e56fe977515d44d7c72c5b6ddd14453 Mon Sep 17 00:00:00 2001 From: Daniil Lyubaev Date: Wed, 7 Oct 2026 13:46:39 +0300 Subject: [PATCH 7/7] chore: adapt legacy fraction sealing to search tests --- frac/fraction_test.go | 181 +++++++++++++++++++++++++++++++++++++----- 1 file changed, 163 insertions(+), 18 deletions(-) diff --git a/frac/fraction_test.go b/frac/fraction_test.go index e160e5f3..b465fb4f 100644 --- a/frac/fraction_test.go +++ b/frac/fraction_test.go @@ -5,10 +5,12 @@ import ( cryptorand "crypto/rand" "encoding/hex" "fmt" + "io" "math" "math/rand/v2" "net/http/httptest" "os" + "os/exec" "path/filepath" "slices" "strings" @@ -20,10 +22,13 @@ import ( "github.com/alecthomas/units" "github.com/johannesboyne/gofakes3" "github.com/johannesboyne/gofakes3/backend/s3mem" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" "github.com/stretchr/testify/suite" "github.com/ozontech/seq-db/cache" "github.com/ozontech/seq-db/compaction" + "github.com/ozontech/seq-db/config" "github.com/ozontech/seq-db/frac" "github.com/ozontech/seq-db/frac/common" "github.com/ozontech/seq-db/frac/processor" @@ -47,11 +52,75 @@ func (testSkipMaskProvider) GetIDsBitmapByFrac(fracName string, minLID, maxLID u } func (testSkipMaskProvider) RemoveFrac(_ string) {} +func testMapping() seq.Mapping { + return seq.Mapping{ + "id": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + "k8s_pod": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + "k8s_namespace": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + "k8s_container": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + "message": seq.NewSingleType(seq.TokenizerTypeText, "", 0), + "level": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + "client_ip": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + "service": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + "pod": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + "status": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + "source": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + "trace_id": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + "request_uri": seq.NewSingleType(seq.TokenizerTypePath, "", 0), + "spans": seq.NewSingleType(seq.TokenizerTypeNested, "", 0), + "spans.span_id": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + "v": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), + } +} + +var testMappingYAML = []byte(`mapping-list: + - type: keyword + name: client_ip + - type: keyword + name: id + - type: keyword + name: k8s_container + - type: keyword + name: k8s_namespace + - type: keyword + name: k8s_pod + - type: keyword + name: level + - type: text + name: message + - type: keyword + name: pod + - type: path + name: request_uri + - type: keyword + name: service + - type: keyword + name: source + - type: nested + name: spans + mapping-list: + - type: keyword + name: span_id + - type: keyword + name: status + - type: keyword + name: trace_id + - type: keyword + name: v +`) + +func TestMappingYAMLMatchesMapping(t *testing.T) { + back, err := seq.ReadMapping(testMappingYAML) + require.NoError(t, err) + assert.Equal(t, testMapping(), back) +} + type FractionTestSuite struct { suite.Suite tmpDir string config *frac.Config mapping seq.Mapping + mappingYAML []byte tokenizers map[seq.TokenizerType]tokenizer.Tokenizer activeIndexer *frac.ActiveIndexer stopIndexer func() @@ -83,24 +152,8 @@ func (s *FractionTestSuite) SetupTestCommon() { seq.TokenizerTypeText: tokenizer.NewTextTokenizer(20, false, true, 100), seq.TokenizerTypePath: tokenizer.NewPathTokenizer(512, false, true), } - s.mapping = seq.Mapping{ - "id": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - "k8s_pod": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - "k8s_namespace": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - "k8s_container": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - "message": seq.NewSingleType(seq.TokenizerTypeText, "", 0), - "level": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - "client_ip": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - "service": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - "pod": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - "status": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - "source": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - "trace_id": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - "request_uri": seq.NewSingleType(seq.TokenizerTypePath, "", 0), - "spans": seq.NewSingleType(seq.TokenizerTypeNested, "", 0), - "spans.span_id": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - "v": seq.NewSingleType(seq.TokenizerTypeKeyword, "", 0), - } + s.mapping = testMapping() + s.mappingYAML = testMappingYAML s.sealParams = common.SealParams{ IDsZstdLevel: 1, LIDsZstdLevel: 1, @@ -3069,6 +3122,90 @@ func (s *CompactedFractionTestSuite) TestFractionInfo() { s.Require().True(info.IndexOnDisk > 0, "index on disk should be non-zero") } +/* +LegacySealedLoadedFractionTestSuite run tests for legacy sealed fraction +*/ +type LegacySealedLoadedFractionTestSuite struct { + FractionTestSuite + + Version config.BinaryDataVersion +} + +func (s *LegacySealedLoadedFractionTestSuite) SetupSuite() { + s.SetupSuiteCommon() +} + +func (s *LegacySealedLoadedFractionTestSuite) SetupTest() { + s.SetupTestCommon() + + s.insertDocuments = func(bulks ...[]string) { + if s.fraction != nil { + s.Require().Fail("can insert docs only once") + } + s.fraction = s.newLegacySealedLoaded(bulks...) + } +} + +func (s *LegacySealedLoadedFractionTestSuite) TearDownTest() { + if sealed, ok := s.fraction.(*frac.Sealed); ok { + sealed.Release() + } else { + s.Require().Nil(s.fraction, "fraction is not of Sealed type") + } + + s.TearDownTestCommon() +} + +func (s *LegacySealedLoadedFractionTestSuite) TearDownSuite() { + s.TearDownSuiteCommon() +} + +func (s *LegacySealedLoadedFractionTestSuite) sealDocs(bulks ...[]string) string { + fracName := filepath.Join(s.T().TempDir(), fmt.Sprintf("frac_v%d", s.Version)) + + mappingPath := filepath.Join(s.T().TempDir(), "mapping.yaml") + err := os.WriteFile(mappingPath, s.mappingYAML, 0o600) + s.Require().NoError(err) + + cmd := exec.Command( + "bash", + "../cmd/fraction/testdata/legacy/seal-fraction.sh", + "--mapping="+mappingPath, + fmt.Sprintf("v%d", s.Version), + fracName, + ) + cmd.Stdin = readerFromBulks(bulks...) + out, err := cmd.CombinedOutput() + s.Require().NoError(err, "seal-fraction.sh failed:\n%s", out) + + return fracName +} + +func (s *LegacySealedLoadedFractionTestSuite) newLegacySealedLoaded(bulks ...[]string) *frac.Sealed { + filename := s.sealDocs(bulks...) + + sealed := frac.NewSealed( + filename, + storage.NewReadLimiter(1, nil), + frac.NewIndexCache(), + cache.NewConcurrentCache[[]byte](nil, nil), + nil, + s.config, + testSkipMaskProvider{}, + ) + + s.fraction = sealed + return sealed +} + +func readerFromBulks(bulks ...[]string) io.Reader { + var readers []io.Reader + for _, batch := range bulks { + readers = append(readers, strings.NewReader(strings.Join(batch, "\n")+"\n")) + } + return io.MultiReader(readers...) +} + func TestActiveFractionTestSuite(t *testing.T) { suite.Run(t, new(ActiveFractionTestSuite)) } @@ -3092,3 +3229,11 @@ func TestRemoteFractionTestSuite(t *testing.T) { func TestCompactedFractionTestSuite(t *testing.T) { suite.Run(t, new(CompactedFractionTestSuite)) } + +func TestLegacySealedFractionTestSuite(t *testing.T) { + for ver := config.BinaryDataV2; ver < config.CurrentFracVersion; ver++ { + t.Run(fmt.Sprintf("v%d", ver), func(t *testing.T) { + suite.Run(t, &LegacySealedLoadedFractionTestSuite{Version: config.BinaryDataVersion(ver)}) + }) + } +}