diff --git a/gorilla-merger/cmd/gorilla-merger/main.go b/gorilla-merger/cmd/gorilla-merger/main.go index cb3aa010f..24cd9e470 100644 --- a/gorilla-merger/cmd/gorilla-merger/main.go +++ b/gorilla-merger/cmd/gorilla-merger/main.go @@ -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)") } diff --git a/gorilla-merger/internal/merger/coldpartstore_test.go b/gorilla-merger/internal/merger/coldpartstore_test.go index c35e2e4ff..833d19deb 100644 --- a/gorilla-merger/internal/merger/coldpartstore_test.go +++ b/gorilla-merger/internal/merger/coldpartstore_test.go @@ -3,6 +3,7 @@ package merger import ( "bytes" "context" + "math" "net/http" "net/http/httptest" "testing" @@ -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{ @@ -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 { @@ -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 diff --git a/gorilla-merger/internal/merger/customstore.go b/gorilla-merger/internal/merger/customstore.go index 315a646d3..b94ede825 100644 --- a/gorilla-merger/internal/merger/customstore.go +++ b/gorilla-merger/internal/merger/customstore.go @@ -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 } @@ -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 }