From 9c547a53ed8dbcaba36eaa405d2e7da8e03c970b Mon Sep 17 00:00:00 2001 From: William Zahary Henderson Date: Mon, 15 Dec 2025 08:48:53 +0000 Subject: [PATCH] Refactor and multitenant compactor --- cmd/thanos/compact.go | 728 +++++++++++++++++++++---------- cmd/thanos/compact_test.go | 854 +++++++++++++++++++++++++++++++++++++ 2 files changed, 1347 insertions(+), 235 deletions(-) create mode 100644 cmd/thanos/compact_test.go diff --git a/cmd/thanos/compact.go b/cmd/thanos/compact.go index 1530c5cf5b0..2311723ac73 100644 --- a/cmd/thanos/compact.go +++ b/cmd/thanos/compact.go @@ -24,8 +24,10 @@ import ( "github.com/prometheus/client_golang/prometheus/promauto" "github.com/prometheus/common/model" "github.com/prometheus/common/route" + "github.com/prometheus/prometheus/model/relabel" "github.com/prometheus/prometheus/storage" "github.com/prometheus/prometheus/tsdb" + "gopkg.in/yaml.v2" "github.com/thanos-io/objstore" "github.com/thanos-io/objstore/client" @@ -99,6 +101,11 @@ func registerCompact(app *extkingpin.App) { }) } +func extractOrdinalFromHostname(hostname string) (int, error) { + parts := strings.Split(hostname, "-") + return strconv.Atoi(parts[len(parts)-1]) +} + type compactMetrics struct { halted prometheus.Gauge retried prometheus.Counter @@ -179,6 +186,16 @@ func runCompact( progressRegistry := compact.NewProgressRegistry(reg, logger) downsampleMetrics := newDownsampleMetrics(reg) + ctx, cancel := context.WithCancel(context.Background()) + ctx = tracing.ContextWithTracer(ctx, tracer) + ctx = objstoretracing.ContextWithTracer(ctx, tracer) + + defer func() { + if rerr != nil { + cancel() + } + }() + httpProbe := prober.NewHTTP() statusProber := prober.Combine( httpProbe, @@ -193,12 +210,10 @@ func runCompact( g.Add(func() error { statusProber.Healthy() - return srv.ListenAndServe() }, func(err error) { statusProber.NotReady(err) defer statusProber.NotHealthy(err) - srv.Shutdown(err) }) @@ -207,33 +222,256 @@ func runCompact( return err } - bkt, err := client.NewBucket(logger, confContentYaml, component.String(), nil) - if conf.enableFolderDeletion { - bkt, err = block.WrapWithAzDataLakeSdk(logger, confContentYaml, bkt) - level.Info(logger).Log("msg", "azdatalake sdk wrapper enabled", "name", bkt.Name()) + var initialBucketConf client.BucketConfig + if err := yaml.Unmarshal(confContentYaml, &initialBucketConf); err != nil { + return errors.Wrap(err, "failed to parse initial bucket configuration") } - if err != nil { - return err + + // Setup tenant partitioning if enabled + var tenantPrefixes []string + var isMultiTenant bool + + if conf.enableTenantPathPrefix { + isMultiTenant = true + + hostname := os.Getenv("HOSTNAME") + ordinal, err := extractOrdinalFromHostname(hostname) + if err != nil { + return errors.Wrapf(err, "failed to extract ordinal from hostname %s", hostname) + } + + totalShards := conf.replicas / conf.replicationFactor + if conf.replicas%conf.replicationFactor != 0 || conf.replicationFactor <= 0 { + return errors.Errorf("replicas %d must be divisible by replication factor %d and total shards must be greater than 0", conf.replicas, conf.replicationFactor) + } + + if ordinal >= totalShards { + return errors.Errorf("ordinal %d is greater than total shards %d", ordinal, totalShards) + } + + tenantWeightsPath := conf.tenantWeightsFile.Path() + if tenantWeightsPath == "" { + return errors.New("tenant weights file is not set") + } + + discoveryBkt, err := client.NewBucket(logger, confContentYaml, component.String(), nil) + if err != nil { + return errors.Wrapf(err, "failed to create discovery bucket") + } + defer runutil.CloseWithLogOnErr(logger, discoveryBkt, "discovery bucket") + + level.Info(logger).Log("msg", "setting up tenant partitioning", "ordinal", ordinal, "total_shards", totalShards) + + tenantAssignments, err := compact.SetupTenantPartitioning(ctx, discoveryBkt, logger, tenantWeightsPath, conf.commonPathPrefix, totalShards) + if err != nil { + return errors.Wrap(err, "failed to setup tenant partitioning") + } + + assignedTenants := tenantAssignments[ordinal] + if len(assignedTenants) == 0 { + level.Warn(logger).Log("msg", "no tenants assigned to this shard", "ordinal", ordinal) + } + + for _, tenant := range assignedTenants { + tenantPrefixes = append(tenantPrefixes, path.Join(conf.commonPathPrefix, tenant)) + } + + level.Info(logger).Log("msg", "tenant partitioning setup complete", "tenant_prefixes", strings.Join(tenantPrefixes, ",")) + } else { + isMultiTenant = false + tenantPrefixes = []string{""} + level.Info(logger).Log("msg", "single tenant mode") } - insBkt := objstoretracing.WrapWithTraces(objstore.WrapWithMetrics(bkt, extprom.WrapRegistererWithPrefix("thanos_", reg), bkt.Name())) + // Parse relabel config once (doesn't depend on tenant) relabelContentYaml, err := conf.selectorRelabelConf.Content() if err != nil { return errors.Wrap(err, "get content of relabel configuration") } - relabelConfig, err := block.ParseRelabelConfig(relabelContentYaml, block.SelectorSupportedRelabelActions) if err != nil { return err } - // Ensure we close up everything properly. + // Get retention policies once (doesn't depend on tenant) + retentionByResolution, retentionByTenant, err := getRetentionPolicies(logger, &conf) + if err != nil { + return errors.Wrap(err, "failed to get retention policies") + } + + // Get compaction levels once (doesn't depend on tenant) + levels, err := getCompactionLevels(logger, &conf) + if err != nil { + return errors.Wrap(err, "failed to get compaction levels") + } + + // Check vertical compaction settings once + enableVerticalCompaction, dedupReplicaLabels := checkVerticalCompaction(logger, &conf) + + // Get merge function once + mergeFunc, err := getMergeFunc(logger, &conf, dedupReplicaLabels) + if err != nil { + return err + } + + // Create working directories once + compactDir := path.Join(conf.dataDir, "compact") + downsamplingDir := path.Join(conf.dataDir, "downsample") + if err := os.MkdirAll(compactDir, os.ModePerm); err != nil { + return errors.Wrap(err, "create working compact directory") + } + if err := os.MkdirAll(downsamplingDir, os.ModePerm); err != nil { + return errors.Wrap(err, "create working downsample directory") + } + + // Create global bucket for web UI + globalBkt, err := client.NewBucket(logger, confContentYaml, component.String(), nil) + if err != nil { + return errors.Wrap(err, "failed to create global bucket") + } + globalInsBkt := objstoretracing.WrapWithTraces(objstore.WrapWithMetrics(globalBkt, extprom.WrapRegistererWithPrefix("thanos_", reg), globalBkt.Name())) + defer func() { - if err != nil { - runutil.CloseWithLogOnErr(logger, insBkt, "bucket client") + if rerr != nil { + runutil.CloseWithLogOnErr(logger, globalInsBkt, "global bucket client") } }() + // Create global block lister and base meta fetcher for web UI + globalBlockLister := getBlockLister(logger, &conf, globalInsBkt) + globalBaseMetaFetcher, err := block.NewBaseFetcher(logger, conf.blockMetaFetchConcurrency, globalInsBkt, globalBlockLister, conf.dataDir, extprom.WrapRegistererWithPrefix("thanos_", reg)) + if err != nil { + return errors.Wrap(err, "create global base meta fetcher") + } + + // Create global API for web UI + globalAPI := blocksAPI.NewBlocksAPI(logger, conf.webConf.disableCORS, conf.label, flagsMap, globalInsBkt) + + // Start compaction for each tenant + for _, tenantPrefix := range tenantPrefixes { + bucketConf := &client.BucketConfig{ + Type: initialBucketConf.Type, + Config: initialBucketConf.Config, + Prefix: path.Join(initialBucketConf.Prefix, tenantPrefix), + } + level.Info(logger).Log("msg", "starting compaction loop", "prefix", bucketConf.Prefix) + + tenantConfYaml, err := yaml.Marshal(bucketConf) + if err != nil { + return errors.Wrap(err, "failed to marshal tenant bucket configuration") + } + + bkt, err := client.NewBucket(logger, tenantConfYaml, component.String(), nil) + if conf.enableFolderDeletion { + bkt, err = block.WrapWithAzDataLakeSdk(logger, tenantConfYaml, bkt) + level.Info(logger).Log("msg", "azdatalake sdk wrapper enabled", "prefix", bucketConf.Prefix, "name", bkt.Name()) + } + if err != nil { + return errors.Wrap(err, "failed to create tenant bucket") + } + + var tenantReg prometheus.Registerer + if isMultiTenant { + tenantReg = prometheus.WrapRegistererWith(prometheus.Labels{"tenant": tenantPrefix}, reg) + } else { + tenantReg = reg + } + insBkt := objstoretracing.WrapWithTraces(objstore.WrapWithMetrics(bkt, extprom.WrapRegistererWithPrefix("thanos_", reg), bkt.Name())) + + var tenantLogger log.Logger + if isMultiTenant { + tenantLogger = log.With(logger, "tenant", tenantPrefix) + } else { + tenantLogger = logger + } + + // Ensure we close up everything properly. + defer func() { + if err != nil { + runutil.CloseWithLogOnErr(tenantLogger, insBkt, "bucket client") + } + }() + + err = runCompactForTenant( + ctx, + g, + tenantLogger, + tracer, + tenantReg, + reg, + insBkt, + deleteDelay, + &conf, + relabelConfig, + retentionByResolution, + retentionByTenant, + levels, + dedupReplicaLabels, + enableVerticalCompaction, + mergeFunc, + compactDir, + downsamplingDir, + compactMetrics, + progressRegistry, + downsampleMetrics, + flagsMap, + cancel, + ) + if err != nil { + return errors.Wrap(err, "failed to run compact for tenant") + } + } + + // Post-process: set up web UI + err = postProcess( + ctx, + g, + logger, + tracer, + srv, + &conf, + reg, + component, + progressRegistry, + globalBaseMetaFetcher, + globalAPI, + statusProber, + cancel, + ) + if err != nil { + return errors.Wrap(err, "failed to post process") + } + + level.Info(logger).Log("msg", "starting compact node") + statusProber.Ready() + return nil +} + +func runCompactForTenant( + ctx context.Context, + g *run.Group, + logger log.Logger, + tracer opentracing.Tracer, + reg prometheus.Registerer, + baseReg *prometheus.Registry, + insBkt objstore.InstrumentedBucket, + deleteDelay time.Duration, + conf *compactConfig, + relabelConfig []*relabel.Config, + retentionByResolution map[compact.ResolutionLevel]time.Duration, + retentionByTenant map[string]compact.RetentionPolicy, + levels []int64, + dedupReplicaLabels []string, + enableVerticalCompaction bool, + mergeFunc storage.VerticalChunkSeriesMergeFunc, + compactDir string, + downsamplingDir string, + compactMetrics *compactMetrics, + progressRegistry *compact.ProgressRegistry, + downsampleMetrics *DownsampleMetrics, + flagsMap map[string]string, + cancel context.CancelFunc, +) error { // While fetching blocks, we filter out blocks that were marked for deletion by using IgnoreDeletionMarkFilter. // The delay of deleteDelay/2 is added to ensure we fetch blocks that are meant to be deleted but do not have a replacement yet. // This is to make sure compactor will not accidentally perform compactions with gap instead. @@ -245,37 +483,17 @@ func runCompact( consistencyDelayMetaFilter := block.NewConsistencyDelayMetaFilter(logger, conf.consistencyDelay, extprom.WrapRegistererWithPrefix("thanos_", reg)) timePartitionMetaFilter := block.NewTimePartitionMetaFilter(conf.filterConf.MinTime, conf.filterConf.MaxTime) - var blockLister block.Lister - switch syncStrategy(conf.blockListStrategy) { - case concurrentDiscovery: - blockLister = block.NewConcurrentLister(logger, insBkt) - case recursiveDiscovery: - blockLister = block.NewRecursiveLister(logger, insBkt) - default: - return errors.Errorf("unknown sync strategy %s", conf.blockListStrategy) - } + blockLister := getBlockLister(logger, conf, insBkt) baseMetaFetcher, err := block.NewBaseFetcher(logger, conf.blockMetaFetchConcurrency, insBkt, blockLister, conf.dataDir, extprom.WrapRegistererWithPrefix("thanos_", reg)) if err != nil { - return errors.Wrap(err, "create meta fetcher") + return errors.Wrap(err, "create base meta fetcher") } - enableVerticalCompaction := conf.enableVerticalCompaction - dedupReplicaLabels := strutil.ParseFlagLabels(conf.dedupReplicaLabels) - if len(dedupReplicaLabels) > 0 { - enableVerticalCompaction = true - level.Info(logger).Log( - "msg", "deduplication.replica-label specified, enabling vertical compaction", "dedupReplicaLabels", strings.Join(dedupReplicaLabels, ","), - ) - } - if enableVerticalCompaction { - level.Info(logger).Log( - "msg", "vertical compaction is enabled", "compact.enable-vertical-compaction", fmt.Sprintf("%v", conf.enableVerticalCompaction), - ) - } - var ( - api = blocksAPI.NewBlocksAPI(logger, conf.webConf.disableCORS, conf.label, flagsMap, insBkt) - sy *compact.Syncer - ) + // Create tenant API + tenantAPI := blocksAPI.NewBlocksAPI(logger, conf.webConf.disableCORS, conf.label, flagsMap, insBkt) + + // Create syncer + var sy *compact.Syncer { filters := []block.MetadataFilter{ timePartitionMetaFilter, @@ -290,10 +508,9 @@ func runCompact( filters = append(filters, noDownsampleMarkerFilter) } // Make sure all compactor meta syncs are done through Syncer.SyncMeta for readability. - cf := baseMetaFetcher.NewMetaFetcher( - extprom.WrapRegistererWithPrefix("thanos_", reg), filters) + cf := baseMetaFetcher.NewMetaFetcher(extprom.WrapRegistererWithPrefix("thanos_", reg), filters) cf.UpdateOnChange(func(blocks []metadata.Meta, err error) { - api.SetLoaded(blocks, err) + tenantAPI.SetLoaded(blocks, err) }) // Still use blockViewerSyncBlockTimeout to retain original behavior before this upstream change: @@ -303,6 +520,7 @@ func runCompact( if !conf.wait { syncMetasTimeout = 0 } + sy, err = compact.NewMetaSyncer( logger, reg, @@ -319,40 +537,7 @@ func runCompact( } } - levels, err := compactions.levels(conf.maxCompactionLevel) - if err != nil { - return errors.Wrap(err, "get compaction levels") - } - - if conf.maxCompactionLevel < compactions.maxLevel() { - level.Warn(logger).Log("msg", "Max compaction level is lower than should be", "current", conf.maxCompactionLevel, "default", compactions.maxLevel()) - } - - ctx, cancel := context.WithCancel(context.Background()) - ctx = tracing.ContextWithTracer(ctx, tracer) - ctx = objstoretracing.ContextWithTracer(ctx, tracer) // objstore tracing uses a different tracer key in context. - - defer func() { - if rerr != nil { - cancel() - } - }() - - var mergeFunc storage.VerticalChunkSeriesMergeFunc - switch conf.dedupFunc { - case compact.DedupAlgorithmPenalty: - mergeFunc = dedup.NewChunkSeriesMerger() - - if len(dedupReplicaLabels) == 0 { - return errors.New("penalty based deduplication needs at least one replica label specified") - } - case "": - mergeFunc = storage.NewCompactingChunkSeriesMerger(storage.ChainedSeriesMerge) - - default: - return errors.Errorf("unsupported deduplication func, got %s", conf.dedupFunc) - } - + // Create TSDB compactor // Instantiate the compactor with different time slices. Timestamps in TSDB // are in milliseconds. comp, err := tsdb.NewLeveledCompactor(ctx, reg, logger, levels, downsample.NewPool(), mergeFunc) @@ -360,19 +545,6 @@ func runCompact( return errors.Wrap(err, "create compactor") } - var ( - compactDir = path.Join(conf.dataDir, "compact") - downsamplingDir = path.Join(conf.dataDir, "downsample") - ) - - if err := os.MkdirAll(compactDir, os.ModePerm); err != nil { - return errors.Wrap(err, "create working compact directory") - } - - if err := os.MkdirAll(downsamplingDir, os.ModePerm); err != nil { - return errors.Wrap(err, "create working downsample directory") - } - grouper := compact.NewDefaultGrouper( logger, insBkt, @@ -401,14 +573,14 @@ func runCompact( planner = largeIndexFilterPlanner } blocksCleaner := compact.NewBlocksCleaner(logger, insBkt, ignoreDeletionMarkFilter, deleteDelay, compactMetrics.blocksCleaned, compactMetrics.blockCleanupFailures) - compactor, err := compact.NewBucketCompactorWithCheckerAndCallback( + bucketCompactor, err := compact.NewBucketCompactorWithCheckerAndCallback( logger, sy, grouper, planner, comp, compact.DefaultBlockDeletableChecker{}, - compact.NewOverlappingCompactionLifecycleCallback(reg, logger, conf.enableOverlappingRemoval), + compact.NewOverlappingCompactionLifecycleCallback(baseReg, logger, conf.enableOverlappingRemoval), compactDir, insBkt, conf.compactionConcurrency, @@ -418,36 +590,6 @@ func runCompact( return errors.Wrap(err, "create bucket compactor") } - retentionByResolution := map[compact.ResolutionLevel]time.Duration{ - compact.ResolutionLevelRaw: time.Duration(conf.retentionRaw), - compact.ResolutionLevel5m: time.Duration(conf.retentionFiveMin), - compact.ResolutionLevel1h: time.Duration(conf.retentionOneHr), - } - - if retentionByResolution[compact.ResolutionLevelRaw].Milliseconds() != 0 { - // If downsampling is enabled, error if raw retention is not sufficient for downsampling to occur (upper bound 10 days for 1h resolution) - if !conf.disableDownsampling && retentionByResolution[compact.ResolutionLevelRaw].Milliseconds() < downsample.ResLevel1DownsampleRange { - return errors.New("raw resolution must be higher than the minimum block size after which 5m resolution downsampling will occur (40 hours)") - } - level.Info(logger).Log("msg", "retention policy of raw samples is enabled", "duration", retentionByResolution[compact.ResolutionLevelRaw]) - } - if retentionByResolution[compact.ResolutionLevel5m].Milliseconds() != 0 { - // If retention is lower than minimum downsample range, then no downsampling at this resolution will be persisted - if !conf.disableDownsampling && retentionByResolution[compact.ResolutionLevel5m].Milliseconds() < downsample.ResLevel2DownsampleRange { - return errors.New("5m resolution retention must be higher than the minimum block size after which 1h resolution downsampling will occur (10 days)") - } - level.Info(logger).Log("msg", "retention policy of 5 min aggregated samples is enabled", "duration", retentionByResolution[compact.ResolutionLevel5m]) - } - if retentionByResolution[compact.ResolutionLevel1h].Milliseconds() != 0 { - level.Info(logger).Log("msg", "retention policy of 1 hour aggregated samples is enabled", "duration", retentionByResolution[compact.ResolutionLevel1h]) - } - - retentionByTenant, err := compact.ParesRetentionPolicyByTenant(logger, *conf.retentionTenants) - if err != nil { - level.Error(logger).Log("msg", "failed to parse retention policy by tenant", "err", err) - return err - } - var cleanMtx sync.Mutex // TODO(GiedriusS): we could also apply retention policies here but the logic would be a bit more complex. cleanPartialMarked := func(progress *compact.Progress) error { @@ -485,7 +627,7 @@ func runCompact( return errors.Wrap(err, "retention by tenant failed") } - if err := compactor.Compact(ctx, progress); err != nil { + if err := bucketCompactor.Compact(ctx, progress); err != nil { return errors.Wrap(err, "whole compaction error") } @@ -574,6 +716,7 @@ func runCompact( return cleanPartialMarked(progress) } + // Add main compaction goroutine g.Add(func() error { defer runutil.CloseWithLogOnErr(logger, insBkt, "bucket client") @@ -616,146 +759,242 @@ func runCompact( cancel() }) - if conf.wait { - if !conf.disableWeb { - r := route.New() - - ins := extpromhttp.NewInstrumentationMiddleware(reg, nil) + // Periodically remove partial blocks and blocks marked for deletion + // since one iteration potentially could take a long time. + if conf.wait && conf.cleanupBlocksInterval > 0 { + g.Add(func() error { + return runutil.Repeat(conf.cleanupBlocksInterval, ctx.Done(), func() error { + err := cleanPartialMarked(progressRegistry.Get(compact.Cleanup)) + if err != nil && compact.IsRetryError(err) { + // The RetryError signals that we hit an retriable error (transient error, no connection). + // You should alert on this being triggered too frequently. + level.Error(logger).Log("msg", "retriable error", "err", err) + compactMetrics.retried.Inc() - global := ui.NewBucketUI(logger, conf.webConf.externalPrefix, conf.webConf.prefixHeaderName, component) - global.Register(r, ins) - - // Configure Request Logging for HTTP calls. - opts := []logging.Option{logging.WithDecider(func(_ string, _ error) logging.Decision { - return logging.NoLogCall - })} - logMiddleware := logging.NewHTTPServerMiddleware(logger, opts...) - api.Register(r.WithPrefix("/api/v1"), tracer, logger, ins, logMiddleware) + return nil + } - // Separate fetcher for global view. - // TODO(bwplotka): Allow Bucket UI to visualize the state of the block as well. - f := baseMetaFetcher.NewMetaFetcher(extprom.WrapRegistererWithPrefix("thanos_bucket_ui", reg), nil, "component", "globalBucketUI") - f.UpdateOnChange(func(blocks []metadata.Meta, err error) { - api.SetGlobal(blocks, err) + return err }) + }, func(error) { + cancel() + }) + } - srv.Handle("/", r) - - g.Add(func() error { - iterCtx, iterCancel := context.WithTimeout(ctx, conf.blockViewerSyncBlockTimeout) - _, _, _ = f.Fetch(iterCtx) - iterCancel() - - // For /global state make sure to fetch periodically. - return runutil.Repeat(conf.blockViewerSyncBlockInterval, ctx.Done(), func() error { - return runutil.RetryWithLog(logger, time.Minute, ctx.Done(), func() error { - progress := progressRegistry.Get(compact.Web) - defer progress.Idle() - iterCtx, iterCancel := context.WithTimeout(ctx, conf.blockViewerSyncBlockTimeout) - defer iterCancel() - progress.Set(compact.SyncMeta) - _, _, err := f.Fetch(iterCtx) - return err - }) - }) - }, func(error) { - cancel() - }) - } + // Periodically calculate the progress of compaction, downsampling and retention. + if conf.wait && conf.progressCalculateInterval > 0 { + g.Add(func() error { + ps := compact.NewCompactionProgressCalculator(reg, tsdbPlanner) + rs := compact.NewRetentionProgressCalculator(reg, retentionByResolution) + var ds *compact.DownsampleProgressCalculator + if !conf.disableDownsampling { + ds = compact.NewDownsampleProgressCalculator(reg) + } - // Periodically remove partial blocks and blocks marked for deletion - // since one iteration potentially could take a long time. - if conf.cleanupBlocksInterval > 0 { - g.Add(func() error { - return runutil.Repeat(conf.cleanupBlocksInterval, ctx.Done(), func() error { - err := cleanPartialMarked(progressRegistry.Get(compact.Cleanup)) - if err != nil && compact.IsRetryError(err) { - // The RetryError signals that we hit an retriable error (transient error, no connection). - // You should alert on this being triggered too frequently. + return runutil.Repeat(conf.progressCalculateInterval, ctx.Done(), func() error { + progress := progressRegistry.Get(compact.Calculate) + defer progress.Idle() + progress.Set(compact.SyncMeta) + if err := sy.SyncMetas(ctx); err != nil { + // The RetryError signals that we hit an retriable error (transient error, no connection). + // You should alert on this being triggered too frequently. + if compact.IsRetryError(err) { level.Error(logger).Log("msg", "retriable error", "err", err) compactMetrics.retried.Inc() return nil } - return err - }) - }, func(error) { - cancel() - }) - } - - // Periodically calculate the progress of compaction, downsampling and retention. - if conf.progressCalculateInterval > 0 { - g.Add(func() error { - ps := compact.NewCompactionProgressCalculator(reg, tsdbPlanner) - rs := compact.NewRetentionProgressCalculator(reg, retentionByResolution) - var ds *compact.DownsampleProgressCalculator - if !conf.disableDownsampling { - ds = compact.NewDownsampleProgressCalculator(reg) + return errors.Wrapf(err, "could not sync metas") } - return runutil.Repeat(conf.progressCalculateInterval, ctx.Done(), func() error { - progress := progressRegistry.Get(compact.Calculate) - defer progress.Idle() - progress.Set(compact.SyncMeta) - if err := sy.SyncMetas(ctx); err != nil { - // The RetryError signals that we hit an retriable error (transient error, no connection). - // You should alert on this being triggered too frequently. - if compact.IsRetryError(err) { - level.Error(logger).Log("msg", "retriable error", "err", err) - compactMetrics.retried.Inc() + metas := sy.Metas() + progress.Set(compact.Grouping) + groups, err := grouper.Groups(metas) + if err != nil { + return errors.Wrapf(err, "could not group metadata for compaction") + } + progress.Set(compact.CalculateProgress) + if err = ps.ProgressCalculate(ctx, groups); err != nil { + return errors.Wrapf(err, "could not calculate compaction progress") + } - return nil - } + progress.Set(compact.Grouping) + retGroups, err := grouper.Groups(metas) + if err != nil { + return errors.Wrapf(err, "could not group metadata for retention") + } - return errors.Wrapf(err, "could not sync metas") - } + progress.Set(compact.CalculateProgress) + if err = rs.ProgressCalculate(ctx, retGroups); err != nil { + return errors.Wrapf(err, "could not calculate retention progress") + } - metas := sy.Metas() + if !conf.disableDownsampling { progress.Set(compact.Grouping) - groups, err := grouper.Groups(metas) + groups, err = grouper.Groups(metas) if err != nil { - return errors.Wrapf(err, "could not group metadata for compaction") + return errors.Wrapf(err, "could not group metadata into downsample groups") } progress.Set(compact.CalculateProgress) - if err = ps.ProgressCalculate(ctx, groups); err != nil { - return errors.Wrapf(err, "could not calculate compaction progress") + if err := ds.ProgressCalculate(ctx, groups); err != nil { + return errors.Wrapf(err, "could not calculate downsampling progress") } + } - progress.Set(compact.Grouping) - retGroups, err := grouper.Groups(metas) - if err != nil { - return errors.Wrapf(err, "could not group metadata for retention") - } + return nil + }) + }, func(error) { + cancel() + }) + } - progress.Set(compact.CalculateProgress) - if err = rs.ProgressCalculate(ctx, retGroups); err != nil { - return errors.Wrapf(err, "could not calculate retention progress") - } + return nil +} - if !conf.disableDownsampling { - progress.Set(compact.Grouping) - groups, err = grouper.Groups(metas) - if err != nil { - return errors.Wrapf(err, "could not group metadata into downsample groups") - } - progress.Set(compact.CalculateProgress) - if err := ds.ProgressCalculate(ctx, groups); err != nil { - return errors.Wrapf(err, "could not calculate downsampling progress") - } - } +func getBlockLister(logger log.Logger, conf *compactConfig, insBkt objstore.InstrumentedBucketReader) block.Lister { + switch syncStrategy(conf.blockListStrategy) { + case concurrentDiscovery: + return block.NewConcurrentLister(logger, insBkt) + case recursiveDiscovery: + return block.NewRecursiveLister(logger, insBkt) + default: + return block.NewConcurrentLister(logger, insBkt) + } +} - return nil - }) - }, func(err error) { - cancel() - }) +func getCompactionLevels(logger log.Logger, conf *compactConfig) ([]int64, error) { + levels, err := compactions.levels(conf.maxCompactionLevel) + if err != nil { + return nil, errors.Wrap(err, "get compaction levels") + } + + if conf.maxCompactionLevel < compactions.maxLevel() { + level.Warn(logger).Log("msg", "Max compaction level is lower than should be", "current", conf.maxCompactionLevel, "default", compactions.maxLevel()) + } + return levels, nil +} + +func checkVerticalCompaction(logger log.Logger, conf *compactConfig) (bool, []string) { + enableVerticalCompaction := conf.enableVerticalCompaction + dedupReplicaLabels := strutil.ParseFlagLabels(conf.dedupReplicaLabels) + if len(dedupReplicaLabels) > 0 { + enableVerticalCompaction = true + level.Info(logger).Log( + "msg", "deduplication.replica-label specified, enabling vertical compaction", "dedupReplicaLabels", strings.Join(dedupReplicaLabels, ","), + ) + } + if enableVerticalCompaction { + level.Info(logger).Log( + "msg", "vertical compaction is enabled", "compact.enable-vertical-compaction", fmt.Sprintf("%v", conf.enableVerticalCompaction), + ) + } + return enableVerticalCompaction, dedupReplicaLabels +} + +func getMergeFunc(logger log.Logger, conf *compactConfig, dedupReplicaLabels []string) (storage.VerticalChunkSeriesMergeFunc, error) { + var mergeFunc storage.VerticalChunkSeriesMergeFunc + switch conf.dedupFunc { + case compact.DedupAlgorithmPenalty: + mergeFunc = dedup.NewChunkSeriesMerger() + if len(dedupReplicaLabels) == 0 { + return nil, errors.New("penalty based deduplication needs at least one replica label specified") } + case "": + mergeFunc = storage.NewCompactingChunkSeriesMerger(storage.ChainedSeriesMerge) + default: + return nil, errors.Errorf("unsupported deduplication func, got %s", conf.dedupFunc) + } + return mergeFunc, nil +} + +func getRetentionPolicies(logger log.Logger, conf *compactConfig) (map[compact.ResolutionLevel]time.Duration, map[string]compact.RetentionPolicy, error) { + retentionByResolution := map[compact.ResolutionLevel]time.Duration{ + compact.ResolutionLevelRaw: time.Duration(conf.retentionRaw), + compact.ResolutionLevel5m: time.Duration(conf.retentionFiveMin), + compact.ResolutionLevel1h: time.Duration(conf.retentionOneHr), } - level.Info(logger).Log("msg", "starting compact node") - statusProber.Ready() + if retentionByResolution[compact.ResolutionLevelRaw].Milliseconds() != 0 { + if !conf.disableDownsampling && retentionByResolution[compact.ResolutionLevelRaw].Milliseconds() < downsample.ResLevel1DownsampleRange { + return nil, nil, errors.New("raw resolution must be higher than the minimum block size after which 5m resolution downsampling will occur (40 hours)") + } + level.Info(logger).Log("msg", "retention policy of raw samples is enabled", "duration", retentionByResolution[compact.ResolutionLevelRaw]) + } + if retentionByResolution[compact.ResolutionLevel5m].Milliseconds() != 0 { + if !conf.disableDownsampling && retentionByResolution[compact.ResolutionLevel5m].Milliseconds() < downsample.ResLevel2DownsampleRange { + return nil, nil, errors.New("5m resolution retention must be higher than the minimum block size after which 1h resolution downsampling will occur (10 days)") + } + level.Info(logger).Log("msg", "retention policy of 5 min aggregated samples is enabled", "duration", retentionByResolution[compact.ResolutionLevel5m]) + } + if retentionByResolution[compact.ResolutionLevel1h].Milliseconds() != 0 { + level.Info(logger).Log("msg", "retention policy of 1 hour aggregated samples is enabled", "duration", retentionByResolution[compact.ResolutionLevel1h]) + } + + retentionByTenant, err := compact.ParesRetentionPolicyByTenant(logger, *conf.retentionTenants) + if err != nil { + level.Error(logger).Log("msg", "failed to parse retention policy by tenant", "err", err) + return nil, nil, errors.Wrap(err, "failed to parse retention policy by tenant") + } + return retentionByResolution, retentionByTenant, nil +} + +func postProcess( + ctx context.Context, + g *run.Group, + logger log.Logger, + tracer opentracing.Tracer, + srv *httpserver.Server, + conf *compactConfig, + reg *prometheus.Registry, + component component.Component, + progressRegistry *compact.ProgressRegistry, + baseMetaFetcher *block.BaseFetcher, + globalAPI *blocksAPI.BlocksAPI, + statusProber prober.Probe, + cancel context.CancelFunc, +) error { + if conf.wait && !conf.disableWeb { + r := route.New() + + ins := extpromhttp.NewInstrumentationMiddleware(reg, nil) + + global := ui.NewBucketUI(logger, conf.webConf.externalPrefix, conf.webConf.prefixHeaderName, component) + global.Register(r, ins) + + opts := []logging.Option{logging.WithDecider(func(_ string, _ error) logging.Decision { + return logging.NoLogCall + })} + logMiddleware := logging.NewHTTPServerMiddleware(logger, opts...) + globalAPI.Register(r.WithPrefix("/api/v1"), tracer, logger, ins, logMiddleware) + + f := baseMetaFetcher.NewMetaFetcher(extprom.WrapRegistererWithPrefix("thanos_bucket_ui", reg), nil, "component", "globalBucketUI") + f.UpdateOnChange(func(blocks []metadata.Meta, err error) { + globalAPI.SetGlobal(blocks, err) + }) + + srv.Handle("/", r) + + g.Add(func() error { + iterCtx, iterCancel := context.WithTimeout(ctx, conf.blockViewerSyncBlockTimeout) + _, _, _ = f.Fetch(iterCtx) + iterCancel() + + return runutil.Repeat(conf.blockViewerSyncBlockInterval, ctx.Done(), func() error { + return runutil.RetryWithLog(logger, time.Minute, ctx.Done(), func() error { + progress := progressRegistry.Get(compact.Web) + defer progress.Idle() + iterCtx, iterCancel := context.WithTimeout(ctx, conf.blockViewerSyncBlockTimeout) + defer iterCancel() + progress.Set(compact.SyncMeta) + _, _, err := f.Fetch(iterCtx) + return err + }) + }) + }, func(error) { + cancel() + }) + } return nil } @@ -797,6 +1036,11 @@ type compactConfig struct { progressCalculateInterval time.Duration filterConf *store.FilterConfig disableAdminOperations bool + tenantWeightsFile extflag.PathOrContent + replicas int + replicationFactor int + commonPathPrefix string + enableTenantPathPrefix bool } func (cc *compactConfig) registerFlag(cmd extkingpin.FlagClause) { @@ -913,6 +1157,20 @@ func (cc *compactConfig) registerFlag(cmd extkingpin.FlagClause) { cc.selectorRelabelConf = *extkingpin.RegisterSelectorRelabelFlags(cmd) + cc.tenantWeightsFile = *extflag.RegisterPathOrContent(cmd, "compact.tenant-weights-file", "YAML file that contains the tenant weights for tenant partitioning.", extflag.WithEnvSubstitution()) + + cmd.Flag("compact.replicas", "Total replicas of the stateful set."). + Default("1").IntVar(&cc.replicas) + + cmd.Flag("compact.replication-factor", "Replication factor of the stateful set."). + Default("1").IntVar(&cc.replicationFactor) + + cmd.Flag("compact.common-path-prefix", "Common path prefix for tenant discovery when using tenant partitioning. This is the prefix before the tenant name in the object storage path."). + Default("v1/raw/").StringVar(&cc.commonPathPrefix) + + cmd.Flag("compact.enable-tenant-path-prefix", "Enable tenant path prefix mode for backward compatibility. When disabled, compactor runs in single-tenant mode."). + Default("false").BoolVar(&cc.enableTenantPathPrefix) + cc.webConf.registerFlag(cmd) cmd.Flag("bucket-web-label", "External block label to use as group title in the bucket web UI").StringVar(&cc.label) diff --git a/cmd/thanos/compact_test.go b/cmd/thanos/compact_test.go new file mode 100644 index 00000000000..603fd32f6a3 --- /dev/null +++ b/cmd/thanos/compact_test.go @@ -0,0 +1,854 @@ +// Copyright (c) The Thanos Authors. +// Licensed under the Apache License 2.0. + +package main + +import ( + "path" + "testing" + "time" + + "github.com/efficientgo/core/testutil" + "github.com/go-kit/log" + "github.com/prometheus/common/model" + "github.com/thanos-io/objstore" + "github.com/thanos-io/objstore/client" + "gopkg.in/yaml.v2" + + "github.com/thanos-io/thanos/pkg/compact" + "github.com/thanos-io/thanos/pkg/store" +) + +func TestExtractOrdinalFromHostname(t *testing.T) { + t.Parallel() + + for _, tcase := range []struct { + name string + hostname string + expectedOrd int + expectedError bool + }{ + { + name: "statefulset hostname with single digit", + hostname: "thanos-compact-0", + expectedOrd: 0, + expectedError: false, + }, + { + name: "statefulset hostname with double digit", + hostname: "thanos-compact-12", + expectedOrd: 12, + expectedError: false, + }, + { + name: "statefulset hostname with triple digit", + hostname: "thanos-compact-123", + expectedOrd: 123, + expectedError: false, + }, + { + name: "kubernetes statefulset with namespace", + hostname: "thanos-compact-shard-5", + expectedOrd: 5, + expectedError: false, + }, + { + name: "complex statefulset name", + hostname: "prod-us-west-thanos-compact-7", + expectedOrd: 7, + expectedError: false, + }, + { + name: "hostname without number", + hostname: "thanos-compact-abc", + expectedOrd: 0, + expectedError: true, + }, + { + name: "hostname with invalid suffix", + hostname: "thanos-compact-", + expectedOrd: 0, + expectedError: true, + }, + { + name: "empty hostname", + hostname: "", + expectedOrd: 0, + expectedError: true, + }, + } { + t.Run(tcase.name, func(t *testing.T) { + ordinal, err := extractOrdinalFromHostname(tcase.hostname) + if tcase.expectedError { + testutil.NotOk(t, err) + } else { + testutil.Ok(t, err) + testutil.Equals(t, tcase.expectedOrd, ordinal) + } + }) + } +} + +func TestTenantPrefixBucketCreation(t *testing.T) { + t.Parallel() + + for _, tcase := range []struct { + name string + initialPrefix string + tenantPrefixes []string + expectedBucketPrefixes []string + }{ + { + name: "single-tenant mode (no prefixes)", + initialPrefix: "", + tenantPrefixes: []string{""}, + expectedBucketPrefixes: []string{""}, + }, + { + name: "tenant partitioning with one tenant", + initialPrefix: "", + tenantPrefixes: []string{"v1/raw/tenant1"}, + expectedBucketPrefixes: []string{"v1/raw/tenant1"}, + }, + { + name: "tenant partitioning with multiple tenants", + initialPrefix: "", + tenantPrefixes: []string{"v1/raw/tenant1", "v1/raw/tenant2", "v1/raw/tenant3"}, + expectedBucketPrefixes: []string{"v1/raw/tenant1", "v1/raw/tenant2", "v1/raw/tenant3"}, + }, + { + name: "tenant partitioning with base prefix and tenant paths", + initialPrefix: "base-prefix", + tenantPrefixes: []string{"v1/raw/tenant1", "v1/raw/tenant2"}, + expectedBucketPrefixes: []string{"base-prefix/v1/raw/tenant1", "base-prefix/v1/raw/tenant2"}, + }, + } { + t.Run(tcase.name, func(t *testing.T) { + initialBucketConf := client.BucketConfig{ + Type: client.FILESYSTEM, + Config: nil, + Prefix: tcase.initialPrefix, + } + + var actualPrefixes []string + for _, tenantPrefix := range tcase.tenantPrefixes { + bucketConf := &client.BucketConfig{ + Type: initialBucketConf.Type, + Config: initialBucketConf.Config, + Prefix: path.Join(initialBucketConf.Prefix, tenantPrefix), + } + actualPrefixes = append(actualPrefixes, bucketConf.Prefix) + } + + testutil.Equals(t, tcase.expectedBucketPrefixes, actualPrefixes) + }) + } +} + +func TestBucketConfigPrefixPreservation(t *testing.T) { + t.Parallel() + + for _, tcase := range []struct { + name string + prefix string + tenantPrefix string + expectedPrefix string + }{ + { + name: "empty prefix", + prefix: "", + tenantPrefix: "", + expectedPrefix: "", + }, + { + name: "simple prefix", + prefix: "data", + tenantPrefix: "", + expectedPrefix: "data", + }, + { + name: "v1/raw prefix", + prefix: "v1/raw", + tenantPrefix: "tenant1", + expectedPrefix: "v1/raw/tenant1", + }, + { + name: "hierarchical prefix", + prefix: "env/prod/region/us-west", + tenantPrefix: "tenant-abc", + expectedPrefix: "env/prod/region/us-west/tenant-abc", + }, + } { + t.Run(tcase.name, func(t *testing.T) { + result := path.Join(tcase.prefix, tcase.tenantPrefix) + testutil.Equals(t, tcase.expectedPrefix, result) + }) + } +} + +func TestBucketConfigFromYAML(t *testing.T) { + t.Parallel() + + for _, tcase := range []struct { + name string + yamlConfig string + expectedType client.ObjProvider + expectedPrefix string + }{ + { + name: "filesystem without prefix", + yamlConfig: ` +type: FILESYSTEM +config: + directory: /tmp/thanos +`, + expectedType: client.FILESYSTEM, + expectedPrefix: "", + }, + { + name: "filesystem with v1/raw prefix", + yamlConfig: ` +type: FILESYSTEM +config: + directory: /tmp/thanos +prefix: v1/raw +`, + expectedType: client.FILESYSTEM, + expectedPrefix: "v1/raw", + }, + { + name: "S3 with data prefix", + yamlConfig: ` +type: S3 +config: + bucket: my-bucket + endpoint: s3.amazonaws.com +prefix: data +`, + expectedType: client.S3, + expectedPrefix: "data", + }, + { + name: "filesystem with tenant path", + yamlConfig: ` +type: FILESYSTEM +config: + directory: /tmp/thanos +prefix: v1/raw/tenant1 +`, + expectedType: client.FILESYSTEM, + expectedPrefix: "v1/raw/tenant1", + }, + } { + t.Run(tcase.name, func(t *testing.T) { + var bucketConf client.BucketConfig + err := yaml.Unmarshal([]byte(tcase.yamlConfig), &bucketConf) + testutil.Ok(t, err) + testutil.Equals(t, tcase.expectedType, bucketConf.Type) + testutil.Equals(t, tcase.expectedPrefix, bucketConf.Prefix) + }) + } +} + +func TestGetBlockLister(t *testing.T) { + t.Parallel() + + logger := log.NewNopLogger() + bkt := objstore.WithNoopInstr(objstore.NewInMemBucket()) + + for _, tcase := range []struct { + name string + blockListStrategy string + expectNonNil bool + }{ + { + name: "concurrent strategy", + blockListStrategy: string(concurrentDiscovery), + expectNonNil: true, + }, + { + name: "recursive strategy", + blockListStrategy: string(recursiveDiscovery), + expectNonNil: true, + }, + { + name: "default strategy", + blockListStrategy: "", + expectNonNil: true, + }, + } { + t.Run(tcase.name, func(t *testing.T) { + conf := &compactConfig{ + blockListStrategy: tcase.blockListStrategy, + } + lister := getBlockLister(logger, conf, bkt) + if tcase.expectNonNil { + testutil.Assert(t, lister != nil, "expected non-nil block lister") + } + }) + } +} + +func TestGetCompactionLevels(t *testing.T) { + t.Parallel() + + logger := log.NewNopLogger() + + for _, tcase := range []struct { + name string + maxCompactionLevel int + expectedLevels int + expectError bool + }{ + { + name: "default max level", + maxCompactionLevel: 4, + expectedLevels: 5, + expectError: false, + }, + { + name: "level 0", + maxCompactionLevel: 0, + expectedLevels: 1, + expectError: false, + }, + { + name: "level 2", + maxCompactionLevel: 2, + expectedLevels: 3, + expectError: false, + }, + { + name: "level too high", + maxCompactionLevel: 10, + expectedLevels: 0, + expectError: true, + }, + } { + t.Run(tcase.name, func(t *testing.T) { + conf := &compactConfig{ + maxCompactionLevel: tcase.maxCompactionLevel, + } + levels, err := getCompactionLevels(logger, conf) + if tcase.expectError { + testutil.NotOk(t, err) + } else { + testutil.Ok(t, err) + testutil.Equals(t, tcase.expectedLevels, len(levels)) + } + }) + } +} + +func TestCheckVerticalCompaction(t *testing.T) { + t.Parallel() + + logger := log.NewNopLogger() + + for _, tcase := range []struct { + name string + enableVerticalCompaction bool + dedupReplicaLabels []string + expectedVerticalCompaction bool + expectedDedupReplicaLabelsCount int + }{ + { + name: "disabled by default", + enableVerticalCompaction: false, + dedupReplicaLabels: nil, + expectedVerticalCompaction: false, + expectedDedupReplicaLabelsCount: 0, + }, + { + name: "explicitly enabled", + enableVerticalCompaction: true, + dedupReplicaLabels: nil, + expectedVerticalCompaction: true, + expectedDedupReplicaLabelsCount: 0, + }, + { + name: "enabled via dedup replica labels", + enableVerticalCompaction: false, + dedupReplicaLabels: []string{"replica"}, + expectedVerticalCompaction: true, + expectedDedupReplicaLabelsCount: 1, + }, + { + name: "multiple dedup replica labels", + enableVerticalCompaction: false, + dedupReplicaLabels: []string{"replica", "prometheus"}, + expectedVerticalCompaction: true, + expectedDedupReplicaLabelsCount: 2, + }, + } { + t.Run(tcase.name, func(t *testing.T) { + conf := &compactConfig{ + enableVerticalCompaction: tcase.enableVerticalCompaction, + dedupReplicaLabels: tcase.dedupReplicaLabels, + } + enabled, labels := checkVerticalCompaction(logger, conf) + testutil.Equals(t, tcase.expectedVerticalCompaction, enabled) + testutil.Equals(t, tcase.expectedDedupReplicaLabelsCount, len(labels)) + }) + } +} + +func TestGetMergeFunc(t *testing.T) { + t.Parallel() + + logger := log.NewNopLogger() + + for _, tcase := range []struct { + name string + dedupFunc string + dedupReplicaLabels []string + expectError bool + }{ + { + name: "default merge func", + dedupFunc: "", + dedupReplicaLabels: nil, + expectError: false, + }, + { + name: "penalty dedup without labels", + dedupFunc: compact.DedupAlgorithmPenalty, + dedupReplicaLabels: nil, + expectError: true, + }, + { + name: "penalty dedup with labels", + dedupFunc: compact.DedupAlgorithmPenalty, + dedupReplicaLabels: []string{"replica"}, + expectError: false, + }, + { + name: "unsupported dedup func", + dedupFunc: "invalid", + dedupReplicaLabels: nil, + expectError: true, + }, + } { + t.Run(tcase.name, func(t *testing.T) { + conf := &compactConfig{ + dedupFunc: tcase.dedupFunc, + } + mergeFunc, err := getMergeFunc(logger, conf, tcase.dedupReplicaLabels) + if tcase.expectError { + testutil.NotOk(t, err) + } else { + testutil.Ok(t, err) + testutil.Assert(t, mergeFunc != nil, "expected non-nil merge func") + } + }) + } +} + +func TestGetRetentionPolicies(t *testing.T) { + t.Parallel() + + logger := log.NewNopLogger() + + for _, tcase := range []struct { + name string + retentionRaw model.Duration + retentionFiveMin model.Duration + retentionOneHr model.Duration + retentionTenants []string + disableDownsampling bool + expectError bool + }{ + { + name: "no retention set", + retentionRaw: 0, + retentionFiveMin: 0, + retentionOneHr: 0, + retentionTenants: nil, + disableDownsampling: false, + expectError: false, + }, + { + name: "valid raw retention", + retentionRaw: model.Duration(30 * 24 * time.Hour), + retentionFiveMin: 0, + retentionOneHr: 0, + retentionTenants: nil, + disableDownsampling: false, + expectError: false, + }, + { + name: "raw retention too low for downsampling", + retentionRaw: model.Duration(1 * time.Hour), + retentionFiveMin: 0, + retentionOneHr: 0, + retentionTenants: nil, + disableDownsampling: false, + expectError: true, + }, + { + name: "raw retention low but downsampling disabled", + retentionRaw: model.Duration(1 * time.Hour), + retentionFiveMin: 0, + retentionOneHr: 0, + retentionTenants: nil, + disableDownsampling: true, + expectError: false, + }, + } { + t.Run(tcase.name, func(t *testing.T) { + tenants := tcase.retentionTenants + conf := &compactConfig{ + retentionRaw: tcase.retentionRaw, + retentionFiveMin: tcase.retentionFiveMin, + retentionOneHr: tcase.retentionOneHr, + retentionTenants: &tenants, + disableDownsampling: tcase.disableDownsampling, + } + byResolution, byTenant, err := getRetentionPolicies(logger, conf) + if tcase.expectError { + testutil.NotOk(t, err) + } else { + testutil.Ok(t, err) + testutil.Assert(t, byResolution != nil, "expected non-nil byResolution") + testutil.Assert(t, byTenant != nil, "expected non-nil byTenant") + } + }) + } +} + +func TestSingleTenantModeConfiguration(t *testing.T) { + t.Parallel() + + // Test that when enableTenantPathPrefix is false, we get single tenant behavior + conf := &compactConfig{ + enableTenantPathPrefix: false, + } + + // Verify single tenant mode configuration + testutil.Equals(t, false, conf.enableTenantPathPrefix) + + // In single tenant mode, tenantPrefixes should be [""] + var tenantPrefixes []string + if conf.enableTenantPathPrefix { + t.Fatal("enableTenantPathPrefix should be false") + } else { + tenantPrefixes = []string{""} + } + + testutil.Equals(t, 1, len(tenantPrefixes)) + testutil.Equals(t, "", tenantPrefixes[0]) +} + +func TestMultiTenantModeConfiguration(t *testing.T) { + t.Parallel() + + // Test that multi-tenant configuration is properly set up + conf := &compactConfig{ + enableTenantPathPrefix: true, + replicas: 6, + replicationFactor: 2, + commonPathPrefix: "v1/raw/", + } + + // Verify multi-tenant mode configuration + testutil.Equals(t, true, conf.enableTenantPathPrefix) + testutil.Equals(t, 6, conf.replicas) + testutil.Equals(t, 2, conf.replicationFactor) + testutil.Equals(t, "v1/raw/", conf.commonPathPrefix) + + // Calculate total shards + totalShards := conf.replicas / conf.replicationFactor + testutil.Equals(t, 3, totalShards) +} + +func TestMultiTenantShardCalculation(t *testing.T) { + t.Parallel() + + for _, tcase := range []struct { + name string + replicas int + replicationFactor int + expectedShards int + expectError bool + }{ + { + name: "6 replicas, 2 replication factor", + replicas: 6, + replicationFactor: 2, + expectedShards: 3, + expectError: false, + }, + { + name: "9 replicas, 3 replication factor", + replicas: 9, + replicationFactor: 3, + expectedShards: 3, + expectError: false, + }, + { + name: "4 replicas, 1 replication factor", + replicas: 4, + replicationFactor: 1, + expectedShards: 4, + expectError: false, + }, + { + name: "non-divisible replicas", + replicas: 5, + replicationFactor: 2, + expectedShards: 0, + expectError: true, + }, + { + name: "zero replication factor", + replicas: 6, + replicationFactor: 0, + expectedShards: 0, + expectError: true, + }, + } { + t.Run(tcase.name, func(t *testing.T) { + if tcase.replicationFactor <= 0 || tcase.replicas%tcase.replicationFactor != 0 { + if !tcase.expectError { + t.Fatal("expected error for invalid configuration") + } + return + } + + totalShards := tcase.replicas / tcase.replicationFactor + if tcase.expectError { + testutil.Assert(t, totalShards <= 0, "expected invalid shard count") + } else { + testutil.Equals(t, tcase.expectedShards, totalShards) + } + }) + } +} + +func TestTenantPrefixGeneration(t *testing.T) { + t.Parallel() + + for _, tcase := range []struct { + name string + commonPathPrefix string + tenants []string + expectedPrefixes []string + }{ + { + name: "v1/raw prefix with tenants", + commonPathPrefix: "v1/raw", + tenants: []string{"tenant1", "tenant2", "tenant3"}, + expectedPrefixes: []string{"v1/raw/tenant1", "v1/raw/tenant2", "v1/raw/tenant3"}, + }, + { + name: "empty prefix with tenants", + commonPathPrefix: "", + tenants: []string{"tenant1"}, + expectedPrefixes: []string{"tenant1"}, + }, + { + name: "nested prefix with tenants", + commonPathPrefix: "org/data/raw", + tenants: []string{"tenant-1", "tenant-2"}, + expectedPrefixes: []string{"org/data/raw/tenant-1", "org/data/raw/tenant-2"}, + }, + } { + t.Run(tcase.name, func(t *testing.T) { + var tenantPrefixes []string + for _, tenant := range tcase.tenants { + tenantPrefixes = append(tenantPrefixes, path.Join(tcase.commonPathPrefix, tenant)) + } + testutil.Equals(t, tcase.expectedPrefixes, tenantPrefixes) + }) + } +} + +func TestCompactConfigDefaults(t *testing.T) { + t.Parallel() + + // Test default values for compactConfig + conf := &compactConfig{} + + // Verify default values + testutil.Equals(t, false, conf.enableTenantPathPrefix) + testutil.Equals(t, 0, conf.replicas) + testutil.Equals(t, 0, conf.replicationFactor) + testutil.Equals(t, "", conf.commonPathPrefix) + testutil.Equals(t, false, conf.enableVerticalCompaction) + testutil.Equals(t, false, conf.disableDownsampling) + testutil.Equals(t, false, conf.disableWeb) + testutil.Equals(t, false, conf.haltOnError) +} + +func TestFilterConfigDefaults(t *testing.T) { + t.Parallel() + + // Verify filter config can be created + conf := &compactConfig{ + filterConf: &store.FilterConfig{}, + } + + testutil.Assert(t, conf.filterConf != nil, "expected non-nil filter config") +} + +func TestBackwardCompatibility_SingleTenantMode(t *testing.T) { + t.Parallel() + + // This test ensures backward compatibility when enableTenantPathPrefix is false + // The compactor should behave as a single-tenant compactor + + conf := &compactConfig{ + enableTenantPathPrefix: false, + dataDir: t.TempDir(), + maxCompactionLevel: 4, + blockMetaFetchConcurrency: 32, + blockFilesConcurrency: 1, + compactionConcurrency: 1, + downsampleConcurrency: 1, + } + + // Verify single tenant mode + testutil.Equals(t, false, conf.enableTenantPathPrefix) + + // In single tenant mode, we should have exactly one tenant prefix (empty string) + var tenantPrefixes []string + var isMultiTenant bool + + if conf.enableTenantPathPrefix { + isMultiTenant = true + } else { + isMultiTenant = false + tenantPrefixes = []string{""} + } + + testutil.Equals(t, false, isMultiTenant) + testutil.Equals(t, 1, len(tenantPrefixes)) + testutil.Equals(t, "", tenantPrefixes[0]) + + // Verify bucket prefix is empty in single tenant mode + initialBucketConf := client.BucketConfig{ + Type: client.FILESYSTEM, + Config: nil, + Prefix: "", + } + + bucketConf := &client.BucketConfig{ + Type: initialBucketConf.Type, + Config: initialBucketConf.Config, + Prefix: path.Join(initialBucketConf.Prefix, tenantPrefixes[0]), + } + + testutil.Equals(t, "", bucketConf.Prefix) +} + +func TestMultiTenantMode_TenantIsolation(t *testing.T) { + t.Parallel() + + // This test ensures that each tenant gets its own isolated bucket prefix + + conf := &compactConfig{ + enableTenantPathPrefix: true, + commonPathPrefix: "v1/raw", + } + + tenants := []string{"tenant-alpha", "tenant-beta", "tenant-gamma"} + + var tenantPrefixes []string + for _, tenant := range tenants { + tenantPrefixes = append(tenantPrefixes, path.Join(conf.commonPathPrefix, tenant)) + } + + // Verify each tenant has a unique, isolated prefix + testutil.Equals(t, 3, len(tenantPrefixes)) + testutil.Equals(t, "v1/raw/tenant-alpha", tenantPrefixes[0]) + testutil.Equals(t, "v1/raw/tenant-beta", tenantPrefixes[1]) + testutil.Equals(t, "v1/raw/tenant-gamma", tenantPrefixes[2]) + + // Verify prefixes are unique + prefixSet := make(map[string]bool) + for _, prefix := range tenantPrefixes { + testutil.Assert(t, !prefixSet[prefix], "duplicate prefix found: "+prefix) + prefixSet[prefix] = true + } +} + +func TestCompactionSetLevels(t *testing.T) { + t.Parallel() + + cs := compactionSet{ + 1 * time.Hour, + 2 * time.Hour, + 8 * time.Hour, + 2 * 24 * time.Hour, + 14 * 24 * time.Hour, + } + + for _, tcase := range []struct { + name string + maxLevel int + expectedLevels int + expectError bool + }{ + { + name: "level 0", + maxLevel: 0, + expectedLevels: 1, + expectError: false, + }, + { + name: "level 2", + maxLevel: 2, + expectedLevels: 3, + expectError: false, + }, + { + name: "max level 4", + maxLevel: 4, + expectedLevels: 5, + expectError: false, + }, + { + name: "level exceeds set", + maxLevel: 5, + expectedLevels: 0, + expectError: true, + }, + } { + t.Run(tcase.name, func(t *testing.T) { + levels, err := cs.levels(tcase.maxLevel) + if tcase.expectError { + testutil.NotOk(t, err) + } else { + testutil.Ok(t, err) + testutil.Equals(t, tcase.expectedLevels, len(levels)) + } + }) + } +} + +func TestCompactionSetMaxLevel(t *testing.T) { + t.Parallel() + + cs := compactionSet{ + 1 * time.Hour, + 2 * time.Hour, + 8 * time.Hour, + 2 * 24 * time.Hour, + 14 * 24 * time.Hour, + } + + testutil.Equals(t, 4, cs.maxLevel()) +} + +func TestCompactionSetString(t *testing.T) { + t.Parallel() + + cs := compactionSet{ + 1 * time.Hour, + 2 * time.Hour, + } + + str := cs.String() + testutil.Assert(t, str != "", "expected non-empty string") + testutil.Assert(t, len(str) > 0, "expected string representation") +}