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)) - }) - } -}