diff --git a/weed/command/scaffold/notification.toml b/weed/command/scaffold/notification.toml index ca82f2c1e..ebbd240f9 100644 --- a/weed/command/scaffold/notification.toml +++ b/weed/command/scaffold/notification.toml @@ -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 diff --git a/weed/notification/event_type.go b/weed/notification/event_type.go new file mode 100644 index 000000000..c322bd5dd --- /dev/null +++ b/weed/notification/event_type.go @@ -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 +} diff --git a/weed/notification/kafka/event_filter.go b/weed/notification/kafka/event_filter.go new file mode 100644 index 000000000..45766ad99 --- /dev/null +++ b/weed/notification/kafka/event_filter.go @@ -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 +} diff --git a/weed/notification/kafka/event_filter_test.go b/weed/notification/kafka/event_filter_test.go new file mode 100644 index 000000000..025f6ff37 --- /dev/null +++ b/weed/notification/kafka/event_filter_test.go @@ -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") + } +} diff --git a/weed/notification/kafka/kafka_queue.go b/weed/notification/kafka/kafka_queue.go index 53f0802f6..b0922d5ed 100644 --- a/weed/notification/kafka/kafka_queue.go +++ b/weed/notification/kafka/kafka_queue.go @@ -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 diff --git a/weed/notification/webhook/types.go b/weed/notification/webhook/types.go index af1fb0242..f177f635a 100644 --- a/weed/notification/webhook/types.go +++ b/weed/notification/webhook/types.go @@ -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)) }