From e8cfb37b76b763852ac4cfb8d0fc5b930d0e6ef5 Mon Sep 17 00:00:00 2001 From: Evan Phoenix Date: Thu, 24 Sep 2026 18:01:13 +0000 Subject: [PATCH 1/3] Clean orphaned entity index entries on list reads Use CAS-guarded bounded deletes and keep index watch deletes quiet for confirmed missing entities. --- pkg/entity/cleanup.go | 71 +++++++++++++- pkg/entity/cleanup_test.go | 93 ++++++++++++++++++ pkg/entity/store.go | 16 +--- pkg/entity/store_test.go | 10 +- servers/entityserver/entityserver.go | 38 +++++--- servers/entityserver/entityserver_test.go | 112 ++++++++++++++++++++++ 6 files changed, 310 insertions(+), 30 deletions(-) diff --git a/pkg/entity/cleanup.go b/pkg/entity/cleanup.go index 0b1392657..a68db8d7e 100644 --- a/pkg/entity/cleanup.go +++ b/pkg/entity/cleanup.go @@ -8,10 +8,79 @@ import ( "strings" "time" + "github.com/mr-tron/base58" "go.etcd.io/etcd/api/v3/mvccpb" clientv3 "go.etcd.io/etcd/client/v3" ) +// CleanupOrphanedIndexEntry removes the entries for one index/id only if the +// backing entity is still absent and the entries have not changed. A list may +// have read an older revision, so both checks belong in the delete transaction. +// It returns true only when it actually removed an entry. +func (s *EtcdStore) CleanupOrphanedIndexEntry(ctx context.Context, attr Attr, id Id) (bool, error) { + prefix, err := s.IndexPrefix(ctx, attr) + if err != nil { + return false, err + } + key := prefix + base58.Encode([]byte(id)) + + plain, err := s.client.Get(ctx, key) + if err != nil { + return false, err + } + sessions, err := s.client.Get(ctx, key+"/", clientv3.WithPrefix()) + if err != nil { + return false, err + } + kvs := append(plain.Kvs, sessions.Kvs...) + if len(kvs) == 0 { + return false, nil + } + + entries := make([]*mvccpb.KeyValue, 0, len(kvs)) + for _, kv := range kvs { + if string(kv.Value) == string(id) { + entries = append(entries, kv) + } + } + removed := false + // One primary compare plus one compare and delete per entry must fit + // etcd's transaction operation limit, even for an old session backlog. + for start := 0; start < len(entries); start += (etcdMaxTxnOps - 1) / 2 { + end := min(start+(etcdMaxTxnOps-1)/2, len(entries)) + batch := entries[start:end] + cmps := []clientv3.Cmp{clientv3.Compare(clientv3.CreateRevision(s.buildKey(id)), "=", 0)} + ops := make([]clientv3.Op, 0, len(batch)) + for _, kv := range batch { + entryKey := string(kv.Key) + cmps = append(cmps, clientv3.Compare(clientv3.ModRevision(entryKey), "=", kv.ModRevision)) + ops = append(ops, clientv3.OpDelete(entryKey)) + } + resp, err := s.client.Txn(ctx).If(cmps...).Then(ops...).Commit() + if err != nil { + return removed, err + } + if resp.Succeeded { + removed = true + continue + } + // A concurrent write changed at least one slot (or recreated the + // entity). Retry individually so unrelated stale slots still drain. + for _, kv := range batch { + entryKey := string(kv.Key) + resp, err := s.client.Txn(ctx).If( + clientv3.Compare(clientv3.CreateRevision(s.buildKey(id)), "=", 0), + clientv3.Compare(clientv3.ModRevision(entryKey), "=", kv.ModRevision), + ).Then(clientv3.OpDelete(entryKey)).Commit() + if err != nil { + return removed, err + } + removed = removed || resp.Succeeded + } + } + return removed, nil +} + // cleanupDeleteBatchSize is how many stale entries we delete between rate-limit // pauses. Deletes are issued one CAS'd transaction at a time (see // CleanupStaleCollectionEntries); this only governs how often we sleep for @@ -281,7 +350,7 @@ func (s *EtcdStore) resolveJustified( batch []collectionEntry, ) (justified map[Id]map[string]bool, unverifiable map[Id]bool, err error) { ids := distinctEntityIDs(batch) - entities, undecodable, err := s.getEntities(ctx, ids, false, 0) + entities, undecodable, err := s.getEntities(ctx, ids, 0) if err != nil { return nil, nil, fmt.Errorf("cleanup: failed to resolve entities: %w", err) } diff --git a/pkg/entity/cleanup_test.go b/pkg/entity/cleanup_test.go index 35de43ba0..020f5ac28 100644 --- a/pkg/entity/cleanup_test.go +++ b/pkg/entity/cleanup_test.go @@ -1,6 +1,7 @@ package entity import ( + "fmt" "log/slog" "testing" @@ -20,6 +21,98 @@ func seedStaleEntry(t *testing.T, store *EtcdStore, client *clientv3.Client, col return key } +func TestCleanupOrphanedIndexEntry_ConcurrentChanges(t *testing.T) { + store, client := setupReindexTestStore(t) + ctx := t.Context() + _, err := store.CreateEntity(ctx, New( + Ident, "test/kind", Doc, "indexed kind", Cardinality, CardinalityOne, + Type, TypeStr, Index, true, + )) + require.NoError(t, err) + index := String(Id("test/kind"), "widget") + id := Id("recreated") + key := seedStaleEntry(t, store, client, collectionSegmentFor(index), id) + + // A new backing entity after the listing makes the observed miss obsolete. + _, err = store.CreateEntity(ctx, New(Ref(DBId, id), index)) + require.NoError(t, err) + removed, err := store.CleanupOrphanedIndexEntry(ctx, index, id) + require.NoError(t, err) + assert.False(t, removed) + response, err := client.Get(ctx, key) + require.NoError(t, err) + assert.Len(t, response.Kvs, 1) + + // A stale listing must not delete a slot now pointing at a different id. + other := Id("other") + _, err = client.Put(ctx, key, other.String()) + require.NoError(t, err) + removed, err = store.CleanupOrphanedIndexEntry(ctx, index, id) + require.NoError(t, err) + assert.False(t, removed) + response, err = client.Get(ctx, key) + require.NoError(t, err) + assert.Equal(t, other.String(), string(response.Kvs[0].Value)) +} + +func TestCleanupOrphanedIndexEntry_LeavesUnreadableEntity(t *testing.T) { + store, client := setupReindexTestStore(t) + ctx := t.Context() + _, err := store.CreateEntity(ctx, New( + Ident, "test/kind", Doc, "indexed kind", Cardinality, CardinalityOne, + Type, TypeStr, Index, true, + )) + require.NoError(t, err) + index := String(Id("test/kind"), "widget") + live, err := store.CreateEntity(ctx, New(Ident, "live", index)) + require.NoError(t, err) + key := store.Prefix() + "/collections/" + collectionSegmentFor(index) + "/" + base58.Encode([]byte(live.Id())) + _, err = client.Put(ctx, store.buildKey(live.Id()), "invalid cbor") + require.NoError(t, err) + removed, err := store.CleanupOrphanedIndexEntry(ctx, index, live.Id()) + require.NoError(t, err) + assert.False(t, removed) + resp, err := client.Get(ctx, key) + require.NoError(t, err) + assert.Len(t, resp.Kvs, 1) +} + +func TestCleanupOrphanedIndexEntry_BatchesSessionVariants(t *testing.T) { + store, client := setupReindexTestStore(t) + ctx := t.Context() + _, err := store.CreateEntity(ctx, New( + Ident, "test/kind", Doc, "indexed kind", Cardinality, CardinalityOne, + Type, TypeStr, Index, true, + )) + require.NoError(t, err) + index := String(Id("test/kind"), "widget") + id := Id("orphan") + key := seedStaleEntry(t, store, client, collectionSegmentFor(index), id) + for i := range 140 { + _, err := client.Put(ctx, fmt.Sprintf("%s/session-%03d", key, i), id.String()) + require.NoError(t, err) + } + // Another slot sharing the byte prefix must not be deleted. + _, err = client.Put(ctx, key+"suffix", "other") + require.NoError(t, err) + + removed, err := store.CleanupOrphanedIndexEntry(ctx, index, id) + require.NoError(t, err) + assert.True(t, removed) + resp, err := client.Get(ctx, key) + require.NoError(t, err) + assert.Empty(t, resp.Kvs) + resp, err = client.Get(ctx, key+"/", clientv3.WithPrefix()) + require.NoError(t, err) + assert.Empty(t, resp.Kvs) + resp, err = client.Get(ctx, key+"suffix") + require.NoError(t, err) + assert.Len(t, resp.Kvs, 1) + removed, err = store.CleanupOrphanedIndexEntry(ctx, index, id) + require.NoError(t, err) + assert.False(t, removed) +} + func TestCleanup_RemovesStaleKeepsLive(t *testing.T) { store, client := setupReindexTestStore(t) ctx := t.Context() diff --git a/pkg/entity/store.go b/pkg/entity/store.go index fa6adf468..1fc506585 100644 --- a/pkg/entity/store.go +++ b/pkg/entity/store.go @@ -487,7 +487,7 @@ func (s *EtcdStore) GetEntity(ctx context.Context, id Id) (*Entity, error) { // alone would hand back an entity missing every session-scoped attribute, // which for a node is its status. func (s *EtcdStore) GetEntityAtRevision(ctx context.Context, id Id, rev int64) (*Entity, error) { - entities, undecodable, err := s.getEntities(ctx, []Id{id}, false, rev) + entities, undecodable, err := s.getEntities(ctx, []Id{id}, rev) if err != nil { return nil, fmt.Errorf("failed to get entity at revision %d: %w", rev, err) } @@ -501,7 +501,7 @@ func (s *EtcdStore) GetEntityAtRevision(ctx context.Context, id Id, rev int64) ( } func (s *EtcdStore) GetEntities(ctx context.Context, ids []Id) ([]*Entity, error) { - entities, _, err := s.getEntities(ctx, ids, true, 0) + entities, _, err := s.getEntities(ctx, ids, 0) return entities, err } @@ -531,7 +531,7 @@ func (s *EtcdStore) ListIndexEntitiesPage( return nil, err } - entities, undecodable, err := s.getEntities(ctx, page.Ids, false, page.Revision) + entities, undecodable, err := s.getEntities(ctx, page.Ids, page.Revision) if err != nil { return nil, err } @@ -547,9 +547,7 @@ func (s *EtcdStore) ListIndexEntitiesPage( } // getEntities reads entities in batches, leaving nil in the result for any id -// that is absent. warnMissing is false for callers where a miss is expected -// input rather than a surprise, such as the index sweep resolving ids read out -// of the index itself. +// that is absent. Callers decide whether a missing entity needs attention. // // A non-zero rev reads every key as of that revision instead of the latest. // @@ -560,7 +558,6 @@ func (s *EtcdStore) ListIndexEntitiesPage( func (s *EtcdStore) getEntities( ctx context.Context, ids []Id, - warnMissing bool, rev int64, ) (entities []*Entity, undecodable map[Id]bool, err error) { undecodable = map[Id]bool{} @@ -615,9 +612,6 @@ func (s *EtcdStore) getEntities( primaryResp := tr.Responses[primaryIdx].GetResponseRange() if len(primaryResp.Kvs) == 0 { // Entity not found, leave nil in the result array - if warnMissing { - s.log.Warn("failed to get primary entity from etcd", "id", batchIds[i]) - } continue } @@ -625,7 +619,7 @@ func (s *EtcdStore) getEntities( err = decoder.Unmarshal(primaryResp.Kvs[0].Value, &entity) if err != nil { // The key is there, so this entity exists; we just cannot read - // it. Always worth a line, whatever warnMissing says. + // it. Always worth a line even though missing keys are expected. s.log.Error("failed to decode entity from etcd", "id", batchIds[i], "error", err) undecodable[batchIds[i]] = true continue diff --git a/pkg/entity/store_test.go b/pkg/entity/store_test.go index e3a303984..865a1d9a4 100644 --- a/pkg/entity/store_test.go +++ b/pkg/entity/store_test.go @@ -3125,7 +3125,7 @@ func TestEtcdStore_getEntitiesAtRevision(t *testing.T) { ids := []Id{Id(created.Id())} - pinned, undecodable, err := store.getEntities(t.Context(), ids, false, before) + pinned, undecodable, err := store.getEntities(t.Context(), ids, before) require.NoError(t, err) require.Len(t, pinned, 1) require.NotNil(t, pinned[0], "the entity existed at that revision") @@ -3135,7 +3135,7 @@ func TestEtcdStore_getEntitiesAtRevision(t *testing.T) { assert.Equal(t, "before", doc.Value.String(), "a pinned read must see the entity as it was, not as it is") - current, _, err := store.getEntities(t.Context(), ids, false, 0) + current, _, err := store.getEntities(t.Context(), ids, 0) require.NoError(t, err) require.Len(t, current, 1) require.NotNil(t, current[0]) @@ -3159,7 +3159,7 @@ func TestEtcdStore_getEntitiesAtRevision(t *testing.T) { // slice and silently shift every entity after it onto the wrong id. ids := []Id{"missing-before", Id(created.Id()), "missing-after"} - entities, undecodable, err := store.getEntities(t.Context(), ids, false, 0) + entities, undecodable, err := store.getEntities(t.Context(), ids, 0) require.NoError(t, err) require.Len(t, entities, 3) assert.Nil(t, entities[0]) @@ -3183,7 +3183,7 @@ func TestEtcdStore_getEntitiesAtRevision(t *testing.T) { require.NoError(t, err) entities, _, err := store.getEntities(t.Context(), - []Id{Id(first.Id()), Id(later.Id())}, false, first.GetRevision()) + []Id{Id(first.Id()), Id(later.Id())}, first.GetRevision()) require.NoError(t, err) require.Len(t, entities, 2) assert.NotNil(t, entities[0]) @@ -3193,7 +3193,7 @@ func TestEtcdStore_getEntitiesAtRevision(t *testing.T) { t.Run("returns an indexable map for an empty request", func(t *testing.T) { store, _ := setupTestEtcdStore(t) - entities, undecodable, err := store.getEntities(t.Context(), nil, false, 0) + entities, undecodable, err := store.getEntities(t.Context(), nil, 0) require.NoError(t, err) assert.Empty(t, entities) diff --git a/servers/entityserver/entityserver.go b/servers/entityserver/entityserver.go index 9d74a8f79..8ada253b3 100644 --- a/servers/entityserver/entityserver.go +++ b/servers/entityserver/entityserver.go @@ -591,13 +591,14 @@ func (e *EntityServer) WatchIndex(ctx context.Context, req *entityserver_v1alpha // and the entity key are removed together in one atomic txn, so // the entity is already gone at this event's revision; read it at // the prior revision to recover what was deleted. - en, err := e.Store.GetEntity(ctx, entityId) - if err != nil { - en, err = e.Store.GetEntityAtRevision(ctx, entityId, event.Kv.ModRevision-1) + en, currentErr := e.Store.GetEntity(ctx, entityId) + readErr := currentErr + if currentErr != nil { + en, readErr = e.Store.GetEntityAtRevision(ctx, entityId, event.Kv.ModRevision-1) } - if err != nil { - e.Log.Error("failed to get entity for delete event", "error", err, "id", entityId) - } else { + if readErr != nil && (!isNotFound(currentErr) || !isNotFound(readErr)) { + e.Log.Error("failed to get entity for delete event", "error", readErr, "id", entityId) + } else if readErr == nil { var rpcEntity entityserver_v1alpha.Entity rpcEntity.SetId(en.Id().String()) rpcEntity.SetCreatedAt(en.GetCreatedAt().UnixMilli()) @@ -670,9 +671,7 @@ func (e *EntityServer) List(ctx context.Context, req *entityserver_v1alpha.Entit var ret []*entityserver_v1alpha.Entity for i, entity := range entities { if entity == nil { - e.Log.Error("entity in index but not in store, skipping", - "id", ids[i], - "index", index) + e.cleanupOrphan(ctx, index, ids[i]) continue } @@ -1048,7 +1047,7 @@ func (e *EntityServer) entityPage( return nil, fmt.Errorf("failed to get entities: %w", err) } - return e.resolve(index, ids, entities, nil, next, total, 0), nil + return e.resolve(ctx, index, ids, entities, nil, next, total, 0), nil } page, err := e.Store.ListIndexEntitiesPage(ctx, index, cursor, limit) @@ -1056,7 +1055,20 @@ func (e *EntityServer) entityPage( return nil, fmt.Errorf("failed to list entities: %w", err) } - return e.resolve(index, page.Ids, page.Entities, page.Undecodable, page.Cursor, page.Total, page.Revision), nil + return e.resolve(ctx, index, page.Ids, page.Entities, page.Undecodable, page.Cursor, page.Total, page.Revision), nil +} + +func (e *EntityServer) cleanupOrphan(ctx context.Context, index entity.Attr, id entity.Id) { + store, ok := e.Store.(*entity.EtcdStore) + if !ok || index.ID == entity.AttrSession || index.ID == entity.DBId { + return + } + removed, err := store.CleanupOrphanedIndexEntry(ctx, index, id) + if err != nil { + e.Log.Warn("failed to clean up orphaned index entry", "id", id, "index", index, "error", err) + } else if removed { + e.Log.Warn("cleaned up orphaned index entry", "id", id, "index", index) + } } // resolve drops the ids the store could not answer for and keeps the reported @@ -1067,6 +1079,7 @@ func (e *EntityServer) entityPage( // work; a key that will not decode is a corrupt entity nobody should assume is // gone. func (e *EntityServer) resolve( + ctx context.Context, index entity.Attr, ids []entity.Id, entities []*entity.Entity, @@ -1086,8 +1099,7 @@ func (e *EntityServer) resolve( e.Log.Error("entity in index cannot be decoded, skipping", "id", ids[i], "index", index) } else { - e.Log.Error("entity in index but not in store, skipping", - "id", ids[i], "index", index) + e.cleanupOrphan(ctx, index, ids[i]) } if total > 0 { diff --git a/servers/entityserver/entityserver_test.go b/servers/entityserver/entityserver_test.go index 067cbb6ed..08634d230 100644 --- a/servers/entityserver/entityserver_test.go +++ b/servers/entityserver/entityserver_test.go @@ -1,15 +1,18 @@ package entityserver import ( + "bytes" "context" "encoding/json" "errors" "fmt" "log/slog" "slices" + "strings" "testing" "time" + "github.com/mr-tron/base58" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" "go.etcd.io/etcd/api/v3/etcdserverpb" @@ -1243,3 +1246,112 @@ func TestEntityServer_ListPage(t *testing.T) { }) } + +func TestEntityServer_ListsCleanOrphansOnce(t *testing.T) { + ctx := t.Context() + client, prefix := setupTestEtcd(t) + store, err := entity.NewEtcdStore(ctx, slog.Default(), client, prefix) + require.NoError(t, err) + index := entity.String(entity.Id("test/kind"), "widget") + _, err = store.CreateEntity(ctx, entity.New( + entity.Ident, "test/kind", entity.Doc, "indexed kind", + entity.Cardinality, entity.CardinalityOne, entity.Type, entity.TypeStr, entity.Index, true, + )) + require.NoError(t, err) + live, err := store.CreateEntity(ctx, entity.New(entity.Ident, "live", index)) + require.NoError(t, err) + indexPrefix, err := store.IndexPrefix(ctx, index) + require.NoError(t, err) + + var logs bytes.Buffer + server, err := NewEntityServer(slog.New(slog.NewTextHandler(&logs, nil)), store) + require.NoError(t, err) + sc := v1alpha.EntityAccessClient{Client: rpc.LocalClient(v1alpha.AdaptEntityAccess(server))} + + for _, paged := range []bool{false, true} { + id := entity.Id(fmt.Sprintf("missing-%t", paged)) + key := indexPrefix + base58.Encode([]byte(id)) + entryKey := key + if paged { + entryKey += "/old-session" + } + _, err := client.Put(ctx, entryKey, id.String()) + require.NoError(t, err) + for range 2 { + if paged { + page, err := sc.ListPage(ctx, index, "", 10) + require.NoError(t, err) + require.Len(t, page.Values(), 1) + require.Equal(t, int64(1), page.Total()) + require.Equal(t, live.Id().String(), page.Values()[0].Id()) + } else { + list, err := sc.List(ctx, index) + require.NoError(t, err) + require.Len(t, list.Values(), 1) + require.Equal(t, live.Id().String(), list.Values()[0].Id()) + } + } + response, err := client.Get(ctx, key, clientv3.WithPrefix()) + require.NoError(t, err) + require.Empty(t, response.Kvs) + } + require.Equal(t, 2, strings.Count(logs.String(), "cleaned up orphaned index entry")) + require.NotContains(t, logs.String(), "entity in index but not in store") +} + +func TestEntityServer_OrphanCleanupWatchDelete(t *testing.T) { + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + client, prefix := setupTestEtcd(t) + var logs bytes.Buffer + log := slog.New(slog.NewTextHandler(&logs, nil)) + store, err := entity.NewEtcdStore(ctx, log, client, prefix) + require.NoError(t, err) + index := entity.String(entity.Id("test/kind"), "widget") + _, err = store.CreateEntity(ctx, entity.New( + entity.Ident, "test/kind", entity.Doc, "indexed kind", + entity.Cardinality, entity.CardinalityOne, entity.Type, entity.TypeStr, entity.Index, true, + )) + require.NoError(t, err) + server, err := NewEntityServer(log, store) + require.NoError(t, err) + sc := v1alpha.EntityAccessClient{Client: rpc.LocalClient(v1alpha.AdaptEntityAccess(server))} + id := entity.Id("orphan") + indexPrefix, err := store.IndexPrefix(ctx, index) + require.NoError(t, err) + _, err = client.Put(ctx, indexPrefix+base58.Encode([]byte(id)), id.String()) + require.NoError(t, err) + listed, err := client.Get(ctx, indexPrefix, clientv3.WithPrefix()) + require.NoError(t, err) + + deletes := make(chan *v1alpha.EntityOp, 1) + watchDone := make(chan error, 1) + go func() { + _, err := sc.WatchIndex(ctx, index, listed.Header.Revision+1, stream.Callback(func(op *v1alpha.EntityOp) error { + if op.Operation() == int64(v1alpha.EntityOperationDelete) { + deletes <- op + } + return nil + })) + watchDone <- err + }() + + result, err := sc.List(ctx, index) + require.NoError(t, err) + require.Empty(t, result.Values()) + select { + case op := <-deletes: + require.Equal(t, id.String(), op.EntityId()) + require.False(t, op.HasEntity()) + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for orphan delete event") + } + cancel() + select { + case <-watchDone: + case <-time.After(5 * time.Second): + t.Fatal("timed out waiting for watch to stop") + } + require.Equal(t, 1, strings.Count(logs.String(), "level=WARN msg=\"cleaned up orphaned index entry\"")) + require.NotContains(t, logs.String(), "level=ERROR") +} From 5d8617a4e6a25705fdd66e0a9a87c56b741edd81 Mon Sep 17 00:00:00 2001 From: Evan Phoenix Date: Thu, 24 Sep 2026 18:21:46 +0000 Subject: [PATCH 2/3] Bound orphan cleanup work per entity listing --- servers/entityserver/entityserver.go | 14 ++++++- servers/entityserver/entityserver_test.go | 48 +++++++++++++++++++++++ 2 files changed, 60 insertions(+), 2 deletions(-) diff --git a/servers/entityserver/entityserver.go b/servers/entityserver/entityserver.go index 8ada253b3..46cdc5bea 100644 --- a/servers/entityserver/entityserver.go +++ b/servers/entityserver/entityserver.go @@ -669,9 +669,13 @@ func (e *EntityServer) List(ctx context.Context, req *entityserver_v1alpha.Entit } var ret []*entityserver_v1alpha.Entity + cleanups := 0 for i, entity := range entities { if entity == nil { - e.cleanupOrphan(ctx, index, ids[i]) + if cleanups < maxIndexCleanupsPerList { + e.cleanupOrphan(ctx, index, ids[i]) + cleanups++ + } continue } @@ -691,6 +695,10 @@ func (e *EntityServer) List(ctx context.Context, req *entityserver_v1alpha.Entit return nil } +// A large legacy backlog should not turn a listing into an unbounded series +// of etcd writes. Later lists and the background sweep drain the rest. +const maxIndexCleanupsPerList = 16 + // ListPage reads a bounded page of entities from an index. func (e *EntityServer) ListPage(ctx context.Context, req *entityserver_v1alpha.EntityAccessListPage) error { args := req.Args() @@ -1088,6 +1096,7 @@ func (e *EntityServer) resolve( total, revision int64, ) *resolvedPage { resolved := make([]*entity.Entity, 0, len(entities)) + cleanups := 0 for i, ent := range entities { if ent != nil { @@ -1098,8 +1107,9 @@ func (e *EntityServer) resolve( if undecodable[ids[i]] { e.Log.Error("entity in index cannot be decoded, skipping", "id", ids[i], "index", index) - } else { + } else if cleanups < maxIndexCleanupsPerList { e.cleanupOrphan(ctx, index, ids[i]) + cleanups++ } if total > 0 { diff --git a/servers/entityserver/entityserver_test.go b/servers/entityserver/entityserver_test.go index 08634d230..0180e26cf 100644 --- a/servers/entityserver/entityserver_test.go +++ b/servers/entityserver/entityserver_test.go @@ -1299,6 +1299,54 @@ func TestEntityServer_ListsCleanOrphansOnce(t *testing.T) { require.NotContains(t, logs.String(), "entity in index but not in store") } +func TestEntityServer_ListCleanupIsBounded(t *testing.T) { + for _, paged := range []bool{false, true} { + t.Run(fmt.Sprintf("paged=%t", paged), func(t *testing.T) { + ctx := t.Context() + client, prefix := setupTestEtcd(t) + store, err := entity.NewEtcdStore(ctx, slog.Default(), client, prefix) + require.NoError(t, err) + index := entity.String(entity.Id("test/kind"), "widget") + _, err = store.CreateEntity(ctx, entity.New( + entity.Ident, "test/kind", entity.Doc, "indexed kind", + entity.Cardinality, entity.CardinalityOne, entity.Type, entity.TypeStr, entity.Index, true, + )) + require.NoError(t, err) + live, err := store.CreateEntity(ctx, entity.New(entity.Ident, "live", index)) + require.NoError(t, err) + indexPrefix, err := store.IndexPrefix(ctx, index) + require.NoError(t, err) + for i := range maxIndexCleanupsPerList + 2 { + id := entity.Id(fmt.Sprintf("orphan-%02d", i)) + _, err := client.Put(ctx, indexPrefix+base58.Encode([]byte(id)), id.String()) + require.NoError(t, err) + } + + var logs bytes.Buffer + server, err := NewEntityServer(slog.New(slog.NewTextHandler(&logs, nil)), store) + require.NoError(t, err) + sc := v1alpha.EntityAccessClient{Client: rpc.LocalClient(v1alpha.AdaptEntityAccess(server))} + for pass, remaining := range []int64{2, 0, 0} { + if paged { + page, err := sc.ListPage(ctx, index, "", 100) + require.NoError(t, err) + require.Len(t, page.Values(), 1) + require.Equal(t, live.Id().String(), page.Values()[0].Id()) + } else { + list, err := sc.List(ctx, index) + require.NoError(t, err) + require.Len(t, list.Values(), 1) + require.Equal(t, live.Id().String(), list.Values()[0].Id()) + } + entries, err := client.Get(ctx, indexPrefix, clientv3.WithPrefix()) + require.NoError(t, err) + require.Equal(t, remaining+1, entries.Count, "pass %d", pass) + } + require.Equal(t, maxIndexCleanupsPerList+2, strings.Count(logs.String(), "cleaned up orphaned index entry")) + }) + } +} + func TestEntityServer_OrphanCleanupWatchDelete(t *testing.T) { ctx, cancel := context.WithCancel(t.Context()) defer cancel() From 81a985d21b56c39e14ea002bfd36d9d6301b448f Mon Sep 17 00:00:00 2001 From: Evan Phoenix Date: Fri, 25 Sep 2026 01:36:58 +0000 Subject: [PATCH 3/3] Keep orphan repair in background GC and quiet list misses --- pkg/entity/cleanup.go | 69 -------------- pkg/entity/cleanup_test.go | 93 ------------------- servers/entityserver/entityserver.go | 36 ++------ servers/entityserver/entityserver_test.go | 108 ++++++---------------- 4 files changed, 34 insertions(+), 272 deletions(-) diff --git a/pkg/entity/cleanup.go b/pkg/entity/cleanup.go index a68db8d7e..f24dfe712 100644 --- a/pkg/entity/cleanup.go +++ b/pkg/entity/cleanup.go @@ -8,79 +8,10 @@ import ( "strings" "time" - "github.com/mr-tron/base58" "go.etcd.io/etcd/api/v3/mvccpb" clientv3 "go.etcd.io/etcd/client/v3" ) -// CleanupOrphanedIndexEntry removes the entries for one index/id only if the -// backing entity is still absent and the entries have not changed. A list may -// have read an older revision, so both checks belong in the delete transaction. -// It returns true only when it actually removed an entry. -func (s *EtcdStore) CleanupOrphanedIndexEntry(ctx context.Context, attr Attr, id Id) (bool, error) { - prefix, err := s.IndexPrefix(ctx, attr) - if err != nil { - return false, err - } - key := prefix + base58.Encode([]byte(id)) - - plain, err := s.client.Get(ctx, key) - if err != nil { - return false, err - } - sessions, err := s.client.Get(ctx, key+"/", clientv3.WithPrefix()) - if err != nil { - return false, err - } - kvs := append(plain.Kvs, sessions.Kvs...) - if len(kvs) == 0 { - return false, nil - } - - entries := make([]*mvccpb.KeyValue, 0, len(kvs)) - for _, kv := range kvs { - if string(kv.Value) == string(id) { - entries = append(entries, kv) - } - } - removed := false - // One primary compare plus one compare and delete per entry must fit - // etcd's transaction operation limit, even for an old session backlog. - for start := 0; start < len(entries); start += (etcdMaxTxnOps - 1) / 2 { - end := min(start+(etcdMaxTxnOps-1)/2, len(entries)) - batch := entries[start:end] - cmps := []clientv3.Cmp{clientv3.Compare(clientv3.CreateRevision(s.buildKey(id)), "=", 0)} - ops := make([]clientv3.Op, 0, len(batch)) - for _, kv := range batch { - entryKey := string(kv.Key) - cmps = append(cmps, clientv3.Compare(clientv3.ModRevision(entryKey), "=", kv.ModRevision)) - ops = append(ops, clientv3.OpDelete(entryKey)) - } - resp, err := s.client.Txn(ctx).If(cmps...).Then(ops...).Commit() - if err != nil { - return removed, err - } - if resp.Succeeded { - removed = true - continue - } - // A concurrent write changed at least one slot (or recreated the - // entity). Retry individually so unrelated stale slots still drain. - for _, kv := range batch { - entryKey := string(kv.Key) - resp, err := s.client.Txn(ctx).If( - clientv3.Compare(clientv3.CreateRevision(s.buildKey(id)), "=", 0), - clientv3.Compare(clientv3.ModRevision(entryKey), "=", kv.ModRevision), - ).Then(clientv3.OpDelete(entryKey)).Commit() - if err != nil { - return removed, err - } - removed = removed || resp.Succeeded - } - } - return removed, nil -} - // cleanupDeleteBatchSize is how many stale entries we delete between rate-limit // pauses. Deletes are issued one CAS'd transaction at a time (see // CleanupStaleCollectionEntries); this only governs how often we sleep for diff --git a/pkg/entity/cleanup_test.go b/pkg/entity/cleanup_test.go index 020f5ac28..35de43ba0 100644 --- a/pkg/entity/cleanup_test.go +++ b/pkg/entity/cleanup_test.go @@ -1,7 +1,6 @@ package entity import ( - "fmt" "log/slog" "testing" @@ -21,98 +20,6 @@ func seedStaleEntry(t *testing.T, store *EtcdStore, client *clientv3.Client, col return key } -func TestCleanupOrphanedIndexEntry_ConcurrentChanges(t *testing.T) { - store, client := setupReindexTestStore(t) - ctx := t.Context() - _, err := store.CreateEntity(ctx, New( - Ident, "test/kind", Doc, "indexed kind", Cardinality, CardinalityOne, - Type, TypeStr, Index, true, - )) - require.NoError(t, err) - index := String(Id("test/kind"), "widget") - id := Id("recreated") - key := seedStaleEntry(t, store, client, collectionSegmentFor(index), id) - - // A new backing entity after the listing makes the observed miss obsolete. - _, err = store.CreateEntity(ctx, New(Ref(DBId, id), index)) - require.NoError(t, err) - removed, err := store.CleanupOrphanedIndexEntry(ctx, index, id) - require.NoError(t, err) - assert.False(t, removed) - response, err := client.Get(ctx, key) - require.NoError(t, err) - assert.Len(t, response.Kvs, 1) - - // A stale listing must not delete a slot now pointing at a different id. - other := Id("other") - _, err = client.Put(ctx, key, other.String()) - require.NoError(t, err) - removed, err = store.CleanupOrphanedIndexEntry(ctx, index, id) - require.NoError(t, err) - assert.False(t, removed) - response, err = client.Get(ctx, key) - require.NoError(t, err) - assert.Equal(t, other.String(), string(response.Kvs[0].Value)) -} - -func TestCleanupOrphanedIndexEntry_LeavesUnreadableEntity(t *testing.T) { - store, client := setupReindexTestStore(t) - ctx := t.Context() - _, err := store.CreateEntity(ctx, New( - Ident, "test/kind", Doc, "indexed kind", Cardinality, CardinalityOne, - Type, TypeStr, Index, true, - )) - require.NoError(t, err) - index := String(Id("test/kind"), "widget") - live, err := store.CreateEntity(ctx, New(Ident, "live", index)) - require.NoError(t, err) - key := store.Prefix() + "/collections/" + collectionSegmentFor(index) + "/" + base58.Encode([]byte(live.Id())) - _, err = client.Put(ctx, store.buildKey(live.Id()), "invalid cbor") - require.NoError(t, err) - removed, err := store.CleanupOrphanedIndexEntry(ctx, index, live.Id()) - require.NoError(t, err) - assert.False(t, removed) - resp, err := client.Get(ctx, key) - require.NoError(t, err) - assert.Len(t, resp.Kvs, 1) -} - -func TestCleanupOrphanedIndexEntry_BatchesSessionVariants(t *testing.T) { - store, client := setupReindexTestStore(t) - ctx := t.Context() - _, err := store.CreateEntity(ctx, New( - Ident, "test/kind", Doc, "indexed kind", Cardinality, CardinalityOne, - Type, TypeStr, Index, true, - )) - require.NoError(t, err) - index := String(Id("test/kind"), "widget") - id := Id("orphan") - key := seedStaleEntry(t, store, client, collectionSegmentFor(index), id) - for i := range 140 { - _, err := client.Put(ctx, fmt.Sprintf("%s/session-%03d", key, i), id.String()) - require.NoError(t, err) - } - // Another slot sharing the byte prefix must not be deleted. - _, err = client.Put(ctx, key+"suffix", "other") - require.NoError(t, err) - - removed, err := store.CleanupOrphanedIndexEntry(ctx, index, id) - require.NoError(t, err) - assert.True(t, removed) - resp, err := client.Get(ctx, key) - require.NoError(t, err) - assert.Empty(t, resp.Kvs) - resp, err = client.Get(ctx, key+"/", clientv3.WithPrefix()) - require.NoError(t, err) - assert.Empty(t, resp.Kvs) - resp, err = client.Get(ctx, key+"suffix") - require.NoError(t, err) - assert.Len(t, resp.Kvs, 1) - removed, err = store.CleanupOrphanedIndexEntry(ctx, index, id) - require.NoError(t, err) - assert.False(t, removed) -} - func TestCleanup_RemovesStaleKeepsLive(t *testing.T) { store, client := setupReindexTestStore(t) ctx := t.Context() diff --git a/servers/entityserver/entityserver.go b/servers/entityserver/entityserver.go index 46cdc5bea..3251a3b9d 100644 --- a/servers/entityserver/entityserver.go +++ b/servers/entityserver/entityserver.go @@ -669,13 +669,10 @@ func (e *EntityServer) List(ctx context.Context, req *entityserver_v1alpha.Entit } var ret []*entityserver_v1alpha.Entity - cleanups := 0 for i, entity := range entities { if entity == nil { - if cleanups < maxIndexCleanupsPerList { - e.cleanupOrphan(ctx, index, ids[i]) - cleanups++ - } + e.Log.Debug("entity in index but not in store, skipping", + "id", ids[i], "index", index) continue } @@ -695,10 +692,6 @@ func (e *EntityServer) List(ctx context.Context, req *entityserver_v1alpha.Entit return nil } -// A large legacy backlog should not turn a listing into an unbounded series -// of etcd writes. Later lists and the background sweep drain the rest. -const maxIndexCleanupsPerList = 16 - // ListPage reads a bounded page of entities from an index. func (e *EntityServer) ListPage(ctx context.Context, req *entityserver_v1alpha.EntityAccessListPage) error { args := req.Args() @@ -1055,7 +1048,7 @@ func (e *EntityServer) entityPage( return nil, fmt.Errorf("failed to get entities: %w", err) } - return e.resolve(ctx, index, ids, entities, nil, next, total, 0), nil + return e.resolve(index, ids, entities, nil, next, total, 0), nil } page, err := e.Store.ListIndexEntitiesPage(ctx, index, cursor, limit) @@ -1063,20 +1056,7 @@ func (e *EntityServer) entityPage( return nil, fmt.Errorf("failed to list entities: %w", err) } - return e.resolve(ctx, index, page.Ids, page.Entities, page.Undecodable, page.Cursor, page.Total, page.Revision), nil -} - -func (e *EntityServer) cleanupOrphan(ctx context.Context, index entity.Attr, id entity.Id) { - store, ok := e.Store.(*entity.EtcdStore) - if !ok || index.ID == entity.AttrSession || index.ID == entity.DBId { - return - } - removed, err := store.CleanupOrphanedIndexEntry(ctx, index, id) - if err != nil { - e.Log.Warn("failed to clean up orphaned index entry", "id", id, "index", index, "error", err) - } else if removed { - e.Log.Warn("cleaned up orphaned index entry", "id", id, "index", index) - } + return e.resolve(index, page.Ids, page.Entities, page.Undecodable, page.Cursor, page.Total, page.Revision), nil } // resolve drops the ids the store could not answer for and keeps the reported @@ -1087,7 +1067,6 @@ func (e *EntityServer) cleanupOrphan(ctx context.Context, index entity.Attr, id // work; a key that will not decode is a corrupt entity nobody should assume is // gone. func (e *EntityServer) resolve( - ctx context.Context, index entity.Attr, ids []entity.Id, entities []*entity.Entity, @@ -1096,7 +1075,6 @@ func (e *EntityServer) resolve( total, revision int64, ) *resolvedPage { resolved := make([]*entity.Entity, 0, len(entities)) - cleanups := 0 for i, ent := range entities { if ent != nil { @@ -1107,9 +1085,9 @@ func (e *EntityServer) resolve( if undecodable[ids[i]] { e.Log.Error("entity in index cannot be decoded, skipping", "id", ids[i], "index", index) - } else if cleanups < maxIndexCleanupsPerList { - e.cleanupOrphan(ctx, index, ids[i]) - cleanups++ + } else { + e.Log.Debug("entity in index but not in store, skipping", + "id", ids[i], "index", index) } if total > 0 { diff --git a/servers/entityserver/entityserver_test.go b/servers/entityserver/entityserver_test.go index 0180e26cf..299ff13e5 100644 --- a/servers/entityserver/entityserver_test.go +++ b/servers/entityserver/entityserver_test.go @@ -1247,10 +1247,12 @@ func TestEntityServer_ListPage(t *testing.T) { } -func TestEntityServer_ListsCleanOrphansOnce(t *testing.T) { +func TestEntityServer_OrphanReadsStayQuietAndPure(t *testing.T) { ctx := t.Context() client, prefix := setupTestEtcd(t) - store, err := entity.NewEtcdStore(ctx, slog.Default(), client, prefix) + var logs bytes.Buffer + log := slog.New(slog.NewTextHandler(&logs, &slog.HandlerOptions{Level: slog.LevelDebug})) + store, err := entity.NewEtcdStore(ctx, log, client, prefix) require.NoError(t, err) index := entity.String(entity.Id("test/kind"), "widget") _, err = store.CreateEntity(ctx, entity.New( @@ -1260,91 +1262,33 @@ func TestEntityServer_ListsCleanOrphansOnce(t *testing.T) { require.NoError(t, err) live, err := store.CreateEntity(ctx, entity.New(entity.Ident, "live", index)) require.NoError(t, err) + server, err := NewEntityServer(log, store) + require.NoError(t, err) + sc := v1alpha.EntityAccessClient{Client: rpc.LocalClient(v1alpha.AdaptEntityAccess(server))} + id := entity.Id("orphan") indexPrefix, err := store.IndexPrefix(ctx, index) require.NoError(t, err) - - var logs bytes.Buffer - server, err := NewEntityServer(slog.New(slog.NewTextHandler(&logs, nil)), store) + key := indexPrefix + base58.Encode([]byte(id)) + _, err = client.Put(ctx, key, id.String()) require.NoError(t, err) - sc := v1alpha.EntityAccessClient{Client: rpc.LocalClient(v1alpha.AdaptEntityAccess(server))} + logs.Reset() - for _, paged := range []bool{false, true} { - id := entity.Id(fmt.Sprintf("missing-%t", paged)) - key := indexPrefix + base58.Encode([]byte(id)) - entryKey := key - if paged { - entryKey += "/old-session" - } - _, err := client.Put(ctx, entryKey, id.String()) + for range 2 { + list, err := sc.List(ctx, index) require.NoError(t, err) - for range 2 { - if paged { - page, err := sc.ListPage(ctx, index, "", 10) - require.NoError(t, err) - require.Len(t, page.Values(), 1) - require.Equal(t, int64(1), page.Total()) - require.Equal(t, live.Id().String(), page.Values()[0].Id()) - } else { - list, err := sc.List(ctx, index) - require.NoError(t, err) - require.Len(t, list.Values(), 1) - require.Equal(t, live.Id().String(), list.Values()[0].Id()) - } - } - response, err := client.Get(ctx, key, clientv3.WithPrefix()) + require.Len(t, list.Values(), 1) + require.Equal(t, live.Id().String(), list.Values()[0].Id()) + page, err := sc.ListPage(ctx, index, "", 10) require.NoError(t, err) - require.Empty(t, response.Kvs) - } - require.Equal(t, 2, strings.Count(logs.String(), "cleaned up orphaned index entry")) - require.NotContains(t, logs.String(), "entity in index but not in store") -} - -func TestEntityServer_ListCleanupIsBounded(t *testing.T) { - for _, paged := range []bool{false, true} { - t.Run(fmt.Sprintf("paged=%t", paged), func(t *testing.T) { - ctx := t.Context() - client, prefix := setupTestEtcd(t) - store, err := entity.NewEtcdStore(ctx, slog.Default(), client, prefix) - require.NoError(t, err) - index := entity.String(entity.Id("test/kind"), "widget") - _, err = store.CreateEntity(ctx, entity.New( - entity.Ident, "test/kind", entity.Doc, "indexed kind", - entity.Cardinality, entity.CardinalityOne, entity.Type, entity.TypeStr, entity.Index, true, - )) - require.NoError(t, err) - live, err := store.CreateEntity(ctx, entity.New(entity.Ident, "live", index)) - require.NoError(t, err) - indexPrefix, err := store.IndexPrefix(ctx, index) - require.NoError(t, err) - for i := range maxIndexCleanupsPerList + 2 { - id := entity.Id(fmt.Sprintf("orphan-%02d", i)) - _, err := client.Put(ctx, indexPrefix+base58.Encode([]byte(id)), id.String()) - require.NoError(t, err) - } - - var logs bytes.Buffer - server, err := NewEntityServer(slog.New(slog.NewTextHandler(&logs, nil)), store) - require.NoError(t, err) - sc := v1alpha.EntityAccessClient{Client: rpc.LocalClient(v1alpha.AdaptEntityAccess(server))} - for pass, remaining := range []int64{2, 0, 0} { - if paged { - page, err := sc.ListPage(ctx, index, "", 100) - require.NoError(t, err) - require.Len(t, page.Values(), 1) - require.Equal(t, live.Id().String(), page.Values()[0].Id()) - } else { - list, err := sc.List(ctx, index) - require.NoError(t, err) - require.Len(t, list.Values(), 1) - require.Equal(t, live.Id().String(), list.Values()[0].Id()) - } - entries, err := client.Get(ctx, indexPrefix, clientv3.WithPrefix()) - require.NoError(t, err) - require.Equal(t, remaining+1, entries.Count, "pass %d", pass) - } - require.Equal(t, maxIndexCleanupsPerList+2, strings.Count(logs.String(), "cleaned up orphaned index entry")) - }) + require.Len(t, page.Values(), 1) + require.Equal(t, int64(1), page.Total()) + entry, err := client.Get(ctx, key) + require.NoError(t, err) + require.Len(t, entry.Kvs, 1, "listing must leave GC to repair the orphan") } + require.Equal(t, 4, strings.Count(logs.String(), "level=DEBUG msg=\"entity in index but not in store, skipping\"")) + require.NotContains(t, logs.String(), "level=WARN") + require.NotContains(t, logs.String(), "level=ERROR") } func TestEntityServer_OrphanCleanupWatchDelete(t *testing.T) { @@ -1387,6 +1331,9 @@ func TestEntityServer_OrphanCleanupWatchDelete(t *testing.T) { result, err := sc.List(ctx, index) require.NoError(t, err) require.Empty(t, result.Values()) + stats, err := store.CleanupStaleCollectionEntries(ctx, log, entity.CleanupOptions{}) + require.NoError(t, err) + require.EqualValues(t, 1, stats.StaleEntriesRemoved) select { case op := <-deletes: require.Equal(t, id.String(), op.EntityId()) @@ -1400,6 +1347,5 @@ func TestEntityServer_OrphanCleanupWatchDelete(t *testing.T) { case <-time.After(5 * time.Second): t.Fatal("timed out waiting for watch to stop") } - require.Equal(t, 1, strings.Count(logs.String(), "level=WARN msg=\"cleaned up orphaned index entry\"")) require.NotContains(t, logs.String(), "level=ERROR") }