Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
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
22 changes: 22 additions & 0 deletions gorilla-merger/internal/merger/coldpartstore.go
Original file line number Diff line number Diff line change
Expand Up @@ -201,6 +201,28 @@ func (s *ColdPartStore) NumParts() int {
return len(s.man.entries)
}

// MinBlockStart returns the smallest block_start_ms across all tracked parts
// and true, or (0,false) when the manifest is empty. The StoreAPI uses it to
// lower its advertised MinTime to the oldest cold data, so thanos-query routes
// queries for old (cold-only) windows to the merger instead of pruning it. The
// per-series [min_ts,max_ts] index entries (not the block range) remain the
// authoritative time filter inside the query path; this is only the advertised
// floor of what the merger MIGHT serve.
func (s *ColdPartStore) MinBlockStart() (int64, bool) {
s.mu.RLock()
defer s.mu.RUnlock()
if len(s.man.entries) == 0 {
return 0, false
}
min := s.man.entries[0].BlockStartMs
for _, e := range s.man.entries[1:] {
if e.BlockStartMs < min {
min = e.BlockStartMs
}
}
return min, true
}

// Bucket exposes the underlying bucket (used by the query path and tests).
func (s *ColdPartStore) Bucket() objstore.Bucket { return s.bkt }

Expand Down
140 changes: 140 additions & 0 deletions gorilla-merger/internal/merger/coldpartstore_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ import (
"net/http"
"net/http/httptest"
"testing"
"time"

"github.com/ProjectASAP/asap-gorilla-go/coldpart"
"github.com/prometheus/prometheus/model/labels"
Expand Down Expand Up @@ -495,3 +496,142 @@ func TestStoreAPIUnionsColdAndTSDB(t *testing.T) {
t.Fatalf("union missing a tier: sawHot=%v sawCold=%v", sawHot, sawCold)
}
}

// TestColdPartStoreMinBlockStart asserts MinBlockStart reports the earliest
// block start across all parts (and false for an empty store / nil querier).
func TestColdPartStoreMinBlockStart(t *testing.T) {
store := NewColdPartStore(objstore.NewInMemBucket(), nil)
if _, ok := store.MinBlockStart(); ok {
t.Fatalf("empty store MinBlockStart ok = true, want false")
}

// Insert parts out of block-start order; MinBlockStart must find the min.
mustPut(t, store, writePartBytes(t, 5000, 6000, []coldpart.Series{
coldSeries(labels.FromStrings(labels.MetricName, "a"), 5000, 1, 2),
}))
mustPut(t, store, writePartBytes(t, 1000, 2000, []coldpart.Series{
coldSeries(labels.FromStrings(labels.MetricName, "b"), 1000, 3, 4),
}))
mustPut(t, store, writePartBytes(t, 9000, 10000, []coldpart.Series{
coldSeries(labels.FromStrings(labels.MetricName, "c"), 9000, 5, 6),
}))
got, ok := store.MinBlockStart()
if !ok || got != 1000 {
t.Fatalf("MinBlockStart = (%d,%v), want (1000,true)", got, ok)
}

// Through the querier wrapper, and the nil-querier guard.
if qmin, ok := NewColdQuerier(store).MinBlockStart(); !ok || qmin != 1000 {
t.Fatalf("ColdQuerier.MinBlockStart = (%d,%v), want (1000,true)", qmin, ok)
}
var nilQ *ColdQuerier
if _, ok := nilQ.MinBlockStart(); ok {
t.Fatalf("nil querier MinBlockStart ok = true, want false")
}
}

// TestCustomStoreTimeRangeIncludesCold is the regression proof for the
// served-empty bug: with a cold querier attached, the customStore's advertised
// timeRange().min must drop to the oldest cold part's block start (which is far
// older than the tsdb head's StartTime). thanos-query uses this advertised
// MinTime to decide whether to route a query to the merger; if it stays at the
// recent tsdb StartTime, queries for old (cold-only) windows are pruned and
// streamColdSeries never runs — the parts are stored but served empty.
func TestCustomStoreTimeRangeIncludesCold(t *testing.T) {
dir := t.TempDir()
st, err := OpenStorage(StorageOptions{Dir: dir})
if err != nil {
t.Fatalf("open storage: %v", err)
}
t.Cleanup(func() { _ = st.Close() })

// The tsdb head's StartTime is roughly "now" (no old data ingested).
tsdbStart, err := st.DB.StartTime()
if err != nil {
t.Fatalf("StartTime: %v", err)
}

// A cold part whose block starts well BEFORE the tsdb StartTime — exactly
// the cold-only window thanos-query would otherwise prune.
coldStart := tsdbStart - int64(24*time.Hour/time.Millisecond)
bkt := objstore.NewInMemBucket()
coldStore := NewColdPartStore(bkt, nil)
mustPut(t, coldStore, writePartBytes(t, coldStart, coldStart+1000, []coldpart.Series{
coldSeries(labels.FromStrings(labels.MetricName, "old_metric"), coldStart, 1, 2),
}))

cs := newCustomStore(st.DB, labels.EmptyLabels(), nil)

// Without the cold querier the advertised min is the (recent) tsdb StartTime.
if min, _ := cs.timeRange(); min != tsdbStart {
t.Fatalf("tsdb-only timeRange min = %d, want tsdb StartTime %d", min, tsdbStart)
}

// With the cold querier attached the advertised min drops to the cold floor.
cs.setColdQuerier(NewColdQuerier(coldStore))
min, max := cs.timeRange()
if min != coldStart {
t.Fatalf("with-cold timeRange min = %d, want cold block start %d", min, coldStart)
}
if max <= min {
t.Fatalf("timeRange max = %d not > min = %d", max, min)
}
}

