diff --git a/pkg/entity/cleanup.go b/pkg/entity/cleanup.go index 0b1392657..f24dfe712 100644 --- a/pkg/entity/cleanup.go +++ b/pkg/entity/cleanup.go @@ -281,7 +281,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/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..3251a3b9d 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,8 @@ 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.Log.Debug("entity in index but not in store, skipping", + "id", ids[i], "index", index) continue } @@ -1086,7 +1086,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", + e.Log.Debug("entity in index but not in store, skipping", "id", ids[i], "index", index) } diff --git a/servers/entityserver/entityserver_test.go b/servers/entityserver/entityserver_test.go index 067cbb6ed..299ff13e5 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,106 @@ func TestEntityServer_ListPage(t *testing.T) { }) } + +func TestEntityServer_OrphanReadsStayQuietAndPure(t *testing.T) { + ctx := t.Context() + client, prefix := setupTestEtcd(t) + 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( + 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) + 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) + key := indexPrefix + base58.Encode([]byte(id)) + _, err = client.Put(ctx, key, id.String()) + require.NoError(t, err) + logs.Reset() + + for range 2 { + 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()) + page, err := sc.ListPage(ctx, index, "", 10) + require.NoError(t, err) + 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) { + 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()) + 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()) + 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.NotContains(t, logs.String(), "level=ERROR") +}