Skip to content

Commit b966209

Browse files
authored
#1008 - nil object in watch event in causes api-server to report errors (#399)
* Remove error watch event when obj is nil * Add test coverage for watcher manager
1 parent 0187ea6 commit b966209

5 files changed

Lines changed: 185 additions & 12 deletions

File tree

pkg/engine/watchermanager.go

Lines changed: 17 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,19 @@ type ObjectWatcher interface {
3333
OnPackageRevisionChange(eventType watch.EventType, obj repository.PackageRevision) bool
3434
}
3535

36+
func isContextError(err error) bool {
37+
return err == context.Canceled || err == context.DeadlineExceeded
38+
}
39+
40+
func logWatcherAction(watcher *watcher, action string, err error) string {
41+
if isContextError(err) {
42+
klog.V(3).Infof("watcher %p %s: %v", watcher, action, err)
43+
return "debug"
44+
}
45+
klog.Infof("watcher %p %s: %v", watcher, action, err)
46+
return "info"
47+
}
48+
3649
func NewWatcherManager() *watcherManager {
3750
return &watcherManager{}
3851
}
@@ -81,10 +94,10 @@ func (r *watcherManager) WatchPackageRevisions(ctx context.Context, filter repos
8194
r.watchers[i] = w
8295
inserted = true
8396
active += 1
84-
klog.Infof("watcher %p finished with: %v and is replaced by watcher %p", watcher, err, w)
97+
logWatcherAction(watcher, "finished and is replaced by watcher", err)
8598
} else {
8699
r.watchers[i] = nil
87-
klog.Infof("watcher %p finished with: %v and is removed", watcher, err)
100+
logWatcherAction(watcher, "finished and is removed", err)
88101
}
89102
} else {
90103
active += 1
@@ -102,7 +115,7 @@ func (r *watcherManager) WatchPackageRevisions(ctx context.Context, filter repos
102115
r.watchers = append(r.watchers, w)
103116
}
104117

105-
klog.Infof("added watcher %p; there are now %d active watchers and %d slots", w, active, len(r.watchers))
118+
klog.V(3).Infof("added watcher %p; there are now %d active watchers and %d slots", w, active, len(r.watchers))
106119
return nil
107120
}
108121

@@ -117,7 +130,7 @@ func (r *watcherManager) NotifyPackageRevisionChange(eventType watch.EventType,
117130
continue
118131
}
119132
if err := watcher.isDoneFunction(); err != nil {
120-
klog.Infof("stopping watcher in response to error %v", err)
133+
logWatcherAction(watcher, "stopping in response to error", err)
121134
r.watchers[i] = nil
122135
continue
123136
}

pkg/engine/watchermanager_test.go

Lines changed: 40 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ package engine
22

33
import (
44
"context"
5+
"errors"
56
"testing"
67

78
"github.com/nephio-project/porch/pkg/externalrepo/fake"
@@ -88,3 +89,42 @@ func countActiveWatchers(manager *watcherManager) (int, int) {
8889
}
8990
return active, len(manager.watchers)
9091
}
92+
93+
func TestIsContextError(t *testing.T) {
94+
tests := []struct {
95+
name string
96+
err error
97+
expected bool
98+
}{
99+
{"context canceled", context.Canceled, true},
100+
{"context deadline", context.DeadlineExceeded, true},
101+
{"regular error", errors.New("test error"), false},
102+
}
103+
104+
for _, tt := range tests {
105+
t.Run(tt.name, func(t *testing.T) {
106+
result := isContextError(tt.err)
107+
assert.Equal(t, tt.expected, result)
108+
})
109+
}
110+
}
111+
112+
func TestLogWatcherAction(t *testing.T) {
113+
tests := []struct {
114+
name string
115+
err error
116+
expected string
117+
}{
118+
{"context canceled", context.Canceled, "debug"},
119+
{"context deadline", context.DeadlineExceeded, "debug"},
120+
{"regular error", errors.New("test error"), "info"},
121+
}
122+
123+
for _, tt := range tests {
124+
t.Run(tt.name, func(t *testing.T) {
125+
w := &watcher{}
126+
result := logWatcherAction(w, "test", tt.err)
127+
assert.Equal(t, tt.expected, result)
128+
})
129+
}
130+
}

pkg/registry/porch/background.go

Lines changed: 3 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -140,7 +140,7 @@ loop:
140140
} else if repository, ok := event.Object.(*configapi.Repository); ok {
141141
if event.Type == watch.Bookmark {
142142
bookmark = repository.ResourceVersion
143-
klog.Infof("Bookmark: %q", bookmark)
143+
klog.V(2).Infof("Bookmark: %q", bookmark)
144144
} else {
145145
if err := b.updateCache(ctx, event.Type, repository); err != nil {
146146
klog.Warningf("error updating cache: %v", err)
@@ -151,7 +151,7 @@ loop:
151151
}
152152

153153
case t := <-ticker.C:
154-
klog.Infof("Background task %s", t)
154+
klog.V(2).Infof("Background task %s", t)
155155
if err := b.runOnce(ctx); err != nil {
156156
klog.Errorf("Periodic repository refresh failed: %v", err)
157157
}
@@ -227,7 +227,7 @@ func (b *background) handleRepositoryEvent(ctx context.Context, repo *configapi.
227227
}
228228

229229
func (b *background) runOnce(ctx context.Context) error {
230-
klog.Infof("background-refreshing repositories")
230+
klog.V(2).Infof("background-refreshing repositories")
231231
repositories := &configapi.RepositoryList{}
232232
if err := b.coreClient.List(ctx, repositories); err != nil {
233233
return fmt.Errorf("error listing repository objects: %w", err)

pkg/registry/porch/watch.go

Lines changed: 15 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -122,11 +122,12 @@ func (w *watcher) listAndWatch(ctx context.Context, r packageReader, filter repo
122122

123123
if err := w.listAndWatchInner(ctx, r, filter); err != nil {
124124
// TODO: We need to populate the object on this error
125-
klog.Warningf("sending error to watch stream: %v", err)
126-
ev := watch.Event{
127-
Type: watch.Error,
125+
if err == context.Canceled || err == context.DeadlineExceeded {
126+
klog.V(3).Infof("sending error to watch stream: %v", err)
127+
} else {
128+
klog.Warningf("sending error to watch stream: %v", err)
128129
}
129-
w.resultChan <- ev
130+
// Don't send error events with nil objects
130131
}
131132
w.cancel()
132133
close(w.resultChan)
@@ -149,6 +150,9 @@ func (w *watcher) listAndWatchInner(ctx context.Context, r packageReader, filter
149150
errorResult <- err
150151
return false
151152
}
153+
if obj == nil {
154+
return true
155+
}
152156

153157
backlog = append(backlog, watch.Event{
154158
Type: eventType,
@@ -173,6 +177,9 @@ func (w *watcher) listAndWatchInner(ctx context.Context, r packageReader, filter
173177
w.mutex.Unlock()
174178
return err
175179
}
180+
if obj == nil {
181+
return nil
182+
}
176183
// TODO: Check resource version?
177184
ev := watch.Event{
178185
Type: watch.Added,
@@ -217,7 +224,7 @@ func (w *watcher) listAndWatchInner(ctx context.Context, r packageReader, filter
217224
w.sendWatchEvent(ev)
218225
}
219226

220-
klog.Infof("watch %p: moving watch into streaming mode after sentAdd %d, sentBacklog %d, sentNewBacklog %d", w, sentAdd, sentBacklog, sentNewBacklog)
227+
klog.V(3).Infof("watch %p: moving watch into streaming mode after sentAdd %d, sentBacklog %d, sentNewBacklog %d", w, sentAdd, sentBacklog, sentNewBacklog)
221228
w.eventCallback = func(eventType watch.EventType, pr repository.PackageRevision) bool {
222229
if w.done {
223230
return false
@@ -228,6 +235,9 @@ func (w *watcher) listAndWatchInner(ctx context.Context, r packageReader, filter
228235
errorResult <- err
229236
return false
230237
}
238+
if obj == nil {
239+
return true
240+
}
231241
// TODO: Check resource version?
232242
ev := watch.Event{
233243
Type: eventType,

pkg/registry/porch/watch_test.go

Lines changed: 110 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -25,6 +25,7 @@ import (
2525
"github.com/nephio-project/porch/pkg/externalrepo/fake"
2626
"github.com/nephio-project/porch/pkg/repository"
2727
metav1 "k8s.io/apimachinery/pkg/apis/meta/v1"
28+
"k8s.io/apimachinery/pkg/runtime"
2829
"k8s.io/apimachinery/pkg/watch"
2930
)
3031

@@ -97,3 +98,112 @@ func (f *fakePackageReader) watchPackages(ctx context.Context, filter repository
9798
func (f *fakePackageReader) listPackageRevisions(ctx context.Context, filter repository.ListPackageRevisionFilter, callback func(ctx context.Context, p repository.PackageRevision) error) error {
9899
return nil
99100
}
101+
func TestWatcherNilObject(t *testing.T) {
102+
tests := []struct {
103+
name string
104+
packages []repository.PackageRevision
105+
waitForStreaming bool
106+
sendEvent bool
107+
}{
108+
{
109+
name: "backlog phase",
110+
packages: nil,
111+
sendEvent: false,
112+
},
113+
{
114+
name: "streaming phase",
115+
packages: nil,
116+
waitForStreaming: true,
117+
sendEvent: true,
118+
},
119+
{
120+
name: "list phase",
121+
packages: []repository.PackageRevision{
122+
&fake.FakePackageRevision{
123+
PackageRevision: &porchapi.PackageRevision{
124+
ObjectMeta: metav1.ObjectMeta{
125+
Labels: make(map[string]string),
126+
},
127+
},
128+
},
129+
},
130+
},
131+
}
132+
133+
for _, tt := range tests {
134+
t.Run(tt.name, func(t *testing.T) {
135+
ctx, cancelFunc := context.WithCancel(context.Background())
136+
defer cancelFunc()
137+
138+
w := &watcher{
139+
cancel: cancelFunc,
140+
resultChan: make(chan watch.Event, 64),
141+
extractor: func(ctx context.Context, pr repository.PackageRevision) (runtime.Object, error) {
142+
return nil, nil
143+
},
144+
}
145+
146+
r := &nilCheckFakeReader{
147+
packages: tt.packages,
148+
sendEventInBacklog: tt.name == "backlog phase",
149+
}
150+
r.Add(1)
151+
var filter repository.ListPackageRevisionFilter
152+
go w.listAndWatch(ctx, r, filter)
153+
r.Wait()
154+
155+
if tt.waitForStreaming {
156+
time.Sleep(100 * time.Millisecond)
157+
}
158+
159+
if tt.sendEvent {
160+
pkgRev := &fake.FakePackageRevision{
161+
PackageRevision: &porchapi.PackageRevision{
162+
ObjectMeta: metav1.ObjectMeta{
163+
Labels: make(map[string]string),
164+
},
165+
},
166+
}
167+
cont := r.callback.OnPackageRevisionChange(watch.Modified, pkgRev)
168+
if !cont {
169+
t.Error("Expected callback to return true for nil object")
170+
}
171+
} else {
172+
time.Sleep(10 * time.Millisecond)
173+
}
174+
})
175+
}
176+
}
177+
178+
type nilCheckFakeReader struct {
179+
sync.WaitGroup
180+
callback engine.ObjectWatcher
181+
packages []repository.PackageRevision
182+
sendEventInBacklog bool
183+
}
184+
185+
func (f *nilCheckFakeReader) watchPackages(ctx context.Context, filter repository.ListPackageRevisionFilter, callback engine.ObjectWatcher) error {
186+
f.callback = callback
187+
if f.sendEventInBacklog {
188+
// Send event synchronously in backlog phase
189+
pkgRev := &fake.FakePackageRevision{
190+
PackageRevision: &porchapi.PackageRevision{
191+
ObjectMeta: metav1.ObjectMeta{
192+
Labels: make(map[string]string),
193+
},
194+
},
195+
}
196+
callback.OnPackageRevisionChange(watch.Modified, pkgRev)
197+
}
198+
f.Done()
199+
return nil
200+
}
201+
202+
func (f *nilCheckFakeReader) listPackageRevisions(ctx context.Context, filter repository.ListPackageRevisionFilter, callback func(ctx context.Context, p repository.PackageRevision) error) error {
203+
for _, pkg := range f.packages {
204+
if err := callback(ctx, pkg); err != nil {
205+
return err
206+
}
207+
}
208+
return nil
209+
}

0 commit comments

Comments
 (0)