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
64 changes: 44 additions & 20 deletions internal/runtime/lifecycle.go
Original file line number Diff line number Diff line change
Expand Up @@ -2434,36 +2434,29 @@ func (r *Runtime) cleanupOrphanedIndexEntries() {

r.logger.Debug("Checking for orphaned index entries")

activeServers := r.upstreamManager.GetAllServerNames()
activeServerMap := make(map[string]bool)
for _, serverName := range activeServers {
activeServerMap[serverName] = true
}

indexedServers, err := r.indexManager.GetAllIndexedServerNames()
orphans, activeCount, indexedCount, err := findOrphanedIndexServers(
r.indexManager.GetAllIndexedServerNames, r.upstreamManager.GetAllServerNames)
if err != nil {
r.logger.Warn("Failed to retrieve indexed server names for orphan cleanup", zap.Error(err))
return
}

var removedCount int
for _, indexedServer := range indexedServers {
if !activeServerMap[indexedServer] {
r.logger.Info("Removing orphaned index entries for server no longer in config",
zap.String("server", indexedServer))
if err := r.indexManager.DeleteServerTools(indexedServer); err != nil {
r.logger.Warn("Failed to delete orphaned index entries",
zap.String("server", indexedServer),
zap.Error(err))
} else {
removedCount++
}
for _, indexedServer := range orphans {
r.logger.Info("Removing orphaned index entries for server no longer in config",
zap.String("server", indexedServer))
if err := r.indexManager.DeleteServerTools(indexedServer); err != nil {
r.logger.Warn("Failed to delete orphaned index entries",
zap.String("server", indexedServer),
zap.Error(err))
} else {
removedCount++
}
}

r.logger.Debug("Orphaned index cleanup completed",
zap.Int("active_servers", len(activeServers)),
zap.Int("indexed_servers", len(indexedServers)),
zap.Int("active_servers", activeCount),
zap.Int("indexed_servers", indexedCount),
zap.Int("orphans_removed", removedCount))

if removedCount > 0 {
Expand All @@ -2472,6 +2465,37 @@ func (r *Runtime) cleanupOrphanedIndexEntries() {
}
}

// findOrphanedIndexServers returns the indexed servers that are not active,
// plus the sizes of both lists for logging.
func findOrphanedIndexServers(
listIndexed func() ([]string, error),
listActive func() []string,
) (orphans []string, activeCount, indexedCount int, err error) {
// Indexed first: a server is registered before its tools are indexed, so a
// server added while this runs is either missing from the indexed list or
// already present in the active list read after it. This covers servers newly
// added during cleanup; index entries persisted from a previous run can still
// be pruned if their server's async re-registration has not finished yet, and
// are re-indexed once discovery runs for it.
indexedServers, err := listIndexed()
if err != nil {
return nil, 0, 0, err
}

activeServers := listActive()
active := make(map[string]bool, len(activeServers))
for _, name := range activeServers {
active[name] = true
}

for _, name := range indexedServers {
if !active[name] {
orphans = append(orphans, name)
}
}
return orphans, len(activeServers), len(indexedServers), nil
}

// supervisorEventForwarder subscribes to supervisor events and emits runtime events
// to notify Web UI via SSE when server connection state changes.
func (r *Runtime) supervisorEventForwarder() {
Expand Down
46 changes: 46 additions & 0 deletions internal/runtime/orphan_cleanup_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,46 @@
package runtime

import (
"testing"

"github.com/stretchr/testify/assert"
"github.com/stretchr/testify/require"
)

// A server registered and indexed while the startup orphan cleanup runs must
// not be treated as an orphan. The event fires right after the cleanup's first
// read, whichever list that is.
func TestFindOrphanedIndexServers_ServerAddedDuringCleanupIsKept(t *testing.T) {
var active, indexed []string
added := false
addServer := func() {
if !added {
added = true
active = append(active, "b")
indexed = append(indexed, "b")
}
}
listActive := func() []string {
defer addServer()
return append([]string(nil), active...)
}
listIndexed := func() ([]string, error) {
defer addServer()
return append([]string(nil), indexed...), nil
}

orphans, _, _, err := findOrphanedIndexServers(listIndexed, listActive)
require.NoError(t, err)
assert.Empty(t, orphans)
}

func TestFindOrphanedIndexServers_RemovedServerIsOrphan(t *testing.T) {
orphans, activeCount, indexedCount, err := findOrphanedIndexServers(
func() ([]string, error) { return []string{"a", "gone"}, nil },
func() []string { return []string{"a"} },
)
require.NoError(t, err)
assert.Equal(t, []string{"gone"}, orphans)
assert.Equal(t, 1, activeCount)
assert.Equal(t, 2, indexedCount)
}
Loading