diff --git a/pkg/cache/v3/delete_resources_test.go b/pkg/cache/v3/delete_resources_test.go new file mode 100644 index 000000000..99978c493 --- /dev/null +++ b/pkg/cache/v3/delete_resources_test.go @@ -0,0 +1,364 @@ +package cache_test + +import ( + "context" + "strconv" + "sync" + "testing" + + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + + clusterv3 "github.com/envoyproxy/go-control-plane/envoy/config/cluster/v3" + core "github.com/envoyproxy/go-control-plane/envoy/config/core/v3" + endpointv3 "github.com/envoyproxy/go-control-plane/envoy/config/endpoint/v3" + listenerv3 "github.com/envoyproxy/go-control-plane/envoy/config/listener/v3" + "github.com/envoyproxy/go-control-plane/pkg/cache/types" + "github.com/envoyproxy/go-control-plane/pkg/cache/v3" + rsrc "github.com/envoyproxy/go-control-plane/pkg/resource/v3" + "github.com/envoyproxy/go-control-plane/pkg/server/stream/v3" +) + +// DeleteResources was a no-op (body commented out in 3ca8faf7). A delete must drop the +// resource and tell a parked wildcard CDS watch via removed_resources. +func TestDeleteResources_RemovesAndNotifiesDeltaWatch(t *testing.T) { + ctx := context.Background() + c := cache.NewSnapshotCache(true, group{}, nil) + const node = "n1" + require.NoError(t, c.UpsertResources(ctx, node, rsrc.ClusterType, map[string]*types.ResourceWithTTL{ + "a": {Resource: &clusterv3.Cluster{Name: "a"}, Version: "1"}, + "b": {Resource: &clusterv3.Cluster{Name: "b"}, Version: "1"}, + })) + + req := &cache.DeltaRequest{TypeUrl: rsrc.ClusterType, Node: &core.Node{Id: node}} + // First watch answers immediately with a+b; adopt its versions as the client's state. + first := make(chan cache.DeltaResponse, 1) + _, delayed := c.CreateDeltaWatch(req, stream.NewStreamState(true, nil), first) + require.False(t, delayed) + state := stream.NewStreamState(true, nil) + state.SetResourceVersions((<-first).GetNextVersionMap()) // also clears "first" + + // Second watch parks: the client is up to date. + second := make(chan cache.DeltaResponse, 1) + _, delayed = c.CreateDeltaWatch(req, state, second) + require.True(t, delayed) + + require.NoError(t, c.DeleteResources(ctx, node, rsrc.ClusterType, []string{"b"})) + + select { + case r := <-second: + dr, err := r.GetDeltaDiscoveryResponse() + require.NoError(t, err) + assert.Equal(t, []string{"b"}, dr.GetRemovedResources()) + default: + t.Fatal("parked watch was not told about the removal") + } + snap, err := c.GetSnapshot(node) + require.NoError(t, err) + _, stillThere := snap.GetResourcesAndTTL(rsrc.ClusterType)["b"] + assert.False(t, stillThere) +} + +// A named (EDS-style) subscription to a CLA that gets deleted must be told via +// removed_resources. +func TestDeleteResources_NamedSubscriptionNotified(t *testing.T) { + ctx := context.Background() + c := cache.NewSnapshotCache(true, group{}, nil) + const node = "n1" + require.NoError(t, c.UpsertResources(ctx, node, rsrc.EndpointType, map[string]*types.ResourceWithTTL{ + "a": {Resource: &endpointv3.ClusterLoadAssignment{ClusterName: "a"}, Version: "1"}, + "b": {Resource: &endpointv3.ClusterLoadAssignment{ClusterName: "b"}, Version: "1"}, + })) + + req := &cache.DeltaRequest{TypeUrl: rsrc.EndpointType, Node: &core.Node{Id: node}} + sub := map[string]struct{}{"b": {}} + + first := stream.NewStreamState(false, nil) + first.SetSubscribedResourceNames(sub) + ch1 := make(chan cache.DeltaResponse, 1) + _, delayed := c.CreateDeltaWatch(req, first, ch1) + require.False(t, delayed) + + state := stream.NewStreamState(false, nil) + state.SetSubscribedResourceNames(sub) + state.SetResourceVersions((<-ch1).GetNextVersionMap()) + ch2 := make(chan cache.DeltaResponse, 1) + _, delayed = c.CreateDeltaWatch(req, state, ch2) + require.True(t, delayed) + + require.NoError(t, c.DeleteResources(ctx, node, rsrc.EndpointType, []string{"b"})) + + select { + case r := <-ch2: + dr, err := r.GetDeltaDiscoveryResponse() + require.NoError(t, err) + assert.Equal(t, []string{"b"}, dr.GetRemovedResources()) + default: + t.Fatal("named watch was not told about the removal") + } +} + +// Deleting the LAST resource of a type while no watch is parked: a client that then +// re-subscribes with its old versions must be answered immediately with the removal. +func TestDeleteResources_LastResourceRemovalOnResubscribe(t *testing.T) { + ctx := context.Background() + c := cache.NewSnapshotCache(true, group{}, nil) + const node = "n1" + require.NoError(t, c.UpsertResources(ctx, node, rsrc.ClusterType, map[string]*types.ResourceWithTTL{ + "b": {Resource: &clusterv3.Cluster{Name: "b"}, Version: "1"}, + })) + req := &cache.DeltaRequest{TypeUrl: rsrc.ClusterType, Node: &core.Node{Id: node}} + first := make(chan cache.DeltaResponse, 1) + _, delayed := c.CreateDeltaWatch(req, stream.NewStreamState(true, nil), first) + require.False(t, delayed) + state := stream.NewStreamState(true, nil) + state.SetResourceVersions((<-first).GetNextVersionMap()) + + require.NoError(t, c.DeleteResources(ctx, node, rsrc.ClusterType, []string{"b"})) + + ch := make(chan cache.DeltaResponse, 1) + _, delayed = c.CreateDeltaWatch(req, state, ch) + require.False(t, delayed, "removal of the last resource must answer immediately") + dr, err := (<-ch).GetDeltaDiscoveryResponse() + require.NoError(t, err) + assert.Equal(t, []string{"b"}, dr.GetRemovedResources()) + + // A fresh stream on the now-empty type must still park, not get an empty response. + fresh := make(chan cache.DeltaResponse, 1) + _, delayed = c.CreateDeltaWatch(req, stream.NewStreamState(true, nil), fresh) + assert.True(t, delayed) +} + +// Unknown names, unknown type URLs and nodes without a snapshot are silent no-ops: no +// error, no version bump, no watch response. +func TestDeleteResources_NoopCases(t *testing.T) { + ctx := context.Background() + c := cache.NewSnapshotCache(true, group{}, nil) + require.NoError(t, c.DeleteResources(ctx, "missing-node", rsrc.ClusterType, []string{"x"})) + + require.NoError(t, c.UpsertResources(ctx, "n1", rsrc.ClusterType, map[string]*types.ResourceWithTTL{ + "a": {Resource: &clusterv3.Cluster{Name: "a"}, Version: "1"}, + })) + req := &cache.DeltaRequest{TypeUrl: rsrc.ClusterType, Node: &core.Node{Id: "n1"}} + first := make(chan cache.DeltaResponse, 1) + _, delayed := c.CreateDeltaWatch(req, stream.NewStreamState(true, nil), first) + require.False(t, delayed) + state := stream.NewStreamState(true, nil) + state.SetResourceVersions((<-first).GetNextVersionMap()) + parked := make(chan cache.DeltaResponse, 1) + _, delayed = c.CreateDeltaWatch(req, state, parked) + require.True(t, delayed) + + before, err := c.GetSnapshot("n1") + require.NoError(t, err) + v := before.GetVersion(rsrc.ClusterType) + require.NoError(t, c.DeleteResources(ctx, "n1", rsrc.ClusterType, []string{"not-there"})) + require.NoError(t, c.DeleteResources(ctx, "n1", "type.googleapis.com/unknown.Type", []string{"a"})) + after, err := c.GetSnapshot("n1") + require.NoError(t, err) + assert.Equal(t, v, after.GetVersion(rsrc.ClusterType)) + assert.Empty(t, parked, "no-op delete must not answer the parked watch") +} + +// A restarted control plane builds types one at a time. A client reconnecting with +// versions for a type the snapshot has never populated must be parked, not told to +// remove what it holds, even when another type is upserted meanwhile. +func TestCreateDeltaWatch_NeverPopulatedTypeParksOnResubscribe(t *testing.T) { + testNeverPopulatedTypeParksOnResubscribe(t, true) +} + +func TestCreateDeltaWatch_NeverPopulatedTypeParksOnResubscribeNonADS(t *testing.T) { + testNeverPopulatedTypeParksOnResubscribe(t, false) +} + +func testNeverPopulatedTypeParksOnResubscribe(t *testing.T, ads bool) { + ctx := context.Background() + c := cache.NewSnapshotCache(ads, group{}, nil) + const node = "n1" + require.NoError(t, c.UpsertResources(ctx, node, rsrc.ClusterType, map[string]*types.ResourceWithTTL{ + "c1": {Resource: &clusterv3.Cluster{Name: "c1"}, Version: "1"}, + })) + + req := &cache.DeltaRequest{TypeUrl: rsrc.ListenerType, Node: &core.Node{Id: node}, InitialResourceVersions: map[string]string{"L": "v1"}} + ch := make(chan cache.DeltaResponse, 1) + _, delayed := c.CreateDeltaWatch(req, stream.NewStreamState(true, map[string]string{"L": "v1"}), ch) + require.True(t, delayed, "never-populated type must park") + require.Empty(t, ch) + + require.NoError(t, c.UpsertResources(ctx, node, rsrc.ClusterType, map[string]*types.ResourceWithTTL{ + "c2": {Resource: &clusterv3.Cluster{Name: "c2"}, Version: "1"}, + })) + require.Empty(t, ch, "upsert of another type must not answer with a removal") + + require.NoError(t, c.UpsertResources(ctx, node, rsrc.ListenerType, map[string]*types.ResourceWithTTL{ + "L": {Resource: &listenerv3.Listener{Name: "L"}, Version: "v2"}, + })) + select { + case r := <-ch: + dr, err := r.GetDeltaDiscoveryResponse() + require.NoError(t, err) + assert.Empty(t, dr.GetRemovedResources()) + require.Len(t, dr.GetResources(), 1) + assert.Equal(t, "L", dr.GetResources()[0].GetName()) + default: + t.Fatal("parked watch was not answered once the type was set") + } +} + +// A type that was populated and emptied by deleting its last resource is a real +// delete: a reconnecting client holding it must be told to remove it. +func TestCreateDeltaWatch_EmptiedTypeRemovesOnResubscribe(t *testing.T) { + ctx := context.Background() + c := cache.NewSnapshotCache(true, group{}, nil) + const node = "n1" + require.NoError(t, c.UpsertResources(ctx, node, rsrc.ListenerType, map[string]*types.ResourceWithTTL{ + "L": {Resource: &listenerv3.Listener{Name: "L"}, Version: "v1"}, + })) + require.NoError(t, c.DeleteResources(ctx, node, rsrc.ListenerType, []string{"L"})) + + req := &cache.DeltaRequest{TypeUrl: rsrc.ListenerType, Node: &core.Node{Id: node}, InitialResourceVersions: map[string]string{"L": "v1"}} + ch := make(chan cache.DeltaResponse, 1) + _, delayed := c.CreateDeltaWatch(req, stream.NewStreamState(true, map[string]string{"L": "v1"}), ch) + require.False(t, delayed) + dr, err := (<-ch).GetDeltaDiscoveryResponse() + require.NoError(t, err) + assert.Equal(t, []string{"L"}, dr.GetRemovedResources()) +} + +// A SOTW CDS client that is up to date parks its watch in the shared cache. Deleting a +// cluster must send that watch a replacement state-of-the-world without the cluster. +func TestDeleteResources_NotifiesParkedSOTWWatch(t *testing.T) { + ctx := context.Background() + c := cache.NewSnapshotCache(true, group{}, nil) + const node = "n1" + require.NoError(t, c.UpsertResources(ctx, node, rsrc.ClusterType, map[string]*types.ResourceWithTTL{ + "a": {Resource: &clusterv3.Cluster{Name: "a"}, Version: "1"}, + "b": {Resource: &clusterv3.Cluster{Name: "b"}, Version: "1"}, + })) + snap, err := c.GetSnapshot(node) + require.NoError(t, err) + version := snap.GetVersion(rsrc.ClusterType) + + req := &cache.Request{TypeUrl: rsrc.ClusterType, Node: &core.Node{Id: node}, VersionInfo: version} + ch := make(chan cache.Response, 1) + cancel := c.CreateWatch(req, stream.NewStreamState(true, nil), ch) + require.NotNil(t, cancel) + defer cancel() + require.Empty(t, ch, "up-to-date SOTW request must park") + + require.NoError(t, c.DeleteResources(ctx, node, rsrc.ClusterType, []string{"b"})) + + select { + case r := <-ch: + raw, ok := r.(*cache.RawResponse) + require.True(t, ok) + names := make([]string, 0, len(raw.Resources)) + for _, res := range raw.Resources { + names = append(names, res.Name) + } + assert.Equal(t, []string{"a"}, names) + assert.NotEqual(t, version, raw.Version) + default: + t.Fatal("parked SOTW watch was not notified of the cluster deletion") + } +} + +// DeleteResources mutates the snapshot's Items map in place under Snapshot.Mu, so +// CreateDeltaWatch must hold it to read that map. Meaningful under -race. +func TestDeleteResources_ConcurrentCreateDeltaWatch(t *testing.T) { + ctx := context.Background() + c := cache.NewSnapshotCache(true, group{}, nil) + const node = "n1" + clusters := make(map[string]*types.ResourceWithTTL) + for i := 0; i < 100; i++ { + name := strconv.Itoa(i) + clusters[name] = &types.ResourceWithTTL{Resource: &clusterv3.Cluster{Name: name}, Version: "1"} + } + require.NoError(t, c.UpsertResources(ctx, node, rsrc.ClusterType, clusters)) + + done := make(chan struct{}) + go func() { + defer close(done) + for name := range clusters { + assert.NoError(t, c.DeleteResources(ctx, node, rsrc.ClusterType, []string{name})) + } + }() + req := &cache.DeltaRequest{TypeUrl: rsrc.ClusterType, Node: &core.Node{Id: node}} + for i := 0; i < 100; i++ { + if cancel, _ := c.CreateDeltaWatch(req, stream.NewStreamState(true, nil), make(chan cache.DeltaResponse, 1)); cancel != nil { + cancel() + } + } + <-done +} + +// An upsert that lands between CreateDeltaWatch reading the snapshot and registering +// its watch must not leave that watch parked while the snapshot has the data. +func TestCreateDeltaWatch_ConcurrentUpsertNotMissed(t *testing.T) { + ctx := context.Background() + req := &cache.DeltaRequest{TypeUrl: rsrc.ClusterType, Node: &core.Node{Id: "n1"}} + for i := 0; i < 10000; i++ { + c := cache.NewSnapshotCache(true, group{}, nil) + ch := make(chan cache.DeltaResponse, 1) + var wg sync.WaitGroup + wg.Add(1) + go func() { + defer wg.Done() + c.CreateDeltaWatch(req, stream.NewStreamState(true, nil), ch) + }() + require.NoError(t, c.UpsertResources(ctx, "n1", rsrc.ClusterType, map[string]*types.ResourceWithTTL{ + "a": {Resource: &clusterv3.Cluster{Name: "a"}, Version: "1"}, + })) + wg.Wait() + require.Len(t, ch, 1, "iteration %d: watch parked although the snapshot has the cluster", i) + } +} + +// SOTW responses read the snapshot's maps, which UpsertResources and DeleteResources +// mutate in place under Snapshot.Mu, while a SOTW client keeps re-parking its watch. +// Meaningful under -race. +func TestDeleteResources_ConcurrentSOTWWatch(t *testing.T) { + ctx := context.Background() + c := cache.NewSnapshotCache(true, group{}, nil) + const node = "n1" + require.NoError(t, c.UpsertResources(ctx, node, rsrc.ClusterType, map[string]*types.ResourceWithTTL{ + "a": {Resource: &clusterv3.Cluster{Name: "a"}, Version: "1"}, + })) + + // Two writers, so one's SOTW responses overlap the other's map writes. + var wg sync.WaitGroup + for _, prefix := range []string{"x", "y"} { + wg.Add(1) + go func(prefix string) { + defer wg.Done() + for i := 0; i < 100; i++ { + name := prefix + strconv.Itoa(i) + assert.NoError(t, c.UpsertResources(ctx, node, rsrc.ClusterType, map[string]*types.ResourceWithTTL{ + name: {Resource: &clusterv3.Cluster{Name: name}, Version: "1"}, + })) + assert.NoError(t, c.DeleteResources(ctx, node, rsrc.ClusterType, []string{name})) + } + }(prefix) + } + done := make(chan struct{}) + go func() { + wg.Wait() + close(done) + }() + version := "" + for i := 0; i < 200; i++ { + ch := make(chan cache.Response, 1) + cancel := c.CreateWatch(&cache.Request{TypeUrl: rsrc.ClusterType, Node: &core.Node{Id: node}, VersionInfo: version}, stream.NewStreamState(true, nil), ch) + select { + case r := <-ch: + v, err := r.GetVersion() + require.NoError(t, err) + version = v + case <-done: + } + if cancel != nil { + cancel() + } + } + <-done +} diff --git a/pkg/cache/v3/simple.go b/pkg/cache/v3/simple.go index 789c1d51f..08d691d3a 100644 --- a/pkg/cache/v3/simple.go +++ b/pkg/cache/v3/simple.go @@ -320,8 +320,7 @@ func (cache *snapshotCache) sendHeartbeats(ctx context.Context, node *core.Node) info.mu.Lock() for id, watch := range info.watches { // Respond with the current version regardless of whether the version has changed. - version := snapshot.GetVersion(watch.Request.GetTypeUrl()) - resources := snapshot.GetResourcesAndTTL(watch.Request.GetTypeUrl()) + version, resources := sotwResources(snapshot, watch.Request.GetTypeUrl()) // TODO(snowp): Construct this once per type instead of once per watch. resourcesWithTTL := map[string]VTMarshaledResource{} @@ -581,36 +580,44 @@ func (cache *snapshotCache) UpdateVirtualHosts(ctx context.Context, _ string, ty } func (cache *snapshotCache) DeleteResources(ctx context.Context, node string, typ string, resourcesToDeleted []string) error { - // cache.mu.Lock() - // defer cache.mu.Unlock() - - // if typ == resource.ClusterType { - // index := GetResponseType(typ) - // snapshot := cache.snapshots[node] - // prevResources := snapshot.(*Snapshot).Resources[index] - // currentVersion := cache.ParseSystemVersionInfo(prevResources.Version) - - // for _, k := range resourcesToDeleted { - // delete(prevResources.Items, k) - // } - - // currentVersion++ - // prevResources.Version = fmt.Sprintf("%d", currentVersion) - // // Update - // snapshot.(*Snapshot).Resources[index] = prevResources - // cache.snapshots[node] = snapshot - - // // Respond deltas - // if info, ok := cache.status[node]; ok { - // info.mu.Lock() - // defer info.mu.Unlock() - - // // Respond to delta watches for the node. - // return cache.respondDeltaWatches(ctx, info, snapshot) - // } + index := GetResponseType(typ) + if index == types.UnknownType { + return nil + } + snapshot := cache.getSnapshot(node) + if snapshot == nil { + return nil // nothing cached for this node, so nothing to remove + } - // } + snapshot.(*Snapshot).Mu.Lock() + currentResources := snapshot.(*Snapshot).Resources[index] + removed := 0 + for _, name := range resourcesToDeleted { + if _, ok := currentResources.Items[name]; ok { + delete(currentResources.Items, name) + removed++ + } + } + if removed == 0 { + snapshot.(*Snapshot).Mu.Unlock() + return nil + } + currentVersion := cache.ParseSystemVersionInfo(currentResources.Version) + currentVersion++ + currentResources.Version = fmt.Sprintf("%d", currentVersion) + snapshot.(*Snapshot).Resources[index] = currentResources + cache.putSnapshot(node, snapshot) + snapshot.(*Snapshot).Mu.Unlock() + // VersionMap needs no patching: respondDeltaWatches rebuilds it via ConstructVersionMap. + if info := cache.getStatus(node); info != nil { + info.mu.Lock() + defer info.mu.Unlock() + if err := cache.respondSOTWWatches(ctx, info, snapshot); err != nil { + return err + } + return cache.respondDeltaWatches(ctx, info, snapshot) + } return nil } @@ -643,10 +650,9 @@ func (cache *snapshotCache) SetSnapshot(ctx context.Context, node string, snapsh func (cache *snapshotCache) respondSOTWWatches(ctx context.Context, info *statusInfo, snapshot ResourceSnapshot) error { // responder callback for SOTW watches respond := func(watch ResponseWatch, id int64) error { - version := snapshot.GetVersion(watch.Request.GetTypeUrl()) + version, resources := sotwResources(snapshot, watch.Request.GetTypeUrl()) if version != watch.Request.GetVersionInfo() { cache.log.Debugf("respond open watch %d %s%v with new version %q", id, watch.Request.GetTypeUrl(), watch.Request.GetResourceNames(), version) - resources := snapshot.GetResourcesAndTTL(watch.Request.GetTypeUrl()) err := cache.respond(ctx, watch.Request, watch.Response, resources, version, false) if err != nil { return err @@ -819,9 +825,10 @@ func (cache *snapshotCache) CreateWatch(request *Request, streamState stream.Str info.mu.Unlock() var version string + var resources map[string]VTMarshaledResource snapshot := cache.getSnapshot(nodeID) if snapshot != nil { - version = snapshot.GetVersion(request.GetTypeUrl()) + version, resources = sotwResources(snapshot, request.GetTypeUrl()) } if snapshot != nil { @@ -837,7 +844,6 @@ func (cache *snapshotCache) CreateWatch(request *Request, streamState stream.Str request.GetTypeUrl(), request.GetResourceNames(), knownResourceNames, diff) if len(diff) > 0 { - resources := snapshot.GetResourcesAndTTL(request.GetTypeUrl()) for _, name := range diff { if _, exists := resources[name]; exists { if err := cache.respond(context.Background(), request, value, resources, version, false); err != nil { @@ -866,7 +872,6 @@ func (cache *snapshotCache) CreateWatch(request *Request, streamState stream.Str } // otherwise, the watch may be responded immediately - resources := snapshot.GetResourcesAndTTL(request.GetTypeUrl()) if err := cache.respond(context.Background(), request, value, resources, version, false); err != nil { cache.log.Errorf("failed to send a response for %s%v to nodeID %q: %s", request.GetTypeUrl(), request.GetResourceNames(), nodeID, err) @@ -892,6 +897,23 @@ func (cache *snapshotCache) cancelWatch(nodeID string, watchID int64) func() { } } +// sotwResources returns the version and a copy of the resources of typeURL, read under +// Snapshot.Mu: UpsertResources and DeleteResources mutate a *Snapshot's maps in place, so a +// SOTW response must not iterate them after the lock is released. +// ponytail: copies the map on every SOTW read; build the response under the lock if that gets hot. +func sotwResources(snapshot ResourceSnapshot, typeURL string) (string, map[string]VTMarshaledResource) { + if s, ok := snapshot.(*Snapshot); ok { + s.Mu.RLock() + defer s.Mu.RUnlock() + } + resources := snapshot.GetResourcesAndTTL(typeURL) + out := make(map[string]VTMarshaledResource, len(resources)) + for name, r := range resources { + out[name] = r + } + return snapshot.GetVersion(typeURL), out +} + // Respond to a watch with the snapshot value. The value channel should have capacity not to block. // TODO(kuat) do not respond always, see issue https://github.com/envoyproxy/go-control-plane/issues/46 func (cache *snapshotCache) respond(ctx context.Context, request *Request, value chan Response, resources map[string]VTMarshaledResource, version string, heartbeat bool) error { @@ -954,11 +976,6 @@ func (cache *snapshotCache) CreateDeltaWatch(request *DeltaRequest, state stream // update last watch request time info.setLastDeltaWatchRequestTime(time.Now()) - // find the current cache snapshot for the provided node - snapshot := cache.getSnapshot(nodeID) - // snapshot exists and we have resources of the typeUrl on the server - exists := snapshot != nil && len(snapshot.GetResourcesAndTTL(request.GetTypeUrl())) > 0 - // There are three different cases that leads to a delayed watch trigger: // - no snapshot exists for the requested nodeID // - a snapshot exists, but we failed to initialize its version map @@ -970,10 +987,24 @@ func (cache *snapshotCache) CreateDeltaWatch(request *DeltaRequest, state stream // exclusive. Without it there is a window after respondDelta concludes there // is nothing to send and before the watch is registered, during which an // upsert can mutate the snapshot, walk the watches, not find this one, and - // silently skip the stream. + // silently skip the stream. Reading the snapshot below is part of that step. info.mu.Lock() defer info.mu.Unlock() + // find the current cache snapshot for the provided node + snapshot := cache.getSnapshot(nodeID) + // snapshot exists and we have resources of the typeUrl on the server + // A client that still holds versions must be answered even when the type is now empty + // (it was populated and its last resource deleted), otherwise that removal is never + // reported. A type never populated is parked by respondDelta instead. + exists := false + if snapshot != nil { + // UpsertResources and DeleteResources mutate the Items map in place under Mu. + snapshot.(*Snapshot).Mu.RLock() + exists = len(snapshot.GetResourcesAndTTL(request.GetTypeUrl())) > 0 || len(state.GetResourceVersions()) > 0 + snapshot.(*Snapshot).Mu.RUnlock() + } + delayedResponse := !exists if exists { err := snapshot.ConstructVersionMap() @@ -1031,8 +1062,16 @@ func (cache *snapshotCache) respondDelta(ctx context.Context, snapshot ResourceS if elapsedLock > 1*time.Millisecond { fmt.Printf("respondDelta took %s to lock\n", elapsedLock) } + resourceMap := snapshot.GetResourcesAndTTL(request.GetTypeUrl()) + // A nil map means the type was never populated (deleting its last resource leaves a + // non-nil empty map). Telling a reconnecting client to remove what it holds would drop + // it until the type is set; park the watch instead. + if resourceMap == nil && len(state.GetResourceVersions()) > 0 { + snapshot.(*Snapshot).Mu.RUnlock() + return nil, nil + } resp := createDeltaResponse(ctx, request, state, resourceContainer{ - resourceMap: snapshot.GetResourcesAndTTL(request.GetTypeUrl()), + resourceMap: resourceMap, versionMap: snapshot.GetVersionMap(request.GetTypeUrl()), systemVersion: snapshot.GetVersion(request.GetTypeUrl()), }) @@ -1109,13 +1148,12 @@ func (cache *snapshotCache) Fetch(ctx context.Context, request *Request) (Respon if snapshot := cache.getSnapshot(nodeID); snapshot != nil { // Respond only if the request version is distinct from the current snapshot state. // It might be beneficial to hold the request since Envoy will re-attempt the refresh. - version := snapshot.GetVersion(request.GetTypeUrl()) + version, resources := sotwResources(snapshot, request.GetTypeUrl()) if request.GetVersionInfo() == version { cache.log.Warnf("skip fetch: version up to date") return nil, &types.SkipFetchError{} } - resources := snapshot.GetResourcesAndTTL(request.GetTypeUrl()) out := createResponse(ctx, request, resources, version, false) return out, nil }