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
2 changes: 1 addition & 1 deletion pkg/entity/cleanup.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand Down
16 changes: 5 additions & 11 deletions pkg/entity/store.go
Original file line number Diff line number Diff line change
Expand Up @@ -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)
}
Expand All @@ -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
}

Expand Down Expand Up @@ -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
}
Expand All @@ -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.
//
Expand All @@ -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{}
Expand Down Expand Up @@ -615,17 +612,14 @@ 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
}

var entity Entity
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
Expand Down
10 changes: 5 additions & 5 deletions pkg/entity/store_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand All @@ -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])
Expand All @@ -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])
Expand All @@ -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])
Expand All @@ -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)

Expand Down
20 changes: 10 additions & 10 deletions servers/entityserver/entityserver.go
Original file line number Diff line number Diff line change
Expand Up @@ -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())
Expand Down Expand Up @@ -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
}

Expand Down Expand Up @@ -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)
}

Expand Down
106 changes: 106 additions & 0 deletions servers/entityserver/entityserver_test.go
Original file line number Diff line number Diff line change
@@ -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"
Expand Down Expand Up @@ -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")
}
Loading