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: 18 additions & 4 deletions gorilla-merger/cmd/gorilla-merger/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -118,10 +118,24 @@ func run(cfg config, logger *slog.Logger, kitLogger kitslog.Logger) error {
coldBucket = bkt
defer func() { _ = coldBucket.Close() }()
coldStore = merger.NewColdPartStore(bkt, kitLogger)
// Rediscover any parts already in the bucket (header/index only).
if rerr := coldStore.Reload(context.Background()); rerr != nil {
logger.Warn("cold manifest reload failed (continuing empty)", "err", rerr)
}
// Rediscover any parts already in the bucket (header/index only) in the
// BACKGROUND. Reload fetches + OpenParts every stored part, so with a large
// accumulated cold tier (thousands of parts) it takes tens of seconds and
// is memory-heavy. Running it synchronously here BLOCKS the StoreAPI and
// HTTP frontend from starting (they launch below), so for the whole reload
// window after a restart thanos-query sees the gRPC endpoint as down and a
// cold (or warm) query returns empty with no streamColdSeries activity —
// the served-empty symptom. The manifest is mutex-guarded and queried under
// RLock, so a concurrent reload is safe: cold queries simply see a smaller
// (growing) manifest until it completes, then the full set. New parts POSTed
// during the reload still register live via Put; Reload's atomic manifest
// swap may briefly drop a part POSTed mid-reload, but the next restart's
// reload (or a re-POST) re-registers it, and the warm/open path is unaffected.
go func() {
if rerr := coldStore.Reload(context.Background()); rerr != nil {
logger.Warn("cold manifest reload failed (continuing empty)", "err", rerr)
}
}()
} else {
logger.Warn("no objstore config provided; cold-part store disabled (decode-on-read unavailable)")
}
Expand Down
71 changes: 56 additions & 15 deletions gorilla-merger/internal/merger/coldpartstore_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -3,6 +3,7 @@ package merger
import (
"bytes"
"context"
"math"
"net/http"
"net/http/httptest"
"testing"
Expand Down Expand Up @@ -545,15 +546,11 @@ func TestCustomStoreTimeRangeIncludesCold(t *testing.T) {
}
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)
// The fresh tsdb head holds no samples, so StartTime() reports the empty-head
// sentinel math.MaxInt64 (see TestCustomStoreTimeRangeEmptyHead). A cold part
// older than "now" stands in for the cold-only window thanos-query would
// otherwise prune.
coldStart := time.Now().UnixMilli() - int64(24*time.Hour/time.Millisecond)
bkt := objstore.NewInMemBucket()
coldStore := NewColdPartStore(bkt, nil)
mustPut(t, coldStore, writePartBytes(t, coldStart, coldStart+1000, []coldpart.Series{
Expand All @@ -562,12 +559,8 @@ func TestCustomStoreTimeRangeIncludesCold(t *testing.T) {

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.
// With the cold querier attached the advertised min drops to the cold floor
// (NOT the empty-head MaxInt64 sentinel, which would prune the merger).
cs.setColdQuerier(NewColdQuerier(coldStore))
min, max := cs.timeRange()
if min != coldStart {
Expand All @@ -578,6 +571,54 @@ func TestCustomStoreTimeRangeIncludesCold(t *testing.T) {
}
}

// TestCustomStoreTimeRangeEmptyHead is the regression proof for the
// served-empty-after-restart bug: a fresh/empty tsdb head reports StartTime ==
// math.MaxInt64, and if that leaks into the advertised StoreAPI MinTime,
// thanos-query prunes the merger from EVERY query (no window can be >=
// MaxInt64) — even the cold parts already reloaded from S3 become unreachable,
// served empty with no streamColdSeries activity. The advertised MinTime must
// therefore be:
// - the cold floor when cold parts exist (empty head must not win), and
// - math.MinInt64 (never MaxInt64) when the store is genuinely empty, so the
// merger stays discoverable and is not pruned while it waits for data.
func TestCustomStoreTimeRangeEmptyHead(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() })

// Sanity: a fresh head really does report the MaxInt64 sentinel.
if got, err := st.DB.StartTime(); err != nil {
t.Fatalf("StartTime: %v", err)
} else if got != math.MaxInt64 {
t.Logf("note: empty-head StartTime = %d (expected MaxInt64); test still asserts no MaxInt64 leak", got)
}

// (1) Empty head, NO cold querier: must advertise MinInt64, not MaxInt64.
csBare := newCustomStore(st.DB, labels.EmptyLabels(), nil)
if min, _ := csBare.timeRange(); min == math.MaxInt64 {
t.Fatalf("empty head (no cold) advertised MinTime = MaxInt64; merger would be pruned from every query")
} else if min != math.MinInt64 {
t.Fatalf("empty head (no cold) timeRange min = %d, want MinInt64", min)
}

// (2) Empty head WITH cold parts already in the store (the post-restart
// reload case): MinTime must drop to the cold floor, not the empty-head
// MaxInt64 sentinel.
coldStart := time.Now().UnixMilli() - int64(48*time.Hour/time.Millisecond)
coldStore := NewColdPartStore(objstore.NewInMemBucket(), nil)
mustPut(t, coldStore, writePartBytes(t, coldStart, coldStart+1000, []coldpart.Series{
coldSeries(labels.FromStrings(labels.MetricName, "reloaded_cold"), coldStart, 1, 2),
}))
csCold := newCustomStore(st.DB, labels.EmptyLabels(), nil)
csCold.setColdQuerier(NewColdQuerier(coldStore))
if min, _ := csCold.timeRange(); min != coldStart {
t.Fatalf("empty head + cold parts: timeRange min = %d, want cold floor %d (MaxInt64 sentinel must not win)", min, coldStart)
}
}

// 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
Expand Down
38 changes: 30 additions & 8 deletions gorilla-merger/internal/merger/customstore.go
Original file line number Diff line number Diff line change
Expand Up @@ -76,15 +76,30 @@ 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.
// The advertised MinTime is the MINIMUM of two lower bounds, whichever exist:
// the tsdb head StartTime and (when a cold querier is attached and tracks parts)
// the oldest cold part's block start. 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.
//
// Two failure modes this guards against, both of which leave stored cold parts
// served EMPTY with no streamColdSeries activity:
//
// - Cold data OLDER than the head: a query for that old (cold-only) window is
// pruned unless the cold floor lowers MinTime to cover it.
//
// - An EMPTY tsdb head: tsdb.DB.StartTime() returns math.MaxInt64 for a head
// holding no samples (e.g. right after a restart, before the first warm
// fragment lands, while cold parts already exist in S3 and were reloaded).
// Advertising MaxInt64 prunes the merger from EVERY query — no window can be
// >= MaxInt64 — so even the cold parts are unreachable. We must NOT let the
// empty-head sentinel win over the cold floor (or survive as the advertised
// MinTime at all), hence taking the min of the two bounds rather than only
// lowering when cold < head.
func (s *customStore) timeRange() (int64, int64) {
var minTime int64 = math.MinInt64
// math.MaxInt64 means "no lower bound from this source": an empty tsdb head
// reports StartTime == MaxInt64, which must NOT become the advertised floor.
minTime := int64(math.MaxInt64)
if st, err := s.db.StartTime(); err == nil {
minTime = st
}
Expand All @@ -93,6 +108,13 @@ func (s *customStore) timeRange() (int64, int64) {
minTime = coldMin
}
}
// If neither source supplied a real lower bound (empty head, no cold parts),
// the store currently holds nothing; advertise MinInt64 so thanos-query does
// not prune it (it simply returns no series for any window, which is correct
// and cheap) and so it stays discoverable for when data arrives.
if minTime == math.MaxInt64 {
minTime = math.MinInt64
}
return minTime, math.MaxInt64
}

Expand Down