server: event file plain read now uses per-fid subscription

Plain reads (without writing a filter first) now get a per-fid
subscription with '*' filter, identical to 'echo * | rdwrs event'.

This fixes event loss — the old EventStream path overwrote events
when multiple arrived before the reader consumed them.
This commit is contained in:
Levi Neely 2026-10-06 13:27:02 +02:00
parent c96ccd51e8
commit cd94cdf259
1 changed files with 24 additions and 23 deletions

View File

@ -598,7 +598,7 @@ func read(root *fs.Tree, cs *connState, fc *plan9.Fcall, ctx context.Context) *p
}
if entry.StreamMode() {
// Special handling for event file with filter
// Special handling for event file — always use per-fid subscription
if path == "event" {
cs.mu.RLock()
f, fidOK := cs.fids[fc.Fid]
@ -609,38 +609,39 @@ func read(root *fs.Tree, cs *connState, fc *plan9.Fcall, ctx context.Context) *p
eventCh := f.eventCh
f.mu.Unlock()
// If filter was written, use filtered subscription
if filter != "" && eventCh == nil {
// Create subscription on first read after write
// Create subscription on first read (filter defaults to "*" for all events)
if eventCh == nil {
if filter == "" {
filter = "*"
}
subCtx, cancel := context.WithCancel(ctx)
ch := session.SubscribeEventsFiltered(subCtx, filter)
f.mu.Lock()
f.eventFilter = filter
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")
// Read from subscription — blocks until event arrives
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")
}
}
}