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)}) + }) + } +}