Repository navigation
fix(cache): make dropped resources and dead streams observable #20
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
Changes from all commits
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,38 @@ | ||
| package cache | ||
|
|
||
| import ( | ||
| "errors" | ||
|
|
||
| "github.com/envoyproxy/go-control-plane/pkg/log" | ||
| ) | ||
|
|
||
| // ErrResponseChannelClosed is returned by respondDelta when the send panicked | ||
| // because the stream's response channel was already closed. | ||
| // | ||
| // It is deliberately distinct from a nil response with a nil error. That pair | ||
| // means "nothing to send, or the channel was full - keep the watch and retry on | ||
| // the next upsert", and retrying is correct there. A closed channel means the | ||
| // stream is gone and no retry will ever happen, so collapsing the two hides a | ||
| // dead watch behind a healthy-looking one. | ||
| var ErrResponseChannelClosed = errors.New("delta response channel is closed") | ||
|
|
||
| // OnResourceMarshalError is called when a resource cannot be marshalled while | ||
| // building or updating a snapshot. The resource is dropped: it is not added to | ||
| // the snapshot, so it is never pushed, and a proxy already holding it sees it | ||
| // disappear on the next diff. | ||
| // | ||
| // None of those call sites can return an error, so this hook is the only way to | ||
| // count the drop. Set it once at startup; it must be safe for concurrent use. | ||
| // Leaving it nil keeps the previous behaviour minus the stdout write. | ||
| var OnResourceMarshalError func(typeURL, name string, err error) | ||
|
|
||
| // reportMarshalError records a dropped resource. logger may be nil. | ||
| func reportMarshalError(logger log.Logger, typeURL, name string, err error) { | ||
| if logger != nil { | ||
| logger.Errorf("dropping resource %q of type %q from the snapshot: MarshalVTStrict failed: %v; "+ | ||
| "it will not be pushed to any proxy", name, typeURL, err) | ||
| } | ||
| if OnResourceMarshalError != nil { | ||
| OnResourceMarshalError(typeURL, name, err) | ||
| } | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,63 @@ | ||
| package cache_test | ||
|
|
||
| import ( | ||
| "errors" | ||
| "sync" | ||
| "testing" | ||
|
|
||
| cluster "github.com/envoyproxy/go-control-plane/envoy/config/cluster/v3" | ||
| "github.com/envoyproxy/go-control-plane/pkg/cache/types" | ||
| "github.com/envoyproxy/go-control-plane/pkg/cache/v3" | ||
| ) | ||
|
|
||
| var errMarshal = errors.New("synthetic marshal failure") | ||
|
|
||
| // unmarshalableResource is a real proto.Message whose MarshalVTStrict always | ||
| // fails, which is the case the cache used to swallow with a fmt.Printf. | ||
| type unmarshalableResource struct { | ||
| *cluster.Cluster | ||
| } | ||
|
|
||
| func (unmarshalableResource) MarshalVTStrict() ([]byte, error) { return nil, errMarshal } | ||
|
|
||
| // A resource that cannot be marshalled is dropped from the snapshot: it is | ||
| // never pushed, and a proxy already holding it sees it removed. Before this | ||
| // change the only trace was an unstructured line on stdout, which is | ||
| // unreadable in an environment logging 55k lines/min against a one-minute | ||
| // kubelet buffer. Nothing could count it. | ||
| func TestDroppedResourceIsReported(t *testing.T) { | ||
| var ( | ||
| mu sync.Mutex | ||
| calls int | ||
| gotErr error | ||
| ) | ||
| original := cache.OnResourceMarshalError | ||
| t.Cleanup(func() { cache.OnResourceMarshalError = original }) | ||
| cache.OnResourceMarshalError = func(_, _ string, err error) { | ||
| mu.Lock() | ||
| defer mu.Unlock() | ||
| calls++ | ||
| gotErr = err | ||
| } | ||
|
|
||
| items := []types.ResourceWithTTL{ | ||
| {Resource: unmarshalableResource{Cluster: &cluster.Cluster{Name: "bad"}}}, | ||
| {Resource: &cluster.Cluster{Name: "good"}}, | ||
| } | ||
| indexed := cache.IndexAndMarshalResourcesByName(items) | ||
|
|
||
| mu.Lock() | ||
| defer mu.Unlock() | ||
| if calls != 1 { | ||
| t.Errorf("OnResourceMarshalError called %d times, want 1", calls) | ||
| } | ||
| if !errors.Is(gotErr, errMarshal) { | ||
| t.Errorf("hook got error %v, want %v", gotErr, errMarshal) | ||
| } | ||
| if _, ok := indexed["bad"]; ok { | ||
| t.Error("unmarshalable resource made it into the snapshot") | ||
| } | ||
| if _, ok := indexed["good"]; !ok { | ||
| t.Error("the good resource was dropped alongside the bad one") | ||
| } | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
| @@ -0,0 +1,77 @@ | ||
| package cache | ||
|
|
||
| import ( | ||
| "context" | ||
| "errors" | ||
| "testing" | ||
|
|
||
| cluster "github.com/envoyproxy/go-control-plane/envoy/config/cluster/v3" | ||
| core "github.com/envoyproxy/go-control-plane/envoy/config/core/v3" | ||
| "github.com/envoyproxy/go-control-plane/pkg/cache/types" | ||
| "github.com/envoyproxy/go-control-plane/pkg/resource/v3" | ||
| "github.com/envoyproxy/go-control-plane/pkg/server/stream/v3" | ||
| ) | ||
|
|
||
| // Sending on a closed channel panics. The recover used to leave the return | ||
| // values at their zero values, so a closed channel reported (nil, nil) - which | ||
| // respondDeltaWatches reads as "no state change, keep the watch" and | ||
| // CreateDeltaWatch reads as "delayedResponse, register the watch". Both then | ||
| // hold a watch against a stream nobody will ever read from, and it looks | ||
| // exactly like a healthy idle watch. | ||
| func TestRespondDeltaReportsAClosedChannel(t *testing.T) { | ||
| snapshot, err := NewSnapshotWithTTLs("v1", map[resource.Type][]types.ResourceWithTTL{ | ||
| resource.ClusterType: {{Resource: &cluster.Cluster{Name: "c1"}, Version: "v1"}}, | ||
| }) | ||
| if err != nil { | ||
| t.Fatalf("NewSnapshot: %v", err) | ||
| } | ||
| if err := snapshot.ConstructVersionMap(); err != nil { | ||
| t.Fatalf("ConstructVersionMap: %v", err) | ||
| } | ||
|
|
||
| c := newSnapshotCache(false, IDHash{}, nil) | ||
|
|
||
| closed := make(chan DeltaResponse, 1) | ||
| close(closed) | ||
|
|
||
| req := &DeltaRequest{Node: &core.Node{Id: "node-1"}, TypeUrl: resource.ClusterType} | ||
| state := stream.NewStreamState(true, nil) | ||
|
|
||
| out, err := c.respondDelta(context.Background(), snapshot, req, closed, state) | ||
|
|
||
| if !errors.Is(err, ErrResponseChannelClosed) { | ||
| t.Errorf("err = %v, want ErrResponseChannelClosed", err) | ||
| } | ||
| if out != nil { | ||
| t.Errorf("response = %v, want nil", out) | ||
| } | ||
| } | ||
|
|
||
| // The other nil-response path must stay distinguishable: a full channel is | ||
| // retryable and deliberately reports (nil, nil) so the caller keeps the watch. | ||
| func TestRespondDeltaFullChannelStillReportsNoError(t *testing.T) { | ||
| snapshot, err := NewSnapshotWithTTLs("v1", map[resource.Type][]types.ResourceWithTTL{ | ||
| resource.ClusterType: {{Resource: &cluster.Cluster{Name: "c1"}, Version: "v1"}}, | ||
| }) | ||
| if err != nil { | ||
| t.Fatalf("NewSnapshot: %v", err) | ||
| } | ||
| if err := snapshot.ConstructVersionMap(); err != nil { | ||
| t.Fatalf("ConstructVersionMap: %v", err) | ||
| } | ||
|
|
||
| c := newSnapshotCache(false, IDHash{}, nil) | ||
|
|
||
| full := make(chan DeltaResponse) // unbuffered, nobody receiving | ||
| req := &DeltaRequest{Node: &core.Node{Id: "node-1"}, TypeUrl: resource.ClusterType} | ||
| state := stream.NewStreamState(true, nil) | ||
|
|
||
| out, err := c.respondDelta(context.Background(), snapshot, req, full, state) | ||
|
|
||
| if err != nil { | ||
| t.Errorf("err = %v, want nil for a full channel", err) | ||
| } | ||
| if out != nil { | ||
| t.Errorf("response = %v, want nil for a full channel", out) | ||
| } | ||
| } |
| Original file line number | Diff line number | Diff line change |
|---|---|---|
|
|
@@ -395,7 +395,7 @@ func (cache *snapshotCache) BatchUpsertResources(ctx context.Context, typ string | |
| for name, r := range resourcesUpserted { | ||
| out, err := r.Resource.MarshalVTStrict() | ||
| if err != nil { | ||
| fmt.Printf("failed to MarshalVTStrict resource %s: %v\n", name, err) | ||
| reportMarshalError(cache.log, typ, name, err) | ||
| continue | ||
| } | ||
| currentResources.Items[name] = VTMarshaledResource{ | ||
|
|
@@ -433,7 +433,7 @@ func (cache *snapshotCache) BatchUpsertResources(ctx context.Context, typ string | |
| wg.Wait() | ||
| elapsed := time.Since(start) | ||
| if elapsed > 50*time.Millisecond { | ||
| fmt.Printf("BatchUpsertResources took %s for %d nodes\n", elapsed, size) | ||
| cache.log.Debugf("BatchUpsertResources took %s for %d nodes", elapsed, size) | ||
| } | ||
| return nil | ||
| } | ||
|
|
@@ -457,7 +457,7 @@ func (cache *snapshotCache) UpsertResources(ctx context.Context, node string, ty | |
| return err | ||
| } | ||
| elapsed := time.Since(start) | ||
| fmt.Printf("UpsertResources took %s\n", elapsed) | ||
| cache.log.Debugf("UpsertResources took %s", elapsed) | ||
| return nil | ||
| } | ||
|
|
||
|
|
@@ -476,7 +476,7 @@ func (cache *snapshotCache) UpsertResources(ctx context.Context, node string, ty | |
| for name, r := range resourcesUpserted { | ||
| out, err := r.Resource.MarshalVTStrict() | ||
| if err != nil { | ||
| fmt.Printf("failed to MarshalVTStrict resource %s: %v\n", name, err) | ||
| reportMarshalError(cache.log, typ, name, err) | ||
| continue | ||
| } | ||
| currentResources.Items[name] = VTMarshaledResource{ | ||
|
|
@@ -722,7 +722,7 @@ func (cache *snapshotCache) respondDeltaWatches(ctx context.Context, info *statu | |
| elapsed := time.Since(start) | ||
| // Log if it takes more than 5ms | ||
| if elapsed > 5*time.Millisecond { | ||
| fmt.Printf("respondDelta took %s\n", elapsed) | ||
| cache.log.Debugf("respondDelta took %s", elapsed) | ||
| } | ||
| if err != nil { | ||
| return | ||
|
|
@@ -743,7 +743,7 @@ func (cache *snapshotCache) respondDeltaWatches(ctx context.Context, info *statu | |
| } | ||
| elapsed := time.Since(start) | ||
| if elapsed > 50*time.Millisecond { | ||
| fmt.Printf("respondDeltaWatches took %s for %d watches and node %s\n", elapsed, deletedCount, info.node.Id) | ||
| cache.log.Debugf("respondDeltaWatches took %s for %d watches and node %s", elapsed, deletedCount, info.node.Id) | ||
| } | ||
| } else { | ||
| for id, watch := range info.deltaWatches { | ||
|
|
@@ -1022,14 +1022,14 @@ func GetEnvoyNodeStr(node *core.Node) string { | |
| } | ||
|
|
||
| // Respond to a delta watch with the provided snapshot value. If the response is nil, there has been no state change. | ||
| func (cache *snapshotCache) respondDelta(ctx context.Context, snapshot ResourceSnapshot, request *DeltaRequest, value chan DeltaResponse, state stream.StreamState) (*RawDeltaResponse, error) { | ||
| func (cache *snapshotCache) respondDelta(ctx context.Context, snapshot ResourceSnapshot, request *DeltaRequest, value chan DeltaResponse, state stream.StreamState) (out *RawDeltaResponse, err error) { | ||
| // Use snapshot.Mu.RLock() to ensure that the snapshot is not modified while we are reading it. | ||
| // Previously, we created copy of resources which was less efficient. | ||
| start := time.Now() | ||
| snapshot.(*Snapshot).Mu.RLock() | ||
| elapsedLock := time.Since(start) | ||
| if elapsedLock > 1*time.Millisecond { | ||
| fmt.Printf("respondDelta took %s to lock\n", elapsedLock) | ||
| cache.log.Debugf("respondDelta waited %s for the snapshot read lock", elapsedLock) | ||
| } | ||
| resp := createDeltaResponse(ctx, request, state, resourceContainer{ | ||
| resourceMap: snapshot.GetResourcesAndTTL(request.GetTypeUrl()), | ||
|
|
@@ -1039,7 +1039,7 @@ func (cache *snapshotCache) respondDelta(ctx context.Context, snapshot ResourceS | |
| snapshot.(*Snapshot).Mu.RUnlock() | ||
| elapsed := time.Since(start) | ||
| if elapsed > 3*time.Millisecond { | ||
| fmt.Printf("createDeltaResponse took %s\n", elapsed) | ||
| cache.log.Debugf("createDeltaResponse took %s", elapsed) | ||
| } | ||
|
|
||
| // Only send a response if there were changes | ||
|
|
@@ -1052,9 +1052,16 @@ func (cache *snapshotCache) respondDelta(ctx context.Context, snapshot ResourceS | |
| request.GetNode().GetId(), request.GetTypeUrl(), GetResourceWithTTLNames(resp.Resources), resp.RemovedResources, state.IsWildcard()) | ||
| } | ||
|
|
||
| // A send on a closed channel panics. Recovering without setting the | ||
| // return values reported (nil, nil), which both callers read as "no | ||
| // state change, keep the watch" - so a watch on a dead stream looked | ||
| // exactly like a healthy idle one. Name the returns and report it. | ||
| defer func() { | ||
| if r := recover(); r != nil { | ||
| fmt.Println("Tried to send on a closed channel") | ||
| out = nil | ||
| err = ErrResponseChannelClosed | ||
| cache.log.Errorf("delta response send panicked for node %q type %s: %v; the stream is gone", | ||
| request.GetNode().GetId(), request.GetTypeUrl(), r) | ||
|
Comment on lines
+1061
to
+1064
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. 🩺 Stability & Availability | 🟠 Major | ⚡ Quick win 🔎 Supported by static analysis🏁 Script executed: rg -n -C 8 'CreateDeltaWatch|respondDeltaWatches|respondDelta\\(' pkg/cache/v3/simple.go
sed -n '1000,1080p' pkg/cache/v3/simple.go
rg -n 'ErrResponseChannelClosed|respondDeltaWatches|CreateDeltaWatch' pkg/cache/v3Repository: ShareChat/go-control-plane Length of output: 10108 🏁 Script executed: #!/bin/bash
sed -n '650,770p' pkg/cache/v3/simple.go
sed -n '920,1015p' pkg/cache/v3/simple.go
sed -n '1,130p' pkg/cache/v3/respond_delta_internal_test.go
sed -n '1,90p' pkg/cache/v3/failures.goRepository: ShareChat/go-control-plane Length of output: 11402 Remove watches after When Treat this error as terminal. Do not register a new watch, delete the affected existing watch, and continue processing other non-ADS watches. Suggested fix response, err := cache.respondDelta(context.Background(), snapshot, request, value, state)
if err != nil {
cache.log.Errorf("failed to respond with delta response: %s", err)
+ if err == ErrResponseChannelClosed {
+ return nil, false
+ }
}
delayedResponse = response == nil if err != nil {
+ if err == ErrResponseChannelClosed {
+ toDeleteCh <- k.ID
+ }
return
} if err != nil {
+ if err == ErrResponseChannelClosed {
+ delete(info.deltaWatches, id)
+ continue
+ }
return err
}🤖 Prompt for AI Agents |
||
| } | ||
| }() | ||
| // Non-blocking. Callers hold info.mu across this (respondDeltaWatches | ||
|
|
||
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
🎯 Functional Correctness | 🟡 Minor | ⚡ Quick win
🔎 Supported by static analysis
🏁 Script executed:
Repository: ShareChat/go-control-plane
Length of output: 9288
🏁 Script executed:
Repository: ShareChat/go-control-plane
Length of output: 42015
Thread the snapshot type URL into marshal-error reporting.
NewSnapshotandNewSnapshotWithTTLsreceivetypfrom a map keyed by resource type URL, but they discard it before callingIndexAndMarshalResourcesByName. A marshal failure then reportstypeURL == "". Because the free function has no logger,OnResourceMarshalErroris the only report and cannot classify the failure by type URL.Pass
typthrough an internal type-aware helper before reporting. Keep the existing public helper signatures compatible. Batch and single upserts already passtyp.🤖 Prompt for AI Agents