// TestCustomStoreSeriesServesOldColdWindow drives the full Series RPC over a
// window that covers ONLY the cold part (older than the tsdb head) and asserts
// the cold series + its samples are returned. This is the end-to-end read-path
// proof that decode-on-read serves stored parts for an old-only window.
func TestCustomStoreSeriesServesOldColdWindow(t *testing.T) {
dir := t.TempDir()
st, err := OpenStorage(StorageOptions{Dir: dir})
if err != nil {
t.Fatalf("open storage: %v", err)
}
t.Cleanup(func() { _ = st.Close() })

const base = int64(1_600_000_000_000) // well in the past
bkt := objstore.NewInMemBucket()
coldStore := NewColdPartStore(bkt, nil)
mustPut(t, coldStore, writePartBytes(t, base, base+1000, []coldpart.Series{
coldSeries(labels.FromStrings(labels.MetricName, "http_requests_total", "job", "api"), base, 11, 22),
}))

cs := newCustomStore(st.DB, labels.EmptyLabels(), nil)
cs.setColdQuerier(NewColdQuerier(coldStore))

// The advertised min must cover this old window (else thanos-query prunes us).
if min, _ := cs.timeRange(); min > base {
t.Fatalf("advertised min %d > query window start %d: store would be pruned", min, base)
}

req := &storepb.SeriesRequest{
MinTime: base - 60_000,
MaxTime: base + 60_000,
Matchers: []storepb.LabelMatcher{
{Type: storepb.LabelMatcher_EQ, Name: labels.MetricName, Value: "http_requests_total"},
},
}
fss := &fakeSeriesServer{ctx: context.Background()}
if err := cs.Series(req, fss); err != nil {
t.Fatalf("Series: %v", err)
}

var results int
for _, r := range fss.responses {
series := r.GetSeries()
if series == nil {
continue
}
results++
want := labels.FromStrings(labels.MetricName, "http_requests_total", "job", "api")
if labels.Compare(labelsOf(t, series.Labels), want) != 0 {
t.Fatalf("cold series labels:\n got %s\n want %s", labelsOf(t, series.Labels).String(), want.String())
}
assertSamples(t, "old-cold", samplesFromChunks(t, series.Chunks),
[]sample{{base, 11}, {base + 1000, 22}})
}
if results != 1 {
t.Fatalf("old-cold window: got %d series, want 1", results)
}
}
12 changes: 12 additions & 0 deletions gorilla-merger/internal/merger/coldquery.go
Original file line number Diff line number Diff line change
Expand Up @@ -43,6 +43,18 @@ func NewColdQuerier(store *ColdPartStore) *ColdQuerier {
return &ColdQuerier{store: store}
}

// MinBlockStart returns the smallest block_start_ms across the cold store's
// tracked parts and true, or (0,false) when there is no cold data (or no
// store). The customStore folds this into its advertised StoreAPI MinTime so
// thanos-query does not prune the merger from a query whose window only covers
// old cold data.
func (q *ColdQuerier) MinBlockStart() (int64, bool) {
if q == nil || q.store == nil {
return 0, false
}
return q.store.MinBlockStart()
}

// Series returns every cold series matching all matchers with at least one
// sample in the inclusive window [mintMs,maxtMs], each as XOR chunk(s) over the
// CLIPPED in-window samples (coldpart.Part.Series returns the WHOLE matched
Expand Down
17 changes: 17 additions & 0 deletions gorilla-merger/internal/merger/customstore.go
Original file line number Diff line number Diff line change
Expand Up @@ -75,11 +75,24 @@ func (s *customStore) setColdQuerier(c *ColdQuerier) { s.cold = c }

// timeRange mirrors TSDBStore.TimeRange: min = head StartTime (the oldest
// sample currently held), max = +inf so the open window is always queried.
//
// When a cold querier is attached, the advertised min is LOWERED to the oldest
// cold part's block start when that is earlier than the tsdb StartTime. This is
// load-bearing: thanos-query prunes a store from a query's fan-out when the
// query window falls entirely below the store's advertised MinTime. The cold
// path serves raw samples OLDER than the tsdb head's StartTime, so without this
// floor a query for old (cold-only) data is never routed to the merger and
// streamColdSeries is never invoked — the parts are stored but served empty.
func (s *customStore) timeRange() (int64, int64) {
var minTime int64 = math.MinInt64
if st, err := s.db.StartTime(); err == nil {
minTime = st
}
if s.cold != nil {
if coldMin, ok := s.cold.MinBlockStart(); ok && coldMin < minTime {
minTime = coldMin
}
}
return minTime, math.MaxInt64
}

Expand Down Expand Up @@ -293,10 +306,14 @@ func (s *customStore) streamColdSeries(
if s.cold == nil {
return nil
}
level.Debug(s.logger).Log("msg", "streamColdSeries: querying cold parts",
"mint", r.MinTime, "maxt", r.MaxTime, "matchers", len(matchers), "skip_chunks", r.SkipChunks)
coldSeries, err := s.cold.Series(ctx, matchers, r.MinTime, r.MaxTime)
if err != nil {
return status.Error(codes.Internal, err.Error())
}
level.Debug(s.logger).Log("msg", "streamColdSeries: emitting cold series",
"cold_series", len(coldSeries), "mint", r.MinTime, "maxt", r.MaxTime)
for _, cs := range coldSeries {
full := completeLabels(cs.Labels, finalExt)
zls := zLabelsCopy(full)
Expand Down