mirror of
https://github.com/seaweedfs/seaweedfs.git
synced 2026-10-10 08:17:44 +02:00
kafka: honor notification.kafka.event_types (#11601)
* kafka: honor notification.kafka.event_types Kafka published every filer event and ignored the filter the webhook notifier already uses. Co-authored-by: Cursor <cursoragent@cursor.com> * notification: share event-type classification between queues Kafka duplicated the webhook's event classification verbatim; move it to the notification package so the two queues cannot drift. Webhook keeps its typed eventType wrappers over the shared helpers. Generated with [Devin](https://devin.ai) Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com> --------- Co-authored-by: Cursor <cursoragent@cursor.com> Co-authored-by: Chris Lu <chrislusf@users.noreply.github.com> Co-authored-by: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
This commit is contained in:
6 files changed
+272
-41
No files matched your search
@@ -20,6 +20,7 @@ hosts = [
|
||||
"localhost:9092"
|
||||
]
|
||||
topic = "seaweedfs_filer"
|
||||
# event_types = ["create", "update", "delete", "rename"] # optional: filter by event types (default: all)
|
||||
offsetFile = "./last.offset"
|
||||
offsetSaveIntervalSeconds = 10
|
||||
# SASL Authentication
|
||||
|
||||
@@ -0,0 +1,53 @@
|
||||
package notification
|
||||
|
||||
import (
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
)
|
||||
|
||||
const (
|
||||
EventTypeCreate = "create"
|
||||
EventTypeDelete = "delete"
|
||||
EventTypeUpdate = "update"
|
||||
EventTypeRename = "rename"
|
||||
)
|
||||
|
||||
// ValidEventType reports whether t is an event type a queue can filter on.
|
||||
func ValidEventType(t string) bool {
|
||||
switch t {
|
||||
case EventTypeCreate, EventTypeDelete, EventTypeUpdate, EventTypeRename:
|
||||
return true
|
||||
default:
|
||||
return false
|
||||
}
|
||||
}
|
||||
|
||||
// DetectEventType classifies an entry-change notification. key is the old
|
||||
// entry's path when one exists.
|
||||
func DetectEventType(key string, notification *filer_pb.EventNotification) string {
|
||||
hasOldEntry := notification.OldEntry != nil
|
||||
hasNewEntry := notification.NewEntry != nil
|
||||
|
||||
if !hasOldEntry && hasNewEntry {
|
||||
return EventTypeCreate
|
||||
}
|
||||
|
||||
if hasOldEntry && !hasNewEntry {
|
||||
return EventTypeDelete
|
||||
}
|
||||
|
||||
if hasOldEntry && hasNewEntry {
|
||||
oldDir, _ := util.FullPath(key).DirAndName()
|
||||
newDir := notification.NewParentPath
|
||||
if newDir == "" {
|
||||
newDir = oldDir
|
||||
}
|
||||
if oldDir != newDir || notification.OldEntry.Name != notification.NewEntry.Name {
|
||||
return EventTypeRename
|
||||
}
|
||||
|
||||
return EventTypeUpdate
|
||||
}
|
||||
|
||||
return EventTypeUpdate
|
||||
}
|
||||
@@ -0,0 +1,40 @@
|
||||
package kafka
|
||||
|
||||
import (
|
||||
"github.com/seaweedfs/seaweedfs/weed/glog"
|
||||
"github.com/seaweedfs/seaweedfs/weed/notification"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"google.golang.org/protobuf/proto"
|
||||
)
|
||||
|
||||
// Empty eventTypes means publish every event. A non-nil map publishes only those types.
|
||||
// The names and classification match the webhook notifier.
|
||||
func (k *KafkaQueue) setEventTypes(types []string) {
|
||||
if len(types) == 0 {
|
||||
k.eventTypes = nil
|
||||
return
|
||||
}
|
||||
|
||||
allowed := make(map[string]struct{}, len(types))
|
||||
for _, et := range types {
|
||||
if !notification.ValidEventType(et) {
|
||||
glog.Warningf("invalid event type: %v", et)
|
||||
continue
|
||||
}
|
||||
allowed[et] = struct{}{}
|
||||
}
|
||||
k.eventTypes = allowed
|
||||
}
|
||||
|
||||
func (k *KafkaQueue) allowsEvent(key string, message proto.Message) bool {
|
||||
if k.eventTypes == nil {
|
||||
return true
|
||||
}
|
||||
|
||||
n, ok := message.(*filer_pb.EventNotification)
|
||||
if !ok || n == nil {
|
||||
return false
|
||||
}
|
||||
_, allowed := k.eventTypes[notification.DetectEventType(key, n)]
|
||||
return allowed
|
||||
}
|
||||
@@ -0,0 +1,160 @@
|
||||
package kafka
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/notification"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"google.golang.org/protobuf/proto"
|
||||
)
|
||||
|
||||
func TestKafkaEventTypes(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
key string
|
||||
eventTypes []string
|
||||
notification *filer_pb.EventNotification
|
||||
wantType string
|
||||
wantAllow bool
|
||||
}{
|
||||
{
|
||||
name: "create event allowed",
|
||||
key: "/test/test.txt",
|
||||
eventTypes: []string{"create", "delete"},
|
||||
notification: &filer_pb.EventNotification{
|
||||
NewEntry: &filer_pb.Entry{Name: "test.txt"},
|
||||
},
|
||||
wantType: "create",
|
||||
wantAllow: true,
|
||||
},
|
||||
{
|
||||
name: "create event filtered out",
|
||||
key: "/test/test.txt",
|
||||
eventTypes: []string{"delete", "update"},
|
||||
notification: &filer_pb.EventNotification{
|
||||
NewEntry: &filer_pb.Entry{Name: "test.txt"},
|
||||
},
|
||||
wantType: "create",
|
||||
wantAllow: false,
|
||||
},
|
||||
{
|
||||
name: "delete event allowed",
|
||||
key: "/test/test.txt",
|
||||
eventTypes: []string{"create", "delete"},
|
||||
notification: &filer_pb.EventNotification{
|
||||
OldEntry: &filer_pb.Entry{Name: "test.txt"},
|
||||
},
|
||||
wantType: "delete",
|
||||
wantAllow: true,
|
||||
},
|
||||
{
|
||||
name: "update event allowed",
|
||||
key: "/test/test.txt",
|
||||
eventTypes: []string{"update"},
|
||||
notification: &filer_pb.EventNotification{
|
||||
OldEntry: &filer_pb.Entry{Name: "test.txt"},
|
||||
NewEntry: &filer_pb.Entry{Name: "test.txt"},
|
||||
NewParentPath: "/test",
|
||||
},
|
||||
wantType: "update",
|
||||
wantAllow: true,
|
||||
},
|
||||
{
|
||||
name: "rename event allowed",
|
||||
key: "/old/path/old.txt",
|
||||
eventTypes: []string{"rename"},
|
||||
notification: &filer_pb.EventNotification{
|
||||
OldEntry: &filer_pb.Entry{Name: "old.txt"},
|
||||
NewEntry: &filer_pb.Entry{Name: "new.txt"},
|
||||
NewParentPath: "/new/path",
|
||||
},
|
||||
wantType: "rename",
|
||||
wantAllow: true,
|
||||
},
|
||||
{
|
||||
name: "rename across directories allowed",
|
||||
key: "/old/path/file.txt",
|
||||
eventTypes: []string{"rename"},
|
||||
notification: &filer_pb.EventNotification{
|
||||
OldEntry: &filer_pb.Entry{Name: "file.txt"},
|
||||
NewEntry: &filer_pb.Entry{Name: "file.txt"},
|
||||
NewParentPath: "/new/path",
|
||||
},
|
||||
wantType: "rename",
|
||||
wantAllow: true,
|
||||
},
|
||||
{
|
||||
name: "rename event filtered out",
|
||||
key: "/old/path/old.txt",
|
||||
eventTypes: []string{"create", "delete", "update"},
|
||||
notification: &filer_pb.EventNotification{
|
||||
OldEntry: &filer_pb.Entry{Name: "old.txt"},
|
||||
NewEntry: &filer_pb.Entry{Name: "new.txt"},
|
||||
NewParentPath: "/new/path",
|
||||
},
|
||||
wantType: "rename",
|
||||
wantAllow: false,
|
||||
},
|
||||
{
|
||||
name: "empty filter publishes every event",
|
||||
key: "/test/test.txt",
|
||||
eventTypes: []string{},
|
||||
notification: &filer_pb.EventNotification{
|
||||
NewEntry: &filer_pb.Entry{Name: "test.txt"},
|
||||
},
|
||||
wantType: "create",
|
||||
wantAllow: true,
|
||||
},
|
||||
{
|
||||
name: "invalid type does not remove a valid one",
|
||||
key: "/test/test.txt",
|
||||
eventTypes: []string{"create", "bogus"},
|
||||
notification: &filer_pb.EventNotification{
|
||||
NewEntry: &filer_pb.Entry{Name: "test.txt"},
|
||||
},
|
||||
wantType: "create",
|
||||
wantAllow: true,
|
||||
},
|
||||
{
|
||||
name: "only invalid types publish nothing",
|
||||
key: "/test/test.txt",
|
||||
eventTypes: []string{"bogus"},
|
||||
notification: &filer_pb.EventNotification{
|
||||
NewEntry: &filer_pb.Entry{Name: "test.txt"},
|
||||
},
|
||||
wantType: "create",
|
||||
wantAllow: false,
|
||||
},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
gotType := notification.DetectEventType(tt.key, tt.notification)
|
||||
if gotType != tt.wantType {
|
||||
t.Errorf("DetectEventType() = %v, want %v", gotType, tt.wantType)
|
||||
}
|
||||
|
||||
q := &KafkaQueue{}
|
||||
q.setEventTypes(tt.eventTypes)
|
||||
if got := q.allowsEvent(tt.key, tt.notification); got != tt.wantAllow {
|
||||
t.Errorf("allowsEvent() = %v, want %v", got, tt.wantAllow)
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestKafkaEventFilterUnsetPublishesOtherMessages(t *testing.T) {
|
||||
q := &KafkaQueue{}
|
||||
if !q.allowsEvent("/test/test.txt", &filer_pb.Entry{Name: "test.txt"}) {
|
||||
t.Fatal("unset filter dropped a message")
|
||||
}
|
||||
}
|
||||
|
||||
func TestKafkaEventFilterDropsUnclassifiedMessages(t *testing.T) {
|
||||
q := &KafkaQueue{}
|
||||
q.setEventTypes([]string{"create"})
|
||||
var message proto.Message = &filer_pb.Entry{Name: "test.txt"}
|
||||
if q.allowsEvent("/test/test.txt", message) {
|
||||
t.Fatal("filtered queue published a message that is not an event notification")
|
||||
}
|
||||
}
|
||||
@@ -15,8 +15,9 @@ func init() {
|
||||
}
|
||||
|
||||
type KafkaQueue struct {
|
||||
topic string
|
||||
producer sarama.AsyncProducer
|
||||
topic string
|
||||
producer sarama.AsyncProducer
|
||||
eventTypes map[string]struct{}
|
||||
}
|
||||
|
||||
func (k *KafkaQueue) GetName() string {
|
||||
@@ -26,6 +27,9 @@ func (k *KafkaQueue) GetName() string {
|
||||
func (k *KafkaQueue) Initialize(configuration util.Configuration, prefix string) (err error) {
|
||||
glog.V(0).Infof("filer.notification.kafka.hosts: %v\n", configuration.GetStringSlice(prefix+"hosts"))
|
||||
glog.V(0).Infof("filer.notification.kafka.topic: %v\n", configuration.GetString(prefix+"topic"))
|
||||
eventTypes := configuration.GetStringSlice(prefix + "event_types")
|
||||
glog.V(0).Infof("filer.notification.kafka.event_types: %v\n", eventTypes)
|
||||
k.setEventTypes(eventTypes)
|
||||
return k.initialize(
|
||||
configuration.GetStringSlice(prefix+"hosts"),
|
||||
configuration.GetString(prefix+"topic"),
|
||||
@@ -63,6 +67,10 @@ func (k *KafkaQueue) initialize(hosts []string, topic string, saslTLS SASLTLSCon
|
||||
}
|
||||
|
||||
func (k *KafkaQueue) SendMessage(key string, message proto.Message) (err error) {
|
||||
if !k.allowsEvent(key, message) {
|
||||
return nil
|
||||
}
|
||||
|
||||
bytes, err := proto.Marshal(message)
|
||||
if err != nil {
|
||||
return
|
||||
|
||||
@@ -3,9 +3,9 @@ package webhook
|
||||
import (
|
||||
"fmt"
|
||||
"net/url"
|
||||
"slices"
|
||||
"strconv"
|
||||
|
||||
"github.com/seaweedfs/seaweedfs/weed/notification"
|
||||
"github.com/seaweedfs/seaweedfs/weed/pb/filer_pb"
|
||||
"github.com/seaweedfs/seaweedfs/weed/util"
|
||||
"google.golang.org/protobuf/proto"
|
||||
@@ -20,21 +20,14 @@ const (
|
||||
type eventType string
|
||||
|
||||
const (
|
||||
eventTypeCreate eventType = "create"
|
||||
eventTypeDelete eventType = "delete"
|
||||
eventTypeUpdate eventType = "update"
|
||||
eventTypeRename eventType = "rename"
|
||||
eventTypeCreate = eventType(notification.EventTypeCreate)
|
||||
eventTypeDelete = eventType(notification.EventTypeDelete)
|
||||
eventTypeUpdate = eventType(notification.EventTypeUpdate)
|
||||
eventTypeRename = eventType(notification.EventTypeRename)
|
||||
)
|
||||
|
||||
func (e eventType) valid() bool {
|
||||
return slices.Contains([]eventType{
|
||||
eventTypeCreate,
|
||||
eventTypeDelete,
|
||||
eventTypeUpdate,
|
||||
eventTypeRename,
|
||||
},
|
||||
e,
|
||||
)
|
||||
return notification.ValidEventType(string(e))
|
||||
}
|
||||
|
||||
var (
|
||||
@@ -157,30 +150,6 @@ func (c *config) validate() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func detectEventType(key string, notification *filer_pb.EventNotification) eventType {
|
||||
hasOldEntry := notification.OldEntry != nil
|
||||
hasNewEntry := notification.NewEntry != nil
|
||||
|
||||
if !hasOldEntry && hasNewEntry {
|
||||
return eventTypeCreate
|
||||
}
|
||||
|
||||
if hasOldEntry && !hasNewEntry {
|
||||
return eventTypeDelete
|
||||
}
|
||||
|
||||
if hasOldEntry && hasNewEntry {
|
||||
oldDir, _ := util.FullPath(key).DirAndName()
|
||||
newDir := notification.NewParentPath
|
||||
if newDir == "" {
|
||||
newDir = oldDir
|
||||
}
|
||||
if oldDir != newDir || notification.OldEntry.Name != notification.NewEntry.Name {
|
||||
return eventTypeRename
|
||||
}
|
||||
|
||||
return eventTypeUpdate
|
||||
}
|
||||
|
||||
return eventTypeUpdate
|
||||
func detectEventType(key string, n *filer_pb.EventNotification) eventType {
|
||||
return eventType(notification.DetectEventType(key, n))
|
||||
}
|
||||
Reference in new issue
Block a user