Redundant event prefix check before event watcher is created. (#32748)

* Revert "Error when redundant prefixes are detected in events. (#32652)"

This reverts commit 84dbc45e1d.

* 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 <edoardo.spadolini@goteleport.com>

---------

Co-authored-by: Edoardo Spadolini <edoardo.spadolini@goteleport.com>
This commit is contained in:
Michael Wilson
2023-09-29 01:22:08 +00:00
committed by GitHub
co-authored by Edoardo Spadolini
parent 2c53f04310
commit e93d165ed1
5 changed files with 16 additions and 114 deletions
-3
View File
@@ -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
+6 -12
View File
@@ -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 {
+1 -1
View File
@@ -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))
}
}
+9 -16
View File
@@ -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,
-82
View File
@@ -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))
})
}
}