From e93d165ed19ee7582d489d5a90bb863f2dd9ff4b Mon Sep 17 00:00:00 2001 From: Michael Wilson Date: Thu, 28 Sep 2023 21:22:08 -0400 Subject: [PATCH] Redundant event prefix check before event watcher is created. (#32748) * Revert "Error when redundant prefixes are detected in events. (#32652)" This reverts commit 84dbc45e1d069cb139fddef79750356b1444c207. * Error when redundant prefixes are detected in events. When creating a new events watcher, redundant prefixes will be detected and produce an error. This should prevent developer mistakes where watched prefixes overlap, causing subsets of events not to be parsed. This has been verified manually. * Fix buffer_test.go * Update lib/services/local/events.go Co-authored-by: Edoardo Spadolini --------- Co-authored-by: Edoardo Spadolini --- lib/backend/backend.go | 3 -- lib/backend/buffer.go | 18 +++---- lib/backend/buffer_test.go | 2 +- lib/services/local/events.go | 25 ++++------ lib/services/local/events_test.go | 82 ------------------------------- 5 files changed, 16 insertions(+), 114 deletions(-) delete mode 100644 lib/services/local/events_test.go diff --git a/lib/backend/backend.go b/lib/backend/backend.go index 23d54b71869..4d2dbf461a5 100644 --- a/lib/backend/backend.go +++ b/lib/backend/backend.go @@ -202,9 +202,6 @@ type Watcher interface { // Close closes the watcher and releases // all associated resources Close() error - - // Prefixes returns the prefixes that the watcher is monitoring. - Prefixes() [][]byte } // GetResult provides the result of GetRange request diff --git a/lib/backend/buffer.go b/lib/backend/buffer.go index ae2c4243205..cb1519e3797 100644 --- a/lib/backend/buffer.go +++ b/lib/backend/buffer.go @@ -205,7 +205,8 @@ func (c *CircularBuffer) fanOutEvent(r Event) { } } -func removeRedundantPrefixes(prefixes [][]byte) [][]byte { +// RemoveRedundantPrefixes will remove redundant prefixes from the given prefix list. +func RemoveRedundantPrefixes(prefixes [][]byte) [][]byte { if len(prefixes) == 0 { return prefixes } @@ -247,7 +248,7 @@ func (c *CircularBuffer) NewWatcher(ctx context.Context, watch Watch) (Watcher, } else { // if watcher's prefixes are redundant, keep only shorter prefixes // to avoid double fan out - watch.Prefixes = removeRedundantPrefixes(watch.Prefixes) + watch.Prefixes = RemoveRedundantPrefixes(watch.Prefixes) } closeCtx, cancel := context.WithCancel(ctx) @@ -305,7 +306,7 @@ type BufferWatcher struct { // String returns user-friendly representation // of the buffer watcher func (w *BufferWatcher) String() string { - return fmt.Sprintf("Watcher(name=%v, prefixes=%v, capacity=%v, size=%v)", w.Name, string(bytes.Join(w.Watch.Prefixes, []byte(", "))), w.capacity, len(w.eventsC)) + return fmt.Sprintf("Watcher(name=%v, prefixes=%v, capacity=%v, size=%v)", w.Name, string(bytes.Join(w.Prefixes, []byte(", "))), w.capacity, len(w.eventsC)) } // Events returns events channel. This method performs internal work and should be re-called after each event @@ -326,13 +327,6 @@ func (w *BufferWatcher) Done() <-chan struct{} { return w.ctx.Done() } -// Prefixes returns the prefixes that the watcher is monitoring. -func (w *BufferWatcher) Prefixes() [][]byte { - prefixes := make([][]byte, len(w.Watch.Prefixes)) - copy(prefixes, w.Watch.Prefixes) - return prefixes -} - // flushBacklog attempts to push any backlogged events into the // event channel. returns true if backlog is empty. func (w *BufferWatcher) flushBacklog() (ok bool) { @@ -438,7 +432,7 @@ type watcherTree struct { // add adds buffer watcher to the tree func (t *watcherTree) add(w *BufferWatcher) { - for _, p := range w.Watch.Prefixes { + for _, p := range w.Prefixes { prefix := string(p) val, ok := t.Tree.Get(prefix) var watchers []*BufferWatcher @@ -456,7 +450,7 @@ func (t *watcherTree) rm(w *BufferWatcher) bool { return false } var found bool - for _, p := range w.Watch.Prefixes { + for _, p := range w.Prefixes { prefix := string(p) val, ok := t.Tree.Get(prefix) if !ok { diff --git a/lib/backend/buffer_test.go b/lib/backend/buffer_test.go index e295a83da5a..2af59ac8496 100644 --- a/lib/backend/buffer_test.go +++ b/lib/backend/buffer_test.go @@ -196,7 +196,7 @@ func TestRemoveRedundantPrefixes(t *testing.T) { }, } for _, tc := range tcs { - require.Empty(t, cmp.Diff(removeRedundantPrefixes(tc.in), tc.out)) + require.Empty(t, cmp.Diff(RemoveRedundantPrefixes(tc.in), tc.out)) } } diff --git a/lib/services/local/events.go b/lib/services/local/events.go index f4e2ed8bb7b..e3cb6394d99 100644 --- a/lib/services/local/events.go +++ b/lib/services/local/events.go @@ -192,6 +192,15 @@ func (e *EventsService) NewWatcher(ctx context.Context, watch types.Watch) (type return nil, trace.BadParameter("none of the requested kinds can be watched") } + origNumPrefixes := len(prefixes) + redundantNumPrefixes := len(backend.RemoveRedundantPrefixes(prefixes)) + if origNumPrefixes != redundantNumPrefixes { + // If you've hit this error, the prefixes in two or more of your parsers probably overlap, meaning + // one prefix will also contain another as a subset. Look into using backend.ExactKey instead of + // backend.Key in your parser. + return nil, trace.BadParameter("redundant prefixes detected in events, which will result in event parsers not aligning with their intended prefix (this is a bug)") + } + w, err := e.backend.NewWatcher(ctx, backend.Watch{ Name: watch.Name, Prefixes: prefixes, @@ -201,25 +210,9 @@ func (e *EventsService) NewWatcher(ctx context.Context, watch types.Watch) (type if err != nil { return nil, trace.Wrap(err) } - - if err := verifyEventWatcherPrefixes(prefixes, w); err != nil { - return nil, trace.Wrap(err) - } - return newWatcher(w, e.Entry, parsers, validKinds), nil } -// verifyEventWatcherPrefixes will ensure that the expected prefixes are all found within the watcher. -func verifyEventWatcherPrefixes(expectedPrefixes [][]byte, watcher backend.Watcher) error { - if len(expectedPrefixes) != len(watcher.Prefixes()) { - // If you've hit this error, the prefixes in two or more of your parsers probably overlap, meaning - // one prefix will also contain another as a subset. Look into using backend.ExactKey instead of - // backend.Key in your parser. - return trace.BadParameter("redundant prefixes detected in events, which will result in event parsers not aligning with their intended prefix") - } - return nil -} - func newWatcher(backendWatcher backend.Watcher, l *logrus.Entry, parsers []resourceParser, kinds []types.WatchKind) *watcher { w := &watcher{ backendWatcher: backendWatcher, diff --git a/lib/services/local/events_test.go b/lib/services/local/events_test.go deleted file mode 100644 index a79b9dcd224..00000000000 --- a/lib/services/local/events_test.go +++ /dev/null @@ -1,82 +0,0 @@ -/* -Copyright 2023 Gravitational, Inc. - -Licensed under the Apache License, Version 2.0 (the "License"); -you may not use this file except in compliance with the License. -You may obtain a copy of the License at - - http://www.apache.org/licenses/LICENSE-2.0 - -Unless required by applicable law or agreed to in writing, software -distributed under the License is distributed on an "AS IS" BASIS, -WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. -See the License for the specific language governing permissions and -limitations under the License. -*/ - -package local - -import ( - "context" - "testing" - - "github.com/gravitational/trace" - "github.com/stretchr/testify/require" - - "github.com/gravitational/teleport/lib/backend" - "github.com/gravitational/teleport/lib/backend/memory" -) - -func TestVerifyEventWatcherPrefxies(t *testing.T) { - t.Parallel() - - tests := []struct { - name string - expectedPrefixes [][]byte - assertErr require.ErrorAssertionFunc - }{ - { - name: "no overlap", - expectedPrefixes: [][]byte{ - backend.Key("one"), - backend.Key("two"), - backend.Key("three"), - }, - assertErr: require.NoError, - }, - { - name: "overlap", - expectedPrefixes: [][]byte{ - backend.Key("one"), - backend.Key("oneoverlap"), - backend.Key("two"), - backend.Key("three"), - }, - assertErr: func(t require.TestingT, err error, i ...interface{}) { - require.True(t, trace.IsBadParameter(err)) - }, - }, - } - - for _, test := range tests { - test := test - t.Run(test.name, func(t *testing.T) { - t.Parallel() - - ctx := context.Background() - - mem, err := memory.New(memory.Config{}) - require.NoError(t, err) - - w, err := mem.NewWatcher(ctx, backend.Watch{ - Name: "test-watcher", - Prefixes: test.expectedPrefixes, - QueueSize: 10, - MetricComponent: "component", - }) - require.NoError(t, err) - - test.assertErr(t, verifyEventWatcherPrefixes(test.expectedPrefixes, w)) - }) - } -}