server: implement filtered event subscriptions
Event file now supports streaming rdwr pattern: - Read: streams all events (unchanged) - Write filter, then read: streams only matching events Filter syntax: - * matches single segment - > matches rest of topic - Examples: session.*.agent.*.state, session.abc.> Usage: echo filter | rdwrs event Per-fid subscription state is cleaned up on clunk.
This commit is contained in:
parent
13831fb9f3
commit
772d77940c
|
|
@ -239,8 +239,8 @@ func buildTreeSpec(cfg *Config) virtfs.FsNodeDecl {
|
|||
}, data)
|
||||
}),
|
||||
),
|
||||
virtfs.FileNode("event", 0444,
|
||||
virtfs.Doc("Streaming read for server events"),
|
||||
virtfs.FileNode("event", 0666,
|
||||
virtfs.Doc("Event stream. Read: all events. Write filter then read: filtered events."),
|
||||
virtfs.Stream(eventRead, eventSignal),
|
||||
),
|
||||
virtfs.FileNode("generate", 0666,
|
||||
|
|
|
|||
|
|
@ -42,6 +42,57 @@ func SubscribeEvents(ctx context.Context, topic string) <-chan Event {
|
|||
return ch
|
||||
}
|
||||
|
||||
// SubscribeEventsFiltered returns a channel that receives events matching the filter pattern.
|
||||
// Filter syntax: * matches single segment, > matches rest of topic.
|
||||
// Examples: "session.*.agent.*.state", "session.abc.>", "*"
|
||||
func SubscribeEventsFiltered(ctx context.Context, filter string) <-chan Event {
|
||||
// Subscribe to all events and filter client-side for complex patterns
|
||||
// (pubsub only supports simple topic matching)
|
||||
ch := make(chan Event, 64)
|
||||
allEvents := SubscribeEvents(ctx, "*")
|
||||
go func() {
|
||||
defer close(ch)
|
||||
for ev := range allEvents {
|
||||
if MatchTopic(filter, ev.Topic) {
|
||||
select {
|
||||
case ch <- ev:
|
||||
default: // drop if full
|
||||
}
|
||||
}
|
||||
}
|
||||
}()
|
||||
return ch
|
||||
}
|
||||
|
||||
// MatchTopic checks if a topic matches a filter pattern.
|
||||
// Filter syntax: * matches single segment, > matches rest of topic.
|
||||
func MatchTopic(filter, topic string) bool {
|
||||
if filter == "*" || filter == ">" {
|
||||
return true
|
||||
}
|
||||
filterParts := strings.Split(filter, ".")
|
||||
topicParts := strings.Split(topic, ".")
|
||||
|
||||
fi, ti := 0, 0
|
||||
for fi < len(filterParts) && ti < len(topicParts) {
|
||||
fp := filterParts[fi]
|
||||
if fp == ">" {
|
||||
return true // > matches rest
|
||||
}
|
||||
if fp == "*" {
|
||||
fi++
|
||||
ti++
|
||||
continue // * matches single segment
|
||||
}
|
||||
if fp != topicParts[ti] {
|
||||
return false
|
||||
}
|
||||
fi++
|
||||
ti++
|
||||
}
|
||||
return fi == len(filterParts) && ti == len(topicParts)
|
||||
}
|
||||
|
||||
// EventStream returns functions suitable for virtfs.Stream - delivers events as they arrive.
|
||||
// Each read returns the next event; blocks until one is available.
|
||||
func EventStream(ch <-chan Event) (readFn func(base string) ([]byte, string, error), signalFn func() <-chan struct{}) {
|
||||
|
|
|
|||
|
|
@ -64,6 +64,11 @@ type fid struct {
|
|||
readCache []byte // cached content for offset-based paging
|
||||
waitBase string
|
||||
dirCache []byte
|
||||
|
||||
// Event stream subscription (for event file with filter)
|
||||
eventFilter string
|
||||
eventCh <-chan session.Event
|
||||
eventCancel context.CancelFunc
|
||||
}
|
||||
|
||||
type connState struct {
|
||||
|
|
@ -593,6 +598,54 @@ func read(root *fs.Tree, cs *connState, fc *plan9.Fcall, ctx context.Context) *p
|
|||
}
|
||||
|
||||
if entry.StreamMode() {
|
||||
// Special handling for event file with filter
|
||||
if path == "event" {
|
||||
cs.mu.RLock()
|
||||
f, fidOK := cs.fids[fc.Fid]
|
||||
cs.mu.RUnlock()
|
||||
if fidOK {
|
||||
f.mu.Lock()
|
||||
filter := f.eventFilter
|
||||
eventCh := f.eventCh
|
||||
f.mu.Unlock()
|
||||
|
||||
// If filter was written, use filtered subscription
|
||||
if filter != "" && eventCh == nil {
|
||||
// Create subscription on first read after write
|
||||
subCtx, cancel := context.WithCancel(ctx)
|
||||
ch := session.SubscribeEventsFiltered(subCtx, filter)
|
||||
f.mu.Lock()
|
||||
f.eventCh = ch
|
||||
f.eventCancel = cancel
|
||||
f.mu.Unlock()
|
||||
eventCh = ch
|
||||
}
|
||||
|
||||
if eventCh != nil {
|
||||
// Read from filtered subscription
|
||||
select {
|
||||
case ev, ok := <-eventCh:
|
||||
if !ok {
|
||||
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0}
|
||||
}
|
||||
line := ev.Topic
|
||||
if ev.Payload != "" {
|
||||
line += " " + ev.Payload
|
||||
}
|
||||
content := []byte(line + "\n")
|
||||
count := int(fc.Count)
|
||||
if count > len(content) {
|
||||
count = len(content)
|
||||
}
|
||||
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: uint32(count), Data: content[:count]}
|
||||
case <-ctx.Done():
|
||||
return errFcall(fc, "interrupted")
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Default stream handling (no filter)
|
||||
cs.mu.RLock()
|
||||
f, fidOK := cs.fids[fc.Fid]
|
||||
var base string
|
||||
|
|
@ -688,7 +741,27 @@ func write(cs *connState, fc *plan9.Fcall) *plan9.Fcall {
|
|||
return errFcall(fc, "bad fid")
|
||||
}
|
||||
f.mu.Lock()
|
||||
path := f.path
|
||||
entry := f.entry
|
||||
|
||||
// Special handling for event file: write sets the filter
|
||||
if path == "/event" || path == "event" {
|
||||
filter := strings.TrimSpace(string(fc.Data))
|
||||
if filter == "" {
|
||||
filter = "*"
|
||||
}
|
||||
f.eventFilter = filter
|
||||
// Clear any existing subscription (will be recreated on next read)
|
||||
if f.eventCancel != nil {
|
||||
f.eventCancel()
|
||||
f.eventCancel = nil
|
||||
}
|
||||
f.eventCh = nil
|
||||
f.mu.Unlock()
|
||||
cs.mu.Unlock()
|
||||
return &plan9.Fcall{Type: plan9.Rwrite, Tag: fc.Tag, Count: uint32(len(fc.Data))}
|
||||
}
|
||||
|
||||
if entry != nil && entry.RdwrMode() {
|
||||
f.readCache = nil
|
||||
}
|
||||
|
|
@ -799,6 +872,10 @@ func clunk(root *fs.Tree, cs *connState, fc *plan9.Fcall) *plan9.Fcall {
|
|||
isRequest = entry.RdwrMode()
|
||||
}
|
||||
}
|
||||
// Cancel any event subscription
|
||||
if f.eventCancel != nil {
|
||||
f.eventCancel()
|
||||
}
|
||||
f.mu.Unlock()
|
||||
delete(cs.fids, fc.Fid)
|
||||
}
|
||||
|
|
|
|||
Loading…
Reference in New Issue