From 230efa19158dc5ed027fc321666cd3f93e36e86d Mon Sep 17 00:00:00 2001 From: jungooji <46273764+zzzz465@users.noreply.github.com> Date: Wed, 1 Jul 2026 16:04:47 +0900 Subject: [PATCH 1/5] feat: filter entry lists by managed hints Signed-off-by: jungooji <46273764+zzzz465@users.noreply.github.com> --- api/v1alpha1/controllermanagerconfig_types.go | 7 ++ cmd/main.go | 7 +- pkg/spireapi/entryapi.go | 41 ++++++- pkg/spireapi/entryapi_test.go | 34 ++++- pkg/spireentry/entry_list_cache.go | 8 +- pkg/spireentry/reconciler.go | 96 ++++++++++++--- pkg/spireentry/reconciler_cache_test.go | 116 +++++++++++++----- 7 files changed, 251 insertions(+), 58 deletions(-) diff --git a/api/v1alpha1/controllermanagerconfig_types.go b/api/v1alpha1/controllermanagerconfig_types.go index 612bd683..17e47a29 100644 --- a/api/v1alpha1/controllermanagerconfig_types.go +++ b/api/v1alpha1/controllermanagerconfig_types.go @@ -193,6 +193,13 @@ type ControllerManagerConfigurationSpec struct { // used when enableEntryListCache is true. Defaults to 1h if unset or zero. // +optional EntryListCacheReloadInterval time.Duration `json:"entryListCacheReloadInterval,omitempty"` + + // EnableEntryListHintFilter limits SPIRE ListEntries calls to hints used by + // managed ClusterSPIFFEID and ClusterStaticEntry resources. If any managed + // resource has an empty hint, the controller falls back to an unfiltered list. + // Defaults to false. + // +optional + EnableEntryListHintFilter bool `json:"enableEntryListHintFilter,omitempty"` } // ReconcileConfig configuration used to enable/disable syncing various types diff --git a/cmd/main.go b/cmd/main.go index 87fa9c6e..e07843fd 100644 --- a/cmd/main.go +++ b/cmd/main.go @@ -259,7 +259,8 @@ func parseConfig() (Config, error) { "entryIDPrefix", retval.ctrlConfig.EntryIDPrefix, "entryIDPrefixCleanup", printCleanup, "enableEntryListCache", retval.ctrlConfig.EnableEntryListCache, - "entryListCacheReloadInterval", retval.ctrlConfig.EntryListCacheReloadInterval) + "entryListCacheReloadInterval", retval.ctrlConfig.EntryListCacheReloadInterval, + "enableEntryListHintFilter", retval.ctrlConfig.EnableEntryListHintFilter) switch { case retval.ctrlConfig.TrustDomain == "": @@ -271,6 +272,8 @@ func parseConfig() (Config, error) { return retval, errors.New("validating webhook configuration name is required configuration") case retval.ctrlConfig.EnableEntryListCache && retval.ctrlConfig.EntryIDPrefix == "": return retval, errors.New("enableEntryListCache requires entryIDPrefix to be set") + case retval.ctrlConfig.EnableEntryListHintFilter && retval.ctrlConfig.EntryIDPrefix == "": + return retval, errors.New("enableEntryListHintFilter requires entryIDPrefix to be set") case retval.ctrlConfig.ControllerManagerConfigurationSpec.Webhook.CertDir != "": setupLog.Info("certDir configuration is ignored", "certDir", retval.ctrlConfig.ControllerManagerConfigurationSpec.Webhook.CertDir) } @@ -393,6 +396,7 @@ func run(mainConfig Config) (err error) { EnableEntryListCache: mainConfig.ctrlConfig.EnableEntryListCache, EntryListCacheReloadInterval: mainConfig.ctrlConfig.EntryListCacheReloadInterval, + EnableEntryListHintFilter: mainConfig.ctrlConfig.EnableEntryListHintFilter, }) } @@ -556,6 +560,7 @@ func staticRun(mainConfig Config) (err error) { EnableEntryListCache: mainConfig.ctrlConfig.EnableEntryListCache, EntryListCacheReloadInterval: mainConfig.ctrlConfig.EntryListCacheReloadInterval, + EnableEntryListHintFilter: mainConfig.ctrlConfig.EnableEntryListHintFilter, }) go func() { err := entryReconciler.Run(ctx) diff --git a/pkg/spireapi/entryapi.go b/pkg/spireapi/entryapi.go index d5d7ee2f..709e60b9 100644 --- a/pkg/spireapi/entryapi.go +++ b/pkg/spireapi/entryapi.go @@ -24,6 +24,7 @@ import ( "github.com/spiffe/spire-api-sdk/proto/spire/api/types" "google.golang.org/grpc" "google.golang.org/grpc/codes" + "google.golang.org/protobuf/types/known/wrapperspb" "sigs.k8s.io/controller-runtime/pkg/log" ) @@ -41,7 +42,7 @@ const ( type Field string type EntryClient interface { - ListEntries(ctx context.Context) ([]Entry, error) + ListEntries(ctx context.Context, hints ...string) ([]Entry, error) CreateEntries(ctx context.Context, entries []Entry) ([]Status, error) UpdateEntries(ctx context.Context, entries []Entry) ([]Status, error) DeleteEntries(ctx context.Context, entryIDs []string) ([]Status, error) @@ -56,11 +57,31 @@ type entryClient struct { api entryv1.EntryClient } -func (c entryClient) ListEntries(ctx context.Context) ([]Entry, error) { +func (c entryClient) ListEntries(ctx context.Context, hints ...string) ([]Entry, error) { + filterHints := uniqueNonEmptyStrings(hints) + if len(filterHints) == 0 { + return c.listEntries(ctx, nil) + } + + var entries []Entry + for _, hint := range filterHints { + hintEntries, err := c.listEntries(ctx, &entryv1.ListEntriesRequest_Filter{ + ByHint: wrapperspb.String(hint), + }) + if err != nil { + return nil, err + } + entries = append(entries, hintEntries...) + } + return entries, nil +} + +func (c entryClient) listEntries(ctx context.Context, filter *entryv1.ListEntriesRequest_Filter) ([]Entry, error) { var entries []*types.Entry var pageToken string for { resp, err := c.api.ListEntries(ctx, &entryv1.ListEntriesRequest{ + Filter: filter, PageToken: pageToken, PageSize: entryListPageSize, }) @@ -76,6 +97,22 @@ func (c entryClient) ListEntries(ctx context.Context) ([]Entry, error) { return entriesFromAPI(entries) } +func uniqueNonEmptyStrings(values []string) []string { + seen := make(map[string]struct{}, len(values)) + out := make([]string, 0, len(values)) + for _, value := range values { + if value == "" { + continue + } + if _, ok := seen[value]; ok { + continue + } + seen[value] = struct{}{} + out = append(out, value) + } + return out +} + func (c entryClient) GetUnsupportedFields(ctx context.Context, td string, entryIDPrefix string) (map[Field]struct{}, error) { entryID := "" if entryIDPrefix != "" { diff --git a/pkg/spireapi/entryapi_test.go b/pkg/spireapi/entryapi_test.go index c900ab7f..7c1e251d 100644 --- a/pkg/spireapi/entryapi_test.go +++ b/pkg/spireapi/entryapi_test.go @@ -92,6 +92,26 @@ func TestEntryAPIListEntries(t *testing.T) { } } +func TestEntryAPIListEntriesByHint(t *testing.T) { + server, client := startEntryAPIServer(t) + withHint := func(entry Entry, hint string) Entry { + entry.Hint = hint + return entry + } + server.setEntries(t, + withHint(entry1, "cluster-a"), + withHint(entry2, "cluster-b"), + withHint(entry3, "cluster-a"), + ) + + actualEntries, err := client.ListEntries(ctx, "cluster-a", "cluster-a", "") + require.NoError(t, err) + assert.ElementsMatch(t, []Entry{ + withHint(entry1, "cluster-a"), + withHint(entry3, "cluster-a"), + }, actualEntries) +} + func TestCreateEntries(t *testing.T) { server, client := startEntryAPIServer(t) @@ -408,8 +428,18 @@ func (s *entryServer) ListEntries(_ context.Context, req *entryv1.ListEntriesReq s.mtx.RLock() defer s.mtx.RUnlock() - start, end, more := listBounds(req.PageToken, int(req.PageSize), len(s.entries), func(i int) string { return s.entries[i].Id }) - for _, entry := range s.entries[start:end] { + entries := s.entries + if hint := req.GetFilter().GetByHint().GetValue(); hint != "" { + entries = nil + for _, entry := range s.entries { + if entry.Hint == hint { + entries = append(entries, entry) + } + } + } + + start, end, more := listBounds(req.PageToken, int(req.PageSize), len(entries), func(i int) string { return entries[i].Id }) + for _, entry := range entries[start:end] { resp.Entries = append(resp.Entries, entry) if more { resp.NextPageToken = entry.Id diff --git a/pkg/spireentry/entry_list_cache.go b/pkg/spireentry/entry_list_cache.go index aa8f3445..2e2db739 100644 --- a/pkg/spireentry/entry_list_cache.go +++ b/pkg/spireentry/entry_list_cache.go @@ -35,11 +35,12 @@ type entryListCache struct { reloadAfter time.Duration entries map[string]spireapi.Entry nextReload time.Time + filterKey string } // fresh reports whether the cache can be served without listing from the server. -func (c *entryListCache) fresh() bool { - return c.entries != nil && time.Now().Before(c.nextReload) +func (c *entryListCache) fresh(filterKey string) bool { + return c.entries != nil && c.filterKey == filterKey && time.Now().Before(c.nextReload) } // snapshot returns the cached entries as a slice. @@ -52,11 +53,12 @@ func (c *entryListCache) snapshot() []spireapi.Entry { } // replace rebuilds the cache from a fresh server list and arms the next reload. -func (c *entryListCache) replace(entries []spireapi.Entry) { +func (c *entryListCache) replace(filterKey string, entries []spireapi.Entry) { c.entries = make(map[string]spireapi.Entry, len(entries)) for _, entry := range entries { c.entries[entry.ID] = entry } + c.filterKey = filterKey c.nextReload = time.Now().Add(c.reloadAfter) } diff --git a/pkg/spireentry/reconciler.go b/pkg/spireentry/reconciler.go index 07f9366d..cc6f28f0 100644 --- a/pkg/spireentry/reconciler.go +++ b/pkg/spireentry/reconciler.go @@ -102,6 +102,12 @@ type ReconcilerConfig struct { // reloaded from the SPIRE server. Only used when EnableEntryListCache is // true. If zero, defaults to defaultEntryListCacheReloadInterval. EntryListCacheReloadInterval time.Duration + + // EnableEntryListHintFilter uses the hints from managed ClusterSPIFFEID and + // ClusterStaticEntry resources to narrow ListEntries calls at the SPIRE + // server. If any managed resource has an empty hint, the reconciler falls + // back to an unfiltered list to preserve compatibility. + EnableEntryListHintFilter bool } const ( @@ -190,19 +196,7 @@ func (r *entryReconciler) reconcile(ctx context.Context) { } unsupportedFields := r.unsupportedFields - // Load current entries from SPIRE server. - currentEntries, deleteOnlyEntries, err := r.listEntries(ctx) - if err != nil { - log.Error(err, "Failed to list SPIRE entries") - return - } - - // Populate the existing state - state := make(entriesState) - for _, entry := range currentEntries { - state.AddCurrent(entry) - } - + var err error clusterStaticEntries := []*ClusterStaticEntry{} if r.config.Reconcile.ClusterStaticEntries { // Load and add entry state for ClusterStaticEntries @@ -211,7 +205,6 @@ func (r *entryReconciler) reconcile(ctx context.Context) { log.Error(err, "Failed to list ClusterStaticEntries") return } - r.addClusterStaticEntryEntriesState(ctx, state, clusterStaticEntries) } clusterSPIFFEIDs := []*ClusterSPIFFEID{} @@ -222,7 +215,28 @@ func (r *entryReconciler) reconcile(ctx context.Context) { log.Error(err, "Failed to list ClusterSPIFFEIDs") return } + } + + entryListHints := r.entryListHints(clusterSPIFFEIDs, clusterStaticEntries) + + // Load current entries from SPIRE server. + currentEntries, deleteOnlyEntries, err := r.listEntries(ctx, entryListHints) + if err != nil { + log.Error(err, "Failed to list SPIRE entries") + return + } + + // Populate the existing state + state := make(entriesState) + for _, entry := range currentEntries { + state.AddCurrent(entry) + } + if r.config.Reconcile.ClusterStaticEntries { + r.addClusterStaticEntryEntriesState(ctx, state, clusterStaticEntries) + } + + if r.config.Reconcile.ClusterSPIFFEIDs { // Pre-load all nodes into a map to avoid per-pod Get() calls, // which each incur a deep copy and mutex lock on the informer cache. nodeMap, err := r.buildNodeMap(ctx) @@ -388,25 +402,71 @@ func (r *entryReconciler) incCounter(name string) { // in-memory cache while fresh, avoiding a full ListEntries RPC to the SPIRE // server on every reconcile; otherwise it lists from the server and (when the // cache is enabled) repopulates it. -func (r *entryReconciler) listEntries(ctx context.Context) ([]spireapi.Entry, []spireapi.Entry, error) { - if r.entryCache != nil && r.entryCache.fresh() { +func (r *entryReconciler) listEntries(ctx context.Context, hints []string) ([]spireapi.Entry, []spireapi.Entry, error) { + filterKey := entryListFilterKey(hints) + if r.entryCache != nil && r.entryCache.fresh(filterKey) { r.incCounter(metrics.EntryListCacheHits) currentEntries, deleteOnlyEntries := r.splitEntries(r.entryCache.snapshot()) return currentEntries, deleteOnlyEntries, nil } r.incCounter(metrics.EntryListServerCalls) - entries, err := r.config.EntryClient.ListEntries(ctx) + entries, err := r.config.EntryClient.ListEntries(ctx, hints...) if err != nil { return nil, nil, err } currentEntries, deleteOnlyEntries := r.splitEntries(entries) if r.entryCache != nil { - r.entryCache.replace(append(currentEntries, deleteOnlyEntries...)) + r.entryCache.replace(filterKey, append(currentEntries, deleteOnlyEntries...)) } return currentEntries, deleteOnlyEntries, nil } +func (r *entryReconciler) entryListHints(clusterSPIFFEIDs []*ClusterSPIFFEID, clusterStaticEntries []*ClusterStaticEntry) []string { + if !r.config.EnableEntryListHintFilter { + return nil + } + + var hints []string + for _, clusterSPIFFEID := range clusterSPIFFEIDs { + if clusterSPIFFEID.Spec.Hint == "" { + return nil + } + hints = append(hints, clusterSPIFFEID.Spec.Hint) + } + for _, clusterStaticEntry := range clusterStaticEntries { + if clusterStaticEntry.Spec.Hint == "" { + return nil + } + hints = append(hints, clusterStaticEntry.Spec.Hint) + } + return uniqueStrings(hints) +} + +func uniqueStrings(values []string) []string { + seen := make(map[string]struct{}, len(values)) + out := make([]string, 0, len(values)) + for _, value := range values { + if value == "" { + continue + } + if _, ok := seen[value]; ok { + continue + } + seen[value] = struct{}{} + out = append(out, value) + } + sort.Strings(out) + return out +} + +func entryListFilterKey(hints []string) string { + if len(hints) == 0 { + return "" + } + return strings.Join(uniqueStrings(hints), "\x00") +} + // splitEntries partitions entries into the set to process and the set marked // only for cleanup deletion, per shouldProcessOrDeleteEntryID. It is recomputed // on every call (including cached reads) so the partition stays consistent with diff --git a/pkg/spireentry/reconciler_cache_test.go b/pkg/spireentry/reconciler_cache_test.go index 20212823..3937f787 100644 --- a/pkg/spireentry/reconciler_cache_test.go +++ b/pkg/spireentry/reconciler_cache_test.go @@ -27,6 +27,7 @@ import ( "github.com/prometheus/client_golang/prometheus" dto "github.com/prometheus/client_model/go" "github.com/spiffe/go-spiffe/v2/spiffeid" + spirev1alpha1 "github.com/spiffe/spire-controller-manager/api/v1alpha1" "github.com/spiffe/spire-controller-manager/pkg/metrics" "github.com/spiffe/spire-controller-manager/pkg/spireapi" "github.com/stretchr/testify/require" @@ -41,6 +42,7 @@ import ( // calls and returns scriptable per-operation status codes (zero value = OK). type fakeEntryClient struct { listCalls int + listHints [][]string list []spireapi.Entry createCode codes.Code @@ -56,8 +58,9 @@ type fakeEntryClient struct { deleted []string } -func (f *fakeEntryClient) ListEntries(context.Context) ([]spireapi.Entry, error) { +func (f *fakeEntryClient) ListEntries(_ context.Context, hints ...string) ([]spireapi.Entry, error) { f.listCalls++ + f.listHints = append(f.listHints, append([]string(nil), hints...)) out := make([]spireapi.Entry, len(f.list)) copy(out, f.list) return out, nil @@ -159,11 +162,11 @@ func TestEntryListCache(t *testing.T) { r := newCacheReconciler(fake, cacheCfg()) ctx := testCtx() - cur1, _, err := r.listEntries(ctx) + cur1, _, err := r.listEntries(ctx, nil) require.NoError(t, err) require.Len(t, cur1, 2) - cur2, _, err := r.listEntries(ctx) + cur2, _, err := r.listEntries(ctx, nil) require.NoError(t, err) require.Len(t, cur2, 2) @@ -177,8 +180,8 @@ func TestEntryListCache(t *testing.T) { r := newCacheReconciler(fake, cfg) require.Nil(t, r.entryCache) - _, _, _ = r.listEntries(testCtx()) - _, _, _ = r.listEntries(testCtx()) + _, _, _ = r.listEntries(testCtx(), nil) + _, _, _ = r.listEntries(testCtx(), nil) require.Equal(t, 2, fake.listCalls) }) @@ -187,7 +190,7 @@ func TestEntryListCache(t *testing.T) { r := newCacheReconciler(fake, cacheCfg()) require.Nil(t, r.entryCache.entries, "cache starts empty") - _, _, err := r.listEntries(testCtx()) + _, _, err := r.listEntries(testCtx(), nil) require.NoError(t, err) require.Equal(t, 1, fake.listCalls) require.NotNil(t, r.entryCache.entries) @@ -198,12 +201,12 @@ func TestEntryListCache(t *testing.T) { r := newCacheReconciler(fake, cacheCfg()) // reloadAfter = 1m ctx := testCtx() - _, _, _ = r.listEntries(ctx) - _, _, _ = r.listEntries(ctx) + _, _, _ = r.listEntries(ctx, nil) + _, _, _ = r.listEntries(ctx, nil) require.Equal(t, 1, fake.listCalls) time.Sleep(2 * time.Minute) // fake clock advances past nextReload - _, _, _ = r.listEntries(ctx) + _, _, _ = r.listEntries(ctx, nil) require.Equal(t, 2, fake.listCalls, "must re-list after reload interval elapses") }) @@ -212,10 +215,10 @@ func TestEntryListCache(t *testing.T) { r := newCacheReconciler(fake, cacheCfg()) ctx := testCtx() - _, _, _ = r.listEntries(ctx) + _, _, _ = r.listEntries(ctx, nil) r.createEntries(ctx, []declaredEntry{declared("test.new")}) - cur, _, err := r.listEntries(ctx) + cur, _, err := r.listEntries(ctx, nil) require.NoError(t, err) require.Equal(t, 1, fake.listCalls, "create result must be reflected in cache, not re-listed") require.ElementsMatch(t, []string{"test.a", "test.new"}, ids(cur)) @@ -226,13 +229,13 @@ func TestEntryListCache(t *testing.T) { r := newCacheReconciler(fake, cacheCfg()) ctx := testCtx() - _, _, _ = r.listEntries(ctx) + _, _, _ = r.listEntries(ctx, nil) r.createEntries(ctx, []declaredEntry{declared("test.new")}) - require.False(t, r.entryCache.fresh(), "drift must invalidate the cache") + require.False(t, r.entryCache.fresh(""), "drift must invalidate the cache") - _, _, _ = r.listEntries(ctx) + _, _, _ = r.listEntries(ctx, nil) require.Equal(t, 2, fake.listCalls) - require.True(t, r.entryCache.fresh(), "cache fresh again after reload") + require.True(t, r.entryCache.fresh(""), "cache fresh again after reload") }) bubble(t, "update NotFound drops entry and triggers resync", func(t *testing.T) { @@ -240,12 +243,12 @@ func TestEntryListCache(t *testing.T) { r := newCacheReconciler(fake, cacheCfg()) ctx := testCtx() - _, _, _ = r.listEntries(ctx) + _, _, _ = r.listEntries(ctx, nil) r.updateEntries(ctx, []declaredEntry{declared("test.a")}) - require.False(t, r.entryCache.fresh()) + require.False(t, r.entryCache.fresh("")) require.NotContains(t, r.entryCache.entries, "test.a") - _, _, _ = r.listEntries(ctx) + _, _, _ = r.listEntries(ctx, nil) require.Equal(t, 2, fake.listCalls) }) @@ -254,10 +257,10 @@ func TestEntryListCache(t *testing.T) { r := newCacheReconciler(fake, cacheCfg()) ctx := testCtx() - _, _, _ = r.listEntries(ctx) + _, _, _ = r.listEntries(ctx, nil) r.deleteEntries(ctx, []spireapi.Entry{testEntry("test.a")}) - cur, _, err := r.listEntries(ctx) + cur, _, err := r.listEntries(ctx, nil) require.NoError(t, err) require.Equal(t, 1, fake.listCalls) require.ElementsMatch(t, []string{"test.b"}, ids(cur)) @@ -268,13 +271,13 @@ func TestEntryListCache(t *testing.T) { r := newCacheReconciler(fake, cacheCfg()) ctx := testCtx() - _, _, _ = r.listEntries(ctx) - require.True(t, r.entryCache.fresh()) + _, _, _ = r.listEntries(ctx, nil) + require.True(t, r.entryCache.fresh("")) r.createEntries(ctx, []declaredEntry{declared("test.new")}) - require.False(t, r.entryCache.fresh(), "RPC error must invalidate the cache") + require.False(t, r.entryCache.fresh(""), "RPC error must invalidate the cache") - _, _, _ = r.listEntries(ctx) + _, _, _ = r.listEntries(ctx, nil) require.Equal(t, 2, fake.listCalls, "next reconcile must reload from server") }) @@ -283,9 +286,9 @@ func TestEntryListCache(t *testing.T) { r := newCacheReconciler(fake, cacheCfg()) ctx := testCtx() - _, _, _ = r.listEntries(ctx) + _, _, _ = r.listEntries(ctx, nil) r.deleteEntries(ctx, []spireapi.Entry{testEntry("test.a")}) - require.False(t, r.entryCache.fresh()) + require.False(t, r.entryCache.fresh("")) require.NotContains(t, r.entryCache.entries, "test.a") require.Len(t, fake.deleted, 1, "must not re-issue the delete") }) @@ -298,10 +301,10 @@ func TestEntryListCache(t *testing.T) { r := newCacheReconciler(fake, cfg) ctx := testCtx() - _, del1, _ := r.listEntries(ctx) + _, del1, _ := r.listEntries(ctx, nil) require.ElementsMatch(t, []string{"old.x"}, ids(del1)) - cur2, del2, err := r.listEntries(ctx) + cur2, del2, err := r.listEntries(ctx, nil) require.NoError(t, err) require.Equal(t, 1, fake.listCalls) require.ElementsMatch(t, []string{"test.a"}, ids(cur2)) @@ -332,13 +335,13 @@ func TestEntryListMetrics(t *testing.T) { ctx := testCtx() // First reconcile: cold cache -> server list. - _, _, err := r.listEntries(ctx) + _, _, err := r.listEntries(ctx, nil) require.NoError(t, err) require.Equal(t, 1.0, counterValue(t, r.promCounter[metrics.EntryListServerCalls])) require.Equal(t, 0.0, counterValue(t, r.promCounter[metrics.EntryListCacheHits])) // Second reconcile: served from the fresh cache -> cache hit. - _, _, err = r.listEntries(ctx) + _, _, err = r.listEntries(ctx, nil) require.NoError(t, err) require.Equal(t, 1.0, counterValue(t, r.promCounter[metrics.EntryListServerCalls])) require.Equal(t, 1.0, counterValue(t, r.promCounter[metrics.EntryListCacheHits])) @@ -352,13 +355,62 @@ func TestEntryListMetrics(t *testing.T) { r.promCounter = newCounters() ctx := testCtx() - _, _, _ = r.listEntries(ctx) - _, _, _ = r.listEntries(ctx) + _, _, _ = r.listEntries(ctx, nil) + _, _, _ = r.listEntries(ctx, nil) require.Equal(t, 2.0, counterValue(t, r.promCounter[metrics.EntryListServerCalls])) require.Equal(t, 0.0, counterValue(t, r.promCounter[metrics.EntryListCacheHits])) }) } +func TestEntryListHintFilter(t *testing.T) { + t.Run("collects unique hints from managed resources", func(t *testing.T) { + r := newCacheReconciler(&fakeEntryClient{}, ReconcilerConfig{EnableEntryListHintFilter: true}) + hints := r.entryListHints( + []*ClusterSPIFFEID{ + {ClusterSPIFFEID: spirev1alpha1.ClusterSPIFFEID{Spec: spirev1alpha1.ClusterSPIFFEIDSpec{Hint: "cluster-a"}}}, + {ClusterSPIFFEID: spirev1alpha1.ClusterSPIFFEID{Spec: spirev1alpha1.ClusterSPIFFEIDSpec{Hint: "cluster-b"}}}, + {ClusterSPIFFEID: spirev1alpha1.ClusterSPIFFEID{Spec: spirev1alpha1.ClusterSPIFFEIDSpec{Hint: "cluster-a"}}}, + }, + []*ClusterStaticEntry{ + {ClusterStaticEntry: spirev1alpha1.ClusterStaticEntry{Spec: spirev1alpha1.ClusterStaticEntrySpec{Hint: "static"}}}, + }, + ) + + require.Equal(t, []string{"cluster-a", "cluster-b", "static"}, hints) + }) + + t.Run("falls back to full list when a managed resource has no hint", func(t *testing.T) { + r := newCacheReconciler(&fakeEntryClient{}, ReconcilerConfig{EnableEntryListHintFilter: true}) + hints := r.entryListHints( + []*ClusterSPIFFEID{ + {ClusterSPIFFEID: spirev1alpha1.ClusterSPIFFEID{Spec: spirev1alpha1.ClusterSPIFFEIDSpec{Hint: "cluster-a"}}}, + {ClusterSPIFFEID: spirev1alpha1.ClusterSPIFFEID{}}, + }, + nil, + ) + + require.Nil(t, hints) + }) + + t.Run("passes hint filter to server list and reloads when hints change", func(t *testing.T) { + fake := &fakeEntryClient{list: []spireapi.Entry{testEntry("test.a")}} + r := newCacheReconciler(fake, cacheCfg()) + ctx := testCtx() + + _, _, err := r.listEntries(ctx, []string{"cluster-a"}) + require.NoError(t, err) + require.Equal(t, [][]string{{"cluster-a"}}, fake.listHints) + + _, _, err = r.listEntries(ctx, []string{"cluster-a"}) + require.NoError(t, err) + require.Equal(t, 1, fake.listCalls) + + _, _, err = r.listEntries(ctx, []string{"cluster-b"}) + require.NoError(t, err) + require.Equal(t, [][]string{{"cluster-a"}, {"cluster-b"}}, fake.listHints) + }) +} + func ids(entries []spireapi.Entry) []string { out := make([]string, 0, len(entries)) for _, e := range entries { From 80b57a611cea14bf092de37764beb61dd3aafbe2 Mon Sep 17 00:00:00 2001 From: jungooji <46273764+zzzz465@users.noreply.github.com> Date: Wed, 1 Jul 2026 16:12:46 +0900 Subject: [PATCH 2/5] perf: parallelize hint entry listing Signed-off-by: jungooji <46273764+zzzz465@users.noreply.github.com> --- go.mod | 1 + go.sum | 2 ++ pkg/spireapi/entryapi.go | 65 +++++++++++++++++++++-------------- pkg/spireapi/entryapi_test.go | 32 +++++++++++++++-- pkg/spireentry/reconciler.go | 16 +++------ 5 files changed, 76 insertions(+), 40 deletions(-) diff --git a/go.mod b/go.mod index 1ad5a8ca..d1106597 100644 --- a/go.mod +++ b/go.mod @@ -11,6 +11,7 @@ require ( github.com/onsi/ginkgo/v2 v2.32.0 github.com/onsi/gomega v1.42.1 github.com/prometheus/client_golang v1.23.2 + github.com/samber/lo v1.53.0 github.com/spiffe/go-spiffe/v2 v2.8.1 github.com/spiffe/spire-api-sdk v1.15.1 github.com/stretchr/testify v1.11.1 diff --git a/go.sum b/go.sum index 5a1a2c92..e222537f 100644 --- a/go.sum +++ b/go.sum @@ -128,6 +128,8 @@ github.com/prometheus/procfs v0.19.2 h1:zUMhqEW66Ex7OXIiDkll3tl9a1ZdilUOd/F6ZXw4 github.com/prometheus/procfs v0.19.2/go.mod h1:M0aotyiemPhBCM0z5w87kL22CxfcH05ZpYlu+b4J7mw= github.com/rogpeppe/go-internal v1.14.1 h1:UQB4HGPB6osV0SQTLymcB4TgvyWu6ZyliaW0tI/otEQ= github.com/rogpeppe/go-internal v1.14.1/go.mod h1:MaRKkUm5W0goXpeCfT7UZI6fk/L7L7so1lCWt35ZSgc= +github.com/samber/lo v1.53.0 h1:t975lj2py4kJPQ6haz1QMgtId2gtmfktACxIXArw3HM= +github.com/samber/lo v1.53.0/go.mod h1:4+MXEGsJzbKGaUEQFKBq2xtfuznW9oz/WrgyzMzRoM0= github.com/spf13/pflag v1.0.10 h1:4EBh2KAYBwaONj6b2Ye1GiHfwjqyROoF4RwYO+vPwFk= github.com/spf13/pflag v1.0.10/go.mod h1:McXfInJRrz4CZXVZOBLb0bTZqETkiAhM9Iw0y3An2Bg= github.com/spiffe/go-spiffe/v2 v2.8.1 h1:eXZMLsu+3MLEPJyGJkolqtVrteZfQdUpOWj6LTiDl/E= diff --git a/pkg/spireapi/entryapi.go b/pkg/spireapi/entryapi.go index 709e60b9..27479da4 100644 --- a/pkg/spireapi/entryapi.go +++ b/pkg/spireapi/entryapi.go @@ -19,7 +19,9 @@ package spireapi import ( "context" "fmt" + "sync" + "github.com/samber/lo" entryv1 "github.com/spiffe/spire-api-sdk/proto/spire/api/server/entry/v1" "github.com/spiffe/spire-api-sdk/proto/spire/api/types" "google.golang.org/grpc" @@ -58,22 +60,51 @@ type entryClient struct { } func (c entryClient) ListEntries(ctx context.Context, hints ...string) ([]Entry, error) { - filterHints := uniqueNonEmptyStrings(hints) + filterHints := lo.Uniq(lo.Filter(hints, func(hint string, _ int) bool { + return hint != "" + })) if len(filterHints) == 0 { return c.listEntries(ctx, nil) } + ctx, cancel := context.WithCancel(ctx) + defer cancel() + + var wg sync.WaitGroup + var errOnce sync.Once + var firstErr error + entriesByHint := make([][]Entry, len(filterHints)) + + for i, hint := range filterHints { + wg.Add(1) + go func(i int, hint string) { + defer wg.Done() + + hintEntries, err := c.listEntries(ctx, &entryv1.ListEntriesRequest_Filter{ + ByHint: wrapperspb.String(hint), + }) + if err != nil { + errOnce.Do(func() { + firstErr = err + cancel() + }) + return + } + entriesByHint[i] = hintEntries + }(i, hint) + } + wg.Wait() + if firstErr != nil { + return nil, firstErr + } + var entries []Entry - for _, hint := range filterHints { - hintEntries, err := c.listEntries(ctx, &entryv1.ListEntriesRequest_Filter{ - ByHint: wrapperspb.String(hint), - }) - if err != nil { - return nil, err - } + for _, hintEntries := range entriesByHint { entries = append(entries, hintEntries...) } - return entries, nil + return lo.UniqBy(entries, func(entry Entry) string { + return entry.ID + }), nil } func (c entryClient) listEntries(ctx context.Context, filter *entryv1.ListEntriesRequest_Filter) ([]Entry, error) { @@ -97,22 +128,6 @@ func (c entryClient) listEntries(ctx context.Context, filter *entryv1.ListEntrie return entriesFromAPI(entries) } -func uniqueNonEmptyStrings(values []string) []string { - seen := make(map[string]struct{}, len(values)) - out := make([]string, 0, len(values)) - for _, value := range values { - if value == "" { - continue - } - if _, ok := seen[value]; ok { - continue - } - seen[value] = struct{}{} - out = append(out, value) - } - return out -} - func (c entryClient) GetUnsupportedFields(ctx context.Context, td string, entryIDPrefix string) (map[Field]struct{}, error) { entryID := "" if entryIDPrefix != "" { diff --git a/pkg/spireapi/entryapi_test.go b/pkg/spireapi/entryapi_test.go index 7c1e251d..a4ad1fd9 100644 --- a/pkg/spireapi/entryapi_test.go +++ b/pkg/spireapi/entryapi_test.go @@ -104,12 +104,26 @@ func TestEntryAPIListEntriesByHint(t *testing.T) { withHint(entry3, "cluster-a"), ) - actualEntries, err := client.ListEntries(ctx, "cluster-a", "cluster-a", "") + actualEntries, err := client.ListEntries(ctx, "cluster-a", "cluster-b", "cluster-a", "") require.NoError(t, err) assert.ElementsMatch(t, []Entry{ withHint(entry1, "cluster-a"), + withHint(entry2, "cluster-b"), withHint(entry3, "cluster-a"), }, actualEntries) + assert.ElementsMatch(t, []string{"cluster-a", "cluster-b"}, server.getListEntryHints()) +} + +func TestEntryAPIListEntriesByHintDeduplicatesEntries(t *testing.T) { + server, client := startEntryAPIServer(t) + entry := entry1 + entry.Hint = "cluster-a" + server.setEntries(t, entry) + server.duplicateListEntries = true + + actualEntries, err := client.ListEntries(ctx, "cluster-a") + require.NoError(t, err) + assert.Equal(t, []Entry{entry}, actualEntries) } func TestCreateEntries(t *testing.T) { @@ -413,8 +427,10 @@ type entryServer struct { mtx sync.RWMutex entries []*apitypes.Entry + hints []string clearUnsupportedFields bool + duplicateListEntries bool listEntriesErr error batchCreateEntriesErr error @@ -425,11 +441,12 @@ type entryServer struct { func (s *entryServer) ListEntries(_ context.Context, req *entryv1.ListEntriesRequest) (*entryv1.ListEntriesResponse, error) { resp := new(entryv1.ListEntriesResponse) - s.mtx.RLock() - defer s.mtx.RUnlock() + s.mtx.Lock() + defer s.mtx.Unlock() entries := s.entries if hint := req.GetFilter().GetByHint().GetValue(); hint != "" { + s.hints = append(s.hints, hint) entries = nil for _, entry := range s.entries { if entry.Hint == hint { @@ -437,6 +454,9 @@ func (s *entryServer) ListEntries(_ context.Context, req *entryv1.ListEntriesReq } } } + if s.duplicateListEntries && len(entries) == 1 { + entries = append(entries, entries[0]) + } start, end, more := listBounds(req.PageToken, int(req.PageSize), len(entries), func(i int) string { return entries[i].Id }) for _, entry := range entries[start:end] { @@ -526,6 +546,12 @@ func (s *entryServer) getEntries(t *testing.T) []Entry { return entries } +func (s *entryServer) getListEntryHints() []string { + s.mtx.Lock() + defer s.mtx.Unlock() + return append([]string(nil), s.hints...) +} + func (s *entryServer) setEntries(t *testing.T, entries ...Entry) { s.clearEntries() for _, entry := range entries { diff --git a/pkg/spireentry/reconciler.go b/pkg/spireentry/reconciler.go index cc6f28f0..e53accdc 100644 --- a/pkg/spireentry/reconciler.go +++ b/pkg/spireentry/reconciler.go @@ -34,6 +34,7 @@ import ( "github.com/google/uuid" lru "github.com/hashicorp/golang-lru/v2" "github.com/prometheus/client_golang/prometheus" + "github.com/samber/lo" "github.com/spiffe/go-spiffe/v2/spiffeid" "google.golang.org/grpc/codes" corev1 "k8s.io/api/core/v1" @@ -444,18 +445,9 @@ func (r *entryReconciler) entryListHints(clusterSPIFFEIDs []*ClusterSPIFFEID, cl } func uniqueStrings(values []string) []string { - seen := make(map[string]struct{}, len(values)) - out := make([]string, 0, len(values)) - for _, value := range values { - if value == "" { - continue - } - if _, ok := seen[value]; ok { - continue - } - seen[value] = struct{}{} - out = append(out, value) - } + out := lo.Uniq(lo.Filter(values, func(value string, _ int) bool { + return value != "" + })) sort.Strings(out) return out } From c2410c3bffe52735c0407f6c1e8b6ebe6317df9f Mon Sep 17 00:00:00 2001 From: jungooji <46273764+zzzz465@users.noreply.github.com> Date: Wed, 1 Jul 2026 16:20:03 +0900 Subject: [PATCH 3/5] refactor: use lo parallel map for hint listing Signed-off-by: jungooji <46273764+zzzz465@users.noreply.github.com> --- pkg/spireapi/entryapi.go | 51 +++++++++++++++++----------------------- 1 file changed, 21 insertions(+), 30 deletions(-) diff --git a/pkg/spireapi/entryapi.go b/pkg/spireapi/entryapi.go index 27479da4..218ade44 100644 --- a/pkg/spireapi/entryapi.go +++ b/pkg/spireapi/entryapi.go @@ -19,9 +19,9 @@ package spireapi import ( "context" "fmt" - "sync" "github.com/samber/lo" + "github.com/samber/lo/parallel" entryv1 "github.com/spiffe/spire-api-sdk/proto/spire/api/server/entry/v1" "github.com/spiffe/spire-api-sdk/proto/spire/api/types" "google.golang.org/grpc" @@ -70,38 +70,29 @@ func (c entryClient) ListEntries(ctx context.Context, hints ...string) ([]Entry, ctx, cancel := context.WithCancel(ctx) defer cancel() - var wg sync.WaitGroup - var errOnce sync.Once - var firstErr error - entriesByHint := make([][]Entry, len(filterHints)) - - for i, hint := range filterHints { - wg.Add(1) - go func(i int, hint string) { - defer wg.Done() - - hintEntries, err := c.listEntries(ctx, &entryv1.ListEntriesRequest_Filter{ - ByHint: wrapperspb.String(hint), - }) - if err != nil { - errOnce.Do(func() { - firstErr = err - cancel() - }) - return - } - entriesByHint[i] = hintEntries - }(i, hint) - } - wg.Wait() - if firstErr != nil { - return nil, firstErr + type listEntriesResult struct { + entries []Entry + err error } - var entries []Entry - for _, hintEntries := range entriesByHint { - entries = append(entries, hintEntries...) + results := parallel.Map(filterHints, func(hint string, _ int) listEntriesResult { + entries, err := c.listEntries(ctx, &entryv1.ListEntriesRequest_Filter{ + ByHint: wrapperspb.String(hint), + }) + if err != nil { + cancel() + } + return listEntriesResult{entries: entries, err: err} + }) + if result, ok := lo.Find(results, func(result listEntriesResult) bool { + return result.err != nil + }); ok { + return nil, result.err } + + entries := lo.FlatMap(results, func(result listEntriesResult, _ int) []Entry { + return result.entries + }) return lo.UniqBy(entries, func(entry Entry) string { return entry.ID }), nil From bb623823826137d411eb51aa7e0e90153807dbf8 Mon Sep 17 00:00:00 2001 From: jungooji <46273764+zzzz465@users.noreply.github.com> Date: Wed, 1 Jul 2026 16:28:27 +0900 Subject: [PATCH 4/5] ci: use repository image tag in PR image test Signed-off-by: jungooji <46273764+zzzz465@users.noreply.github.com> --- .github/workflows/pr_build.yaml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/.github/workflows/pr_build.yaml b/.github/workflows/pr_build.yaml index 400c3818..230f399d 100644 --- a/.github/workflows/pr_build.yaml +++ b/.github/workflows/pr_build.yaml @@ -85,7 +85,7 @@ jobs: make load-images - name: Test image run: | - docker tag ghcr.io/spiffe/spire-controller-manager:devel ghcr.io/spiffe/spire-controller-manager:nightly + docker tag "ghcr.io/${GITHUB_REPOSITORY}:devel" ghcr.io/spiffe/spire-controller-manager:devel (cd demo; ./test.sh) success: From 173e9ca87324fc3e0a2938077b08122c295bbcd1 Mon Sep 17 00:00:00 2001 From: jungooji <46273764+zzzz465@users.noreply.github.com> Date: Wed, 1 Jul 2026 16:50:43 +0900 Subject: [PATCH 5/5] refactor: simplify entry list hint handling Signed-off-by: jungooji <46273764+zzzz465@users.noreply.github.com> --- pkg/spireentry/reconciler.go | 41 +++++++++++------------------------- 1 file changed, 12 insertions(+), 29 deletions(-) diff --git a/pkg/spireentry/reconciler.go b/pkg/spireentry/reconciler.go index e53accdc..a565092a 100644 --- a/pkg/spireentry/reconciler.go +++ b/pkg/spireentry/reconciler.go @@ -404,7 +404,7 @@ func (r *entryReconciler) incCounter(name string) { // server on every reconcile; otherwise it lists from the server and (when the // cache is enabled) repopulates it. func (r *entryReconciler) listEntries(ctx context.Context, hints []string) ([]spireapi.Entry, []spireapi.Entry, error) { - filterKey := entryListFilterKey(hints) + filterKey := strings.Join(hints, "\x00") if r.entryCache != nil && r.entryCache.fresh(filterKey) { r.incCounter(metrics.EntryListCacheHits) currentEntries, deleteOnlyEntries := r.splitEntries(r.entryCache.snapshot()) @@ -428,35 +428,18 @@ func (r *entryReconciler) entryListHints(clusterSPIFFEIDs []*ClusterSPIFFEID, cl return nil } - var hints []string - for _, clusterSPIFFEID := range clusterSPIFFEIDs { - if clusterSPIFFEID.Spec.Hint == "" { - return nil - } - hints = append(hints, clusterSPIFFEID.Spec.Hint) - } - for _, clusterStaticEntry := range clusterStaticEntries { - if clusterStaticEntry.Spec.Hint == "" { - return nil - } - hints = append(hints, clusterStaticEntry.Spec.Hint) - } - return uniqueStrings(hints) -} - -func uniqueStrings(values []string) []string { - out := lo.Uniq(lo.Filter(values, func(value string, _ int) bool { - return value != "" - })) - sort.Strings(out) - return out -} - -func entryListFilterKey(hints []string) string { - if len(hints) == 0 { - return "" + hints := lo.Map(clusterSPIFFEIDs, func(clusterSPIFFEID *ClusterSPIFFEID, _ int) string { + return clusterSPIFFEID.Spec.Hint + }) + hints = append(hints, lo.Map(clusterStaticEntries, func(clusterStaticEntry *ClusterStaticEntry, _ int) string { + return clusterStaticEntry.Spec.Hint + })...) + if lo.Contains(hints, "") { + return nil } - return strings.Join(uniqueStrings(hints), "\x00") + hints = lo.Uniq(hints) + sort.Strings(hints) + return hints } // splitEntries partitions entries into the set to process and the set marked