Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
141 commits
Select commit Hold shift + click to select a range
9d19dbf
WIP: add planner
u-veles-a Jun 24, 2026
331070c
WIP: add planner filters
u-veles-a Jun 24, 2026
704276d
WIP: do group
u-veles-a Jun 25, 2026
12c6b78
WIP: tcompactor
u-veles-a Jun 26, 2026
4ae7a71
WIP: improvements
u-veles-a Jun 26, 2026
be68664
Merge branch 'pp' into tcompactor
u-veles-a Jun 26, 2026
87fdf1b
WIP: add write thanos meta
u-veles-a Jun 29, 2026
b822344
WIP: fork block and lcompactor
u-veles-a Jul 1, 2026
93bc9a3
WIP: copy code
u-veles-a Jul 1, 2026
2c14ad2
WIP: add test
u-veles-a Jul 1, 2026
0896c9c
WIP: add test
u-veles-a Jul 2, 2026
de83da1
Merge branch 'pp' into tcompactor
u-veles-a Jul 2, 2026
58d1f13
move block manager
u-veles-a Jul 2, 2026
bfdfe9a
revert block compactor
u-veles-a Jul 2, 2026
feeef45
WIP: add OverlappingBlocks test
u-veles-a Jul 2, 2026
421afcb
WIP: add test
u-veles-a Jul 2, 2026
c1473bf
add test
u-veles-a Jul 2, 2026
a5bdca0
small fix
u-veles-a Jul 2, 2026
af0e46c
fix mmap, pool, test
u-veles-a Jul 3, 2026
d265d62
Merge branch 'pp' into tcompactor
u-veles-a Jul 6, 2026
e28acbf
remove GatherIndexHealthStats
u-veles-a Jul 6, 2026
58810e2
fix linter
u-veles-a Jul 6, 2026
d2a5dfe
Merge branch 'pp' into tcompactor
u-veles-a Jul 7, 2026
6923a97
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Jul 8, 2026
ba0d923
WIP: need cleaner
u-veles-a Jul 9, 2026
7328bb3
add ExpirationPolicy
u-veles-a Jul 10, 2026
872fcd7
Merge branch 'pp' into tcompactor
u-veles-a Jul 13, 2026
7c4632c
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Jul 13, 2026
6f9a768
fix linter
u-veles-a Jul 13, 2026
648b221
Merge branch 'pp' into tcompactor
u-veles-a Jul 14, 2026
a227346
fix after merge
u-veles-a Jul 14, 2026
18c87ef
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Jul 14, 2026
353613a
WIP: merge
u-veles-a Jul 14, 2026
bbaa160
fix after merge
u-veles-a Jul 14, 2026
093329a
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Jul 14, 2026
eee1f65
Merge branch 'pp' into tcompactor
u-veles-a Jul 15, 2026
98a7640
add test
u-veles-a Jul 15, 2026
f1cb0ba
add test for ExpirationPolicy
u-veles-a Jul 15, 2026
f5ef682
Merge branch 'pp' into tcompactor
u-veles-a Jul 17, 2026
d7e8094
Merge branch 'pp' into tcompactor
u-veles-a Jul 21, 2026
88cd0b2
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Jul 21, 2026
c75a645
fix after merge
u-veles-a Jul 21, 2026
4db19d8
fix test
u-veles-a Jul 21, 2026
69a92b7
tcompactor: review fixes (dead code cleanup) (#428)
vporoshok Jul 22, 2026
c5e93c3
fix select downsampling
u-veles-a Jul 22, 2026
d4bebe3
Merge branch 'pp' into tcompactor
u-veles-a Jul 22, 2026
b42daea
Merge branch 'tcompactor' of github.com:deckhouse/prompp into tcompactor
u-veles-a Jul 22, 2026
1740c9e
fix review
u-veles-a Jul 23, 2026
7148cb6
Merge branch 'pp' into tcompactor
u-veles-a Jul 23, 2026
882351e
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Jul 23, 2026
b4e4076
go mod tidy
u-veles-a Jul 23, 2026
41b6a94
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Jul 23, 2026
0b33c82
go mod tidy
u-veles-a Jul 23, 2026
2d311f4
Merge branch 'pp' into tcompactor
u-veles-a Jul 23, 2026
53f5abb
fix linter
u-veles-a Jul 23, 2026
f5fac02
Merge branch 'pp' into tcompactor
u-veles-a Jul 23, 2026
c1df71a
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Jul 23, 2026
011df26
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Jul 23, 2026
e2dbcb2
fix linter
u-veles-a Jul 23, 2026
f3511d9
Merge branch 'pp' into tcompactor
u-veles-a Jul 23, 2026
9794c3c
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Jul 23, 2026
ef99d99
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Jul 23, 2026
62f28c2
Merge branch 'pp' into tcompactor
u-veles-a Jul 24, 2026
bdd7313
fix linter
u-veles-a Jul 24, 2026
80e28d0
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Jul 24, 2026
4e60da9
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Jul 24, 2026
5ad98f6
Merge branch 'pp' into tcompactor
u-veles-a Jul 24, 2026
40a85bc
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Jul 24, 2026
c4e6e2a
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Jul 24, 2026
ac3815f
fix review
u-veles-a Jul 28, 2026
cb81223
Merge branch 'pp' into tcompactor
u-veles-a Jul 28, 2026
e34e69f
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Jul 28, 2026
ffda194
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Jul 29, 2026
2fe9b81
fix timeInterval
u-veles-a Jul 29, 2026
1c2c79d
fix set_items_count
u-veles-a Aug 3, 2026
91e47b7
Merge branch 'pp' into tcompactor
u-veles-a Aug 6, 2026
ceb0bb5
Merge branch 'pp' into tcompactor
u-veles-a Aug 6, 2026
b4a82f9
go mod tidy
u-veles-a Aug 6, 2026
6ca1c21
Merge branch 'pp' into tcompactor
u-veles-a Aug 6, 2026
cbd70e5
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Aug 6, 2026
1be7a5a
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Aug 6, 2026
9ccfe2d
fix after merge
u-veles-a Aug 6, 2026
0d268f7
add expirationpolicy
u-veles-a Aug 6, 2026
80c4e1e
add: manager excludes downsampling blocks from the query
u-veles-a Aug 6, 2026
ff3048c
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Aug 6, 2026
751ea1a
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Aug 6, 2026
d518fa4
fix linter
u-veles-a Aug 6, 2026
4774256
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Aug 6, 2026
b37f909
small fix
u-veles-a Aug 6, 2026
c88bfa5
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Aug 6, 2026
5963124
refactor: simplify block deletion logic and improve error handling
u-veles-a Aug 6, 2026
ead2e90
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Aug 6, 2026
74f2d9e
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Aug 7, 2026
ba4b8f9
fix review
u-veles-a Aug 10, 2026
07c3213
Merge branch 'pp' into tcompactor
u-veles-a Aug 10, 2026
5ba3382
revert fix memory
u-veles-a Aug 10, 2026
8f665f2
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Aug 10, 2026
423700b
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Aug 10, 2026
71e691e
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Aug 10, 2026
b7aea7f
feat(PROMPP_FEATURES): refactor related logic
u-veles-a Aug 11, 2026
9358c58
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Aug 11, 2026
fc3fb5b
fix after merge
u-veles-a Aug 11, 2026
e0d777d
fix: standardize downsampling variable naming to camel case
u-veles-a Aug 11, 2026
ddcc972
feat(upsampler): add upsampler package with NeedsUpsampling function …
u-veles-a Aug 11, 2026
0bfc9c1
Implement upsampler package with Querier, Series, and SeriesSet for s…
u-veles-a Aug 12, 2026
617fcde
fix review
u-veles-a Aug 12, 2026
609b673
Merge branch 'tcompactor' into downsampling_tcompactor
u-veles-a Aug 12, 2026
ef5f41d
feat(upsampler): add benchmark tests for iterator and querier perform…
u-veles-a Aug 12, 2026
ec95b11
fix: standardize rangeMs to rangeMS in Series and SeriesSet
u-veles-a Aug 12, 2026
7c0133d
fix(benchmarks): optimize label creation in iterator and series bench…
u-veles-a Aug 12, 2026
2c94a92
chore: remove .claudecodeignore file
u-veles-a Aug 12, 2026
7b2bb04
feat(manager): integrate upsampler for downsampling blocks in querier
u-veles-a Aug 12, 2026
b7602af
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Aug 12, 2026
10799ca
refactor(manager): simplify block skipping logic in Querier and Chunk…
u-veles-a Aug 12, 2026
0b669b3
feat(upsampler): implement downsampling logic in adapter and querier,…
u-veles-a Aug 12, 2026
731be59
feat(benchmarks): add benchmarks for head querier with and without do…
u-veles-a Aug 12, 2026
e952e62
feat(querier): bypass downsampling for specific functions and update …
u-veles-a Aug 12, 2026
7e29fcb
fix(Iterator): fixed holes in the graph
u-veles-a Aug 18, 2026
9808ca3
fix description
u-veles-a Aug 18, 2026
3b6d51e
fix(upsampler): interpolate only gaps within rangeMS*2
u-veles-a Aug 19, 2026
4bdb2ba
feat(fanout): fork fanout storage and centralize upsampling by max qu…
u-veles-a Aug 19, 2026
4fe79e6
fix(upsampler): interpolate only gaps within resolutionMS*2
u-veles-a Aug 19, 2026
f3b8cab
test(fanout): cover resolution selection and cross-storage interpolation
u-veles-a Aug 19, 2026
33e726a
feat(manager): report loaded downsampled blocks in a separate metric
u-veles-a Aug 19, 2026
ca3f48c
clearing
u-veles-a Aug 19, 2026
8cd4df3
clearing
u-veles-a Aug 19, 2026
4c3bbcc
clearing
u-veles-a Aug 19, 2026
54cce3c
clearing
u-veles-a Aug 19, 2026
7758231
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Aug 19, 2026
5e9fecf
fix(upsampler): hold a value drop flat for counter functions only
u-veles-a Aug 19, 2026
f703050
perf(upsampler): keep the gap thresholds as uint32 durations
u-veles-a Aug 19, 2026
fc74c2e
fix linter
u-veles-a Aug 19, 2026
9df11bd
upsampler: move gapThresholds next to its caller in series_set.go
u-veles-a Aug 19, 2026
ec758a1
fix(downsampling): base lookback-delta and needDownsampling decision…
u-veles-a Sep 2, 2026
c3174a4
adds(upsampler): add to allowedFuncs _over_time funcs
u-veles-a Sep 2, 2026
b4ca4d3
adds(upsampler): test and description
u-veles-a Sep 2, 2026
d4fb9f7
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Sep 2, 2026
ea083f2
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Sep 2, 2026
4ab88b5
feat(upsampler): hold the last value on gaps for _over_time funcs
u-veles-a Sep 3, 2026
bc5acdf
fix(linter): iterator.go
u-veles-a Sep 3, 2026
5c979c5
Merge branch 'agg_series_set' into downsampling_tcompactor
u-veles-a Sep 4, 2026
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
47 changes: 39 additions & 8 deletions cmd/prometheus/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -60,10 +60,12 @@ import (

"github.com/prometheus/prometheus/pp-pkg/blocks/block"
"github.com/prometheus/prometheus/pp-pkg/blocks/expirationpolicy"
"github.com/prometheus/prometheus/pp-pkg/blocks/fanout"
"github.com/prometheus/prometheus/pp-pkg/blocks/lcompactor"
"github.com/prometheus/prometheus/pp-pkg/blocks/manager"
"github.com/prometheus/prometheus/pp-pkg/blocks/tcompactor" // PP_CHANGES.md: rebuild on cpp
"github.com/prometheus/prometheus/pp-pkg/blocks/manager" // PP_CHANGES.md: rebuild on cpp
"github.com/prometheus/prometheus/pp-pkg/blocks/tcompactor"
"github.com/prometheus/prometheus/pp-pkg/featuresflags"
v1 "github.com/prometheus/prometheus/web/api/v1"

// PP_CHANGES.md: rebuild on cpp
"github.com/prometheus/prometheus/pp-pkg/localstorageobserver"
Expand All @@ -75,10 +77,8 @@ import (
pp_pkg_remote "github.com/prometheus/prometheus/pp-pkg/storage/remote" // PP_CHANGES.md: rebuild on cpp
pp_pkg_tsdb "github.com/prometheus/prometheus/pp-pkg/tsdb" // PP_CHANGES.md: rebuild on cpp

pp_storage "github.com/prometheus/prometheus/pp/go/storage" // PP_CHANGES.md: rebuild on cpp
// PP_CHANGES.md: rebuild on cpp
"github.com/prometheus/prometheus/pp/go/storage/catalog" // PP_CHANGES.md: rebuild on cpp
// PP_CHANGES.md: rebuild on cpp
pp_storage "github.com/prometheus/prometheus/pp/go/storage" // PP_CHANGES.md: rebuild on cpp
"github.com/prometheus/prometheus/pp/go/storage/catalog" // PP_CHANGES.md: rebuild on cpp
"github.com/prometheus/prometheus/pp/go/storage/querier" // PP_CHANGES.md: rebuild on cpp
"github.com/prometheus/prometheus/pp/go/storage/ready" // PP_CHANGES.md: rebuild on cpp
"github.com/prometheus/prometheus/pp/go/storage/remotewriter" // PP_CHANGES.md: rebuild on cpp
Expand Down Expand Up @@ -183,6 +183,8 @@ type flagConfig struct {
HeadRetentionTimeout model.Duration
UseBlockManagerStorage bool

Downsampling model.Duration

featureList []string
memlimitRatio float64
// These options are extracted from featureList
Expand Down Expand Up @@ -293,6 +295,11 @@ func (c *flagConfig) DisableBlockManagerStorage() {
c.UseBlockManagerStorage = false
}

// SetDownsampling sets the downsampling duration for the flagConfig.
func (c *flagConfig) SetDownsampling(downsampling model.Duration) {
c.Downsampling = downsampling
}

func main() {
if os.Getenv("DEBUG") != "" {
runtime.SetBlockProfileRate(20)
Expand Down Expand Up @@ -657,6 +664,17 @@ func main() {
cfg.tsdb.OutOfOrderTimeWindow = cfgFile.StorageConfig.TSDBConfig.OutOfOrderTimeWindow
}

// check if the downsampling value is greater than the scrape interval
if cfg.Downsampling != 0 && cfg.Downsampling < cfgFile.GlobalConfig.ScrapeInterval {
_ = level.Error(logger).Log(
"msg", "The downsampling value must be greater than the scrape interval",
"scrape_interval", cfgFile.GlobalConfig.ScrapeInterval,
"downsampling", cfg.Downsampling,
)

os.Exit(2)
}

// Now that the validity of the config is established, set the config
// success metrics accordingly, although the config isn't really loaded
// yet. This will happen later (including setting these metrics again),
Expand Down Expand Up @@ -803,6 +821,7 @@ func main() {
&pp_storage.Options{
Seed: cfgFile.GlobalConfig.ExternalLabels.Hash(),
BlockDuration: time.Duration(cfg.tsdb.MinBlockDuration),
Downsampling: time.Duration(cfg.Downsampling),
CommitInterval: time.Duration(cfg.WalCommitInterval),
MaxRetentionPeriod: ppRetentionPeriod,
HeadRetentionPeriod: time.Duration(cfg.HeadRetentionTimeout),
Expand Down Expand Up @@ -832,10 +851,16 @@ func main() {
prometheus.DefaultRegisterer,
)

retentionMS := int64(time.Duration(cfg.tsdb.RetentionDuration) / time.Millisecond)
downsamplingMS := int64(time.Duration(cfg.Downsampling) / time.Millisecond)
adapter := pp_pkg_storage.NewAdapter(
clock,
hManager.Proxy(),
hManager.Builder(),
&pp_pkg_storage.AdapterOptions{
RetentionMS: retentionMS,
DownsamplingMS: downsamplingMS,
},
hManager.MergeOutOfOrderChunks,
prometheus.DefaultRegisterer,
)
Expand All @@ -861,6 +886,7 @@ func main() {
tsdbHistorical *tsdbHistoricalStorage
persistedStorage storage.Storage = localStorage
startTimeFn = localStorage.StartTime
minTimeBlocks v1.MinTimeBlocks
)
if !agentMode {
// Storage is constructed eagerly here for both schemes. The historical
Expand Down Expand Up @@ -892,7 +918,6 @@ func main() {
"CorruptedRetentionDuration", cfg.tsdb.CorruptedRetentionDuration,
"EnableOverlappingCompaction", cfg.tsdb.EnableOverlappingCompaction,
)
retentionMS := int64(time.Duration(cfg.tsdb.RetentionDuration) / time.Millisecond)

chunkPool := chunkenc.NewPool()
compactCtx, compactCancel := context.WithCancel(context.Background())
Expand Down Expand Up @@ -928,6 +953,7 @@ func main() {
localStoragePath,
&manager.Options{
RetentionDuration: retentionMS,
DownsamplingMS: downsamplingMS,
CorruptedRetentionDuration: time.Duration(cfg.tsdb.CorruptedRetentionDuration),
EnableOverlappingCompaction: cfg.tsdb.EnableOverlappingCompaction,
},
Expand All @@ -948,6 +974,8 @@ func main() {
os.Exit(1)
}

minTimeBlocks = blockManager.MinTimeBlocks

bs := &blockStorage{m: blockManager, onClose: func() error {
// Cancel any in-flight leveled compaction first so the manager
// loop can return promptly, then stop the loop and close blocks.
Expand Down Expand Up @@ -984,7 +1012,7 @@ func main() {
log.With(logger, "component", "remote"),
startTimeFn,
)
fanoutStorage := storage.NewFanout(
fanoutStorage := fanout.New(
logger,
adapter,
persistedStorage,
Expand Down Expand Up @@ -1174,6 +1202,9 @@ func main() {
cfg.web.RuleManager = ruleManager
cfg.web.Notifier = notifierManager
cfg.web.LookbackDelta = time.Duration(cfg.lookbackDelta)
// x2 because the needed minimum lookback delta is 2*downsampling
cfg.web.DownsamplingLookbackDelta = time.Duration(cfg.Downsampling * 2)
cfg.web.MinTimeBlocks = minTimeBlocks
cfg.web.IsAgent = agentMode
cfg.web.AppName = modeAppName

Expand Down
137 changes: 137 additions & 0 deletions cmd/prompptool/blocks.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,137 @@
package main

import (
"context"
"encoding/json"
"errors"
"fmt"
"os"
"slices"
"strconv"
"text/tabwriter"
"time"

"github.com/alecthomas/kingpin/v2"
"github.com/alecthomas/units"
"github.com/go-kit/log"
"github.com/prometheus/client_golang/prometheus"

"github.com/prometheus/prometheus/pp-pkg/blocks/block"
)

// cmdBlocks is the command for listing blocks.
type cmdBlocks struct {
humanReadable bool
}

// registerCmdBlocks registers the blocks command.
func registerCmdBlocks(cmd *cmdBlocks, clause *kingpin.CmdClause) {
clause.Flag(
"human-readable",
"Print human readable values. Default is false.",
).Default("false").Short('r').BoolVar(&cmd.humanReadable)
}

// Do lists the blocks in the directory.
func (cmd *cmdBlocks) Do(
ctx context.Context,
workingDir string,
logger log.Logger,
registerer prometheus.Registerer,
) error {
return listBlocks(logger, workingDir, cmd.humanReadable)
}

// listBlocks lists the blocks in the directory.
func listBlocks(logger log.Logger, path string, humanReadable bool) error {
blocks, _, err := block.OpenBlocks(logger, path, nil, nil)
if err != nil {
return err
}
defer func() {
err = errors.Join(err, block.CloseAll(blocks))
}()

slices.SortFunc(blocks, func(a, b *block.Block) int {
switch {
case a.Meta().MinTime < b.Meta().MinTime:
return -1
case a.Meta().MinTime > b.Meta().MinTime:
return 1
default:
return 0
}
})

printBlocks(blocks, true, humanReadable)

return nil
}

// printBlocks prints the blocks to the console.
func printBlocks(blocks []*block.Block, writeHeader, humanReadable bool) {
tw := tabwriter.NewWriter(os.Stdout, 13, 0, 2, ' ', 0)
defer tw.Flush()

if writeHeader {
fmt.Fprintln(
tw,
"BLOCK ULID\tMIN TIME\tMAX TIME\tDURATION\tNUM SAMPLES\tNUM CHUNKS\tNUM SERIES\tSIZE\tRESOLUTION\tLABELS",
)
}

for _, b := range blocks {
meta := b.Metadata()

fmt.Fprintf(tw,
"%v\t%v\t%v\t%v\t%v\t%v\t%v\t%v\t%v\t%v\n",
meta.ULID,
getFormatedTime(meta.MinTime, humanReadable),
getFormatedTime(meta.MaxTime, humanReadable),
time.Duration(meta.MaxTime-meta.MinTime)*time.Millisecond,
meta.Stats.NumSamples,
meta.Stats.NumChunks,
meta.Stats.NumSeries,
getFormatedBytes(b.Size(), humanReadable),
getFormatedDuration(meta.Thanos.Downsample.Resolution, humanReadable),
labelsToString(meta.Thanos.Labels),
)
}
}

// getFormatedTime converts a timestamp to a human readable string.
func getFormatedTime(timestamp int64, humanReadable bool) string {
if humanReadable {
return time.Unix(timestamp/1000, 0).UTC().String()
}

return strconv.FormatInt(timestamp, 10)
}

// getFormatedBytes converts a number of bytes to a human readable string.
func getFormatedBytes(bytes int64, humanReadable bool) string {
if humanReadable {
return units.Base2Bytes(bytes).String()
}

return strconv.FormatInt(bytes, 10)
}

// getFormatedDuration converts a duration to a human readable string.
func getFormatedDuration(duration int64, humanReadable bool) string {
if humanReadable {
return (time.Duration(duration) * time.Millisecond).String()
}

return strconv.FormatInt(duration, 10)
}

// labelsToString converts a map[string]string to a JSON string.
func labelsToString(ls map[string]string) string {
data, err := json.Marshal(ls)
if err != nil {
return err.Error()
}

return string(data)
}
13 changes: 13 additions & 0 deletions cmd/prompptool/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (
"os/signal"
"path/filepath"

"github.com/prometheus/common/model"
"github.com/prometheus/common/promlog"
"github.com/prometheus/common/version"

Expand Down Expand Up @@ -47,6 +48,10 @@ func main() {
var persistHeadCmd cmdPersistHead
registerCmdPersistHead(&persistHeadCmd, persistHeadClause)

blocksClause := app.Command("blocks", "Converting prometheus wal to tsdb-blocks.")
var blocksCmd cmdBlocks
registerCmdBlocks(&blocksCmd, blocksClause)

cmd := kingpin.MustParse(app.Parse(os.Args[1:]))
logger := initLogger(*verbose)
logger = log.With(logger, "cmd", cmd)
Expand All @@ -70,6 +75,11 @@ func main() {
level.Error(logger).Log("msg", "fail to persist head", "error", err)
os.Exit(1)
}
case blocksClause.FullCommand():
if err := blocksCmd.Do(ctx, workingDir, logger, nil); err != nil {
level.Error(logger).Log("msg", "fail to list blocks", "error", err)
os.Exit(1)
}
}
}

Expand Down Expand Up @@ -103,3 +113,6 @@ type noopFlagConfig struct{}

// DisableBlockManagerStorage is a no-op implementation of the FlagConfig interface, used when no feature flags are set.
func (noopFlagConfig) DisableBlockManagerStorage() {}

// SetDownsampling sets the downsampling duration for the flagConfig.
func (noopFlagConfig) SetDownsampling(model.Duration) {}
1 change: 1 addition & 0 deletions cmd/prompptool/persisthead.go
Original file line number Diff line number Diff line change
Expand Up @@ -130,6 +130,7 @@ func (cmd *cmdPersistHead) Do(
bw := block.NewWriter[*shard.Shard](
outputDir,
block.DefaultChunkSegmentSize,
cppbridge.NoDownsampling,
time.Duration(cmd.blockDuration),
0, // no retention filtering: persist all heads regardless of age
clockwork.NewRealClock(),
Expand Down
1 change: 1 addition & 0 deletions cmd/prompptool/walpp.go
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,7 @@ func (cmd *cmdWALPPToBlock) Do(
bw := block.NewWriter[*shard.Shard](
workingDir,
block.DefaultChunkSegmentSize,
cppbridge.NoDownsampling,
time.Duration(cmd.blockDuration),
0, // no retention filtering: persist all heads regardless of age
clock,
Expand Down
25 changes: 17 additions & 8 deletions pp-pkg/blocks/block/block.go
Original file line number Diff line number Diff line change
Expand Up @@ -270,6 +270,13 @@ func (pb *Block) Metadata() metadata.Meta {
return pb.meta
}

// MinTime returns the minimum time of the block.
func (pb *Block) MinTime() int64 {
pb.mtxMeta.RLock()
defer pb.mtxMeta.RUnlock()
return pb.meta.MinTime
}

// OverlapsClosedInterval returns true if the block overlaps [mint, maxt].
func (pb *Block) OverlapsClosedInterval(mint, maxt int64) bool {
// The block itself is a half-open interval
Expand Down Expand Up @@ -558,14 +565,16 @@ func (o Overlaps) String() string {
))
}

res = append(res, fmt.Sprintf(
"[key: %s, mint: %d, maxt: %d, range: %s, blocks: %d]: %s",
r.Key,
r.Min, r.Max,
(time.Duration((r.Max-r.Min)/1000)*time.Second).String(),
len(overlaps),
strings.Join(groups, ", "),
))
res = append(
res, fmt.Sprintf(
"[key: %s, mint: %d, maxt: %d, range: %s, blocks: %d]: %s",
r.Key,
r.Min, r.Max,
(time.Duration((r.Max-r.Min)/1000)*time.Second).String(),
len(overlaps),
strings.Join(groups, ", "),
),
)
}

return strings.Join(res, "\n")
Expand Down
2 changes: 1 addition & 1 deletion pp-pkg/blocks/block/functions.go
Original file line number Diff line number Diff line change
Expand Up @@ -24,7 +24,7 @@ const checkContextEveryNIterations = 100
//

// CloseAll closes all given closers.
func CloseAll(cs []io.Closer) error {
func CloseAll[TCloser io.Closer](cs []TCloser) error {
errs := make([]error, 0, len(cs))
for _, c := range cs {
errs = append(errs, c.Close())
Expand Down
Loading
Loading