Unified Storage: Mark tenant as orphaned when it doesnt exist in tenant api (#123534)

* mark tenant as orphaned when it doesnt exist in tenant api

* make gen-jsonnet

* when reconciling tenant only include span if over 1s duration so its not spammy when poller goes through 100k+ tenants
This commit is contained in:
owensmallwood
2026-04-24 14:22:43 -06:00
committed by GitHub
parent 6b51cb39cd
commit 47af657fc9
3 changed files with 143 additions and 10 deletions
+1
View File
@@ -103,6 +103,7 @@
"section-variables": (import '../dev-dashboards/section-variables/section-variables.json'),
"shared_queries": (import '../dev-dashboards/panel-common/shared_queries.json'),
"slow_queries_and_annotations": (import '../dev-dashboards/scenarios/slow_queries_and_annotations.json'),
"smoke": (import '../dev-dashboards/transforms/smoke.json'),
"status-history-thresholds-mappings": (import '../dev-dashboards/panel-status-history/status-history-thresholds-mappings.json'),
"table_footer": (import '../dev-dashboards/panel-table/table_footer.json'),
"table_kitchen_sink": (import '../dev-dashboards/panel-table/table_kitchen_sink.json'),
+109 -9
View File
@@ -12,6 +12,9 @@ import (
"time"
authnlib "github.com/grafana/authlib/authn"
"go.opentelemetry.io/otel/attribute"
"go.opentelemetry.io/otel/codes"
"go.opentelemetry.io/otel/trace"
apierrors "k8s.io/apimachinery/pkg/api/errors"
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
"k8s.io/apimachinery/pkg/apis/meta/v1/unstructured"
@@ -46,6 +49,11 @@ const (
defaultPollInterval = 1 * time.Hour
pollPageSize = 2000
// reconcileTraceThreshold controls whether a span is emitted for
// reconcileTenantPendingDelete. Reconciles below this duration are noise in
// traces (mostly just the LabelingComplete fast-path).
reconcileTraceThreshold = 1 * time.Second
)
// TenantWatcher watches Tenant CRDs and syncs pending-delete state to the KV
@@ -278,21 +286,29 @@ func (tw *TenantWatcher) startPolling(ctx context.Context) {
// reconcile each, then clear KV records for tenants that dropped out of the
// filtered view. Memory is bounded to one page plus the liveNames set.
func (tw *TenantWatcher) runPollCycle(ctx context.Context) {
ctx, span := tracer.Start(ctx, "resource.TenantWatcher.runPollCycle", trace.WithAttributes(
attribute.Int("page_size", pollPageSize),
))
defer span.End()
start := time.Now()
liveNames := make(map[string]struct{})
listCtx, listSpan := tracer.Start(ctx, "resource.TenantWatcher.runPollCycle.list")
var continueToken string
var pageCount int
for {
select {
case <-ctx.Done():
case <-listCtx.Done():
listSpan.End()
return
case <-tw.stopCh:
listSpan.End()
return
default:
}
page, err := tw.client.Resource(tenantGVR).List(ctx, metav1.ListOptions{
page, err := tw.client.Resource(tenantGVR).List(listCtx, metav1.ListOptions{
Limit: pollPageSize,
Continue: continueToken,
LabelSelector: labelPendingDelete + "=true",
@@ -300,45 +316,78 @@ func (tw *TenantWatcher) runPollCycle(ctx context.Context) {
if err != nil {
tw.log.Error("tenant watcher poll cycle: list failed, skipping clear phase",
"error", err, "pages_so_far", pageCount, "names_so_far", len(liveNames))
listSpan.RecordError(err)
listSpan.SetStatus(codes.Error, "list failed")
listSpan.End()
span.RecordError(err)
span.SetStatus(codes.Error, "list failed")
return
}
pageCount++
for i := range page.Items {
item := &page.Items[i]
liveNames[item.GetName()] = struct{}{}
tw.handleTenant(ctx, item)
tw.handleTenant(listCtx, item)
}
continueToken = page.GetContinue()
if continueToken == "" {
break
}
}
listDuration := time.Since(start)
listSpan.SetAttributes(
attribute.Int("pages", pageCount),
attribute.Int("live_tenants", len(liveNames)),
attribute.Int64("duration_ms", listDuration.Milliseconds()),
)
listSpan.End()
// An empty list should never happen. If it does something is very wrong, and we don't want to nuke all the
// pending delete records as a consequence.
if len(liveNames) == 0 {
tw.log.Warn("tenant watcher poll cycle: zero live tenants, skipping clear phase",
"pages", pageCount, "list_duration", listDuration)
span.AddEvent("clear phase skipped: zero live tenants")
span.SetAttributes(
attribute.Int("live_tenants", 0),
attribute.Int("pages", pageCount),
attribute.Int64("list_duration_ms", listDuration.Milliseconds()),
)
return
}
// Go through all pending delete records and reconcile against the pending-delete tenants from the List above
clearCtx, clearSpan := tracer.Start(ctx, "resource.TenantWatcher.runPollCycle.clear")
clearStart := time.Now()
var cleared, leftForDeleter, scanned, raced, orphanedSkipped int
for name, err := range tw.pendingDeleteStore.Names(ctx) {
var cleared, leftForDeleter, scanned, raced, orphanedSkipped, markedOrphaned int
for name, err := range tw.pendingDeleteStore.Names(clearCtx) {
if err != nil {
tw.log.Error("tenant watcher poll cycle: failed to list kv records", "error", err)
clearSpan.RecordError(err)
clearSpan.SetStatus(codes.Error, "failed to list kv records")
clearSpan.End()
span.RecordError(err)
span.SetStatus(codes.Error, "failed to list kv records")
return
}
select {
case <-clearCtx.Done():
clearSpan.AddEvent("clear phase cancelled", trace.WithAttributes(attribute.Int("scanned", scanned)))
clearSpan.End()
return
case <-tw.stopCh:
clearSpan.AddEvent("clear phase stopped", trace.WithAttributes(attribute.Int("scanned", scanned)))
clearSpan.End()
return
default:
}
scanned++
if _, live := liveNames[name]; live {
continue
}
// Orphaned records can't be cleared by the watcher, so skip the tenant
// API GET entirely — clearTenantPendingDelete would refuse to clear them.
record, err := tw.pendingDeleteStore.Get(ctx, name)
record, err := tw.pendingDeleteStore.Get(clearCtx, name)
if err != nil {
tw.log.Warn("tenant watcher poll cycle: failed to read pending delete record, leaving record",
"tenant", name, "error", err)
@@ -352,6 +401,13 @@ func (tw *TenantWatcher) runPollCycle(ctx context.Context) {
got, err := tw.client.Resource(tenantGVR).Get(ctx, name, metav1.GetOptions{})
// if not found then the tenant was deleted in the tenant api - so keep the record
if apierrors.IsNotFound(err) {
record.Orphaned = true
if err := tw.pendingDeleteStore.Upsert(clearCtx, name, record); err != nil {
tw.log.Warn("tenant watcher poll cycle: failed to mark record orphaned, leaving record",
"tenant", name, "error", err)
} else {
markedOrphaned++
}
leftForDeleter++
continue
}
@@ -368,9 +424,20 @@ func (tw *TenantWatcher) runPollCycle(ctx context.Context) {
raced++
continue
}
tw.clearTenantPendingDelete(ctx, name)
tw.clearTenantPendingDelete(clearCtx, name)
cleared++
}
clearDuration := time.Since(clearStart)
clearSpan.SetAttributes(
attribute.Int("kv_records_scanned", scanned),
attribute.Int("cleared", cleared),
attribute.Int("left_for_deleter", leftForDeleter),
attribute.Int("marked_orphaned", markedOrphaned),
attribute.Int("raced", raced),
attribute.Int("orphaned_skipped", orphanedSkipped),
attribute.Int64("duration_ms", clearDuration.Milliseconds()),
)
clearSpan.End()
tw.log.Info("tenant watcher poll cycle complete",
"live_tenants", len(liveNames),
@@ -378,10 +445,11 @@ func (tw *TenantWatcher) runPollCycle(ctx context.Context) {
"kv_records_scanned", scanned,
"cleared", cleared,
"left_for_deleter", leftForDeleter,
"marked_orphaned", markedOrphaned,
"raced", raced,
"orphaned_skipped", orphanedSkipped,
"list_duration", listDuration,
"clear_duration", time.Since(clearStart),
"clear_duration", clearDuration,
"total_duration", time.Since(start),
)
}
@@ -416,6 +484,23 @@ func (tw *TenantWatcher) handleTenant(ctx context.Context, tenant *unstructured.
// reconcileTenantPendingDelete ensures a pending-delete record exists for the
// tenant and that all of its resources have been labelled.
func (tw *TenantWatcher) reconcileTenantPendingDelete(ctx context.Context, name string, deleteAfter string) {
// only emit span when over threshold or else this gets really spammy
reconcileStart := time.Now()
defer func() {
d := time.Since(reconcileStart)
if d < reconcileTraceThreshold {
return
}
_, span := tracer.Start(ctx, "resource.TenantWatcher.reconcileTenantPendingDelete",
trace.WithTimestamp(reconcileStart),
trace.WithAttributes(
attribute.String("tenant", name),
attribute.Int64("duration_ms", d.Milliseconds()),
),
)
span.End()
}()
// Fast path: if the record exists and labelling is complete, nothing to do.
record, err := tw.pendingDeleteStore.Get(ctx, name)
if err == nil && record.LabelingComplete {
@@ -454,11 +539,19 @@ func (tw *TenantWatcher) reconcileTenantPendingDelete(ctx context.Context, name
// tenantResourcesEditPendingDeleteLabel iterates every resource belonging to
// the given tenant and adds or removes the pending-delete label.
func (tw *TenantWatcher) tenantResourcesEditPendingDeleteLabel(ctx context.Context, tenantName string, addLabel bool) error {
ctx, span := tracer.Start(ctx, "resource.TenantWatcher.tenantResourcesEditPendingDeleteLabel", trace.WithAttributes(
attribute.String("tenant", tenantName),
attribute.Bool("add_label", addLabel),
))
defer span.End()
groupResources, err := tw.dataStore.getGroupResources(ctx)
if err != nil {
return fmt.Errorf("getting group resources: %w", err)
}
span.SetAttributes(attribute.Int("group_resources", len(groupResources)))
var resourcesEdited int
for _, gr := range groupResources {
listKey := ListRequestKey{
Group: gr.Group,
@@ -474,8 +567,10 @@ func (tw *TenantWatcher) tenantResourcesEditPendingDeleteLabel(ctx context.Conte
if err := tw.editResourceLabel(ctx, dataKey, addLabel); err != nil {
return fmt.Errorf("editing label on %s: %w", dataKey.String(), err)
}
resourcesEdited++
}
}
span.SetAttributes(attribute.Int("resources_edited", resourcesEdited))
return nil
}
@@ -579,6 +674,11 @@ func (tw *TenantWatcher) doEditResourceLabel(ctx context.Context, dataKey DataKe
// clearTenantPendingDelete removes the pending-delete record for a tenant from
// the KV store. No-op if the tenant has no record.
func (tw *TenantWatcher) clearTenantPendingDelete(ctx context.Context, name string) {
ctx, span := tracer.Start(ctx, "resource.TenantWatcher.clearTenantPendingDelete", trace.WithAttributes(
attribute.String("tenant", name),
))
defer span.End()
record, err := tw.pendingDeleteStore.Get(ctx, name)
if errors.Is(err, kvpkg.ErrNotFound) {
return
@@ -650,7 +650,7 @@ func TestRunPollCycle(t *testing.T) {
"pending-delete label should have been removed")
})
t.Run("stale record with tenant NotFound preserves record for deleter", func(t *testing.T) {
t.Run("stale record with tenant NotFound is promoted to orphaned", func(t *testing.T) {
tw := newTestTenantWatcher(t)
collector := &writeEventCollector{}
tw.writeEvent = collector.append
@@ -671,6 +671,7 @@ func TestRunPollCycle(t *testing.T) {
r, err := tw.pendingDeleteStore.Get(t.Context(), "tenant-deleted")
require.NoError(t, err, "KV record should be preserved for the deleter")
assert.True(t, r.LabelingComplete)
assert.True(t, r.Orphaned, "record should be promoted to orphaned so future cycles skip the tenant API GET")
})
t.Run("empty live set skips clear phase to avoid wiping records", func(t *testing.T) {
@@ -797,6 +798,37 @@ func TestRunPollCycle(t *testing.T) {
assert.True(t, r.LabelingComplete)
assert.Empty(t, collector.all(), "no unlabel writes for raced tenant")
})
t.Run("orphaned records skip tenant API GET on subsequent cycles", func(t *testing.T) {
tw := newTestTenantWatcher(t)
// First cycle: tenant is 404 in tenant API, so record gets promoted to orphaned.
require.NoError(t, tw.pendingDeleteStore.Upsert(t.Context(), "tenant-deleted", PendingDeleteRecord{
DeleteAfter: "2026-03-01T00:00:00Z",
LabelingComplete: true,
}))
fake := newFakeTenantClient(t, pendingDeleteTenant("tenant-active", "2026-03-01T00:00:00Z"))
var getCount int
fake.PrependReactor("get", "tenants", func(action k8stesting.Action) (bool, runtime.Object, error) {
if ga, ok := action.(k8stesting.GetAction); ok && ga.GetName() == "tenant-deleted" {
getCount++
}
return false, nil, nil
})
tw.client = fake
tw.runPollCycle(t.Context())
require.Equal(t, 1, getCount, "first cycle should GET the tenant once to discover 404")
r, err := tw.pendingDeleteStore.Get(t.Context(), "tenant-deleted")
require.NoError(t, err)
require.True(t, r.Orphaned)
// Second cycle: orphaned record must be skipped before the tenant API GET.
tw.runPollCycle(t.Context())
assert.Equal(t, 1, getCount, "orphaned record should not trigger a tenant API GET on subsequent cycles")
})
}
// newFakeTenantClient builds a fake dynamic client seeded with the given