replace pubsub library with simple fan-out event hub
The pubsub library had issues: - Published to literal '*' topic (nonsensical) - Used TrySend which drops events - Complex hierarchical wildcard publishing New implementation: - Simple eventHub with map of subscribers - PublishEvent fans out to all subscribers (blocking send) - SubscribeEvents returns channel, cleaned up on ctx cancel - SubscribeEventsFiltered filters client-side with MatchTopic - Removed simonfxr/pubsub dependency
This commit is contained in:
parent
63e7ed5c65
commit
d7d9e9654a
|
|
@ -116,9 +116,6 @@ func buildTreeSpec(cfg *Config) virtfs.FsNodeDecl {
|
|||
treeSpecUID = u.Username
|
||||
}
|
||||
|
||||
// Event subscription for event (lives for server lifetime).
|
||||
eventRead, eventSignal := session.EventStream(session.SubscribeEvents(cfg.Ctx, "*"))
|
||||
|
||||
return virtfs.DirNode("/",
|
||||
virtfs.UID(treeSpecUID),
|
||||
virtfs.GID(treeSpecGID),
|
||||
|
|
@ -241,7 +238,12 @@ func buildTreeSpec(cfg *Config) virtfs.FsNodeDecl {
|
|||
),
|
||||
virtfs.FileNode("event", 0666,
|
||||
virtfs.Doc("Event stream. Read: all events. Write filter then read: filtered events."),
|
||||
virtfs.Stream(eventRead, eventSignal),
|
||||
// Stream mode marker — actual handling is in server.go Tread handler
|
||||
virtfs.StreamRaw(func(ctx context.Context, base string) ([]byte, string, error) {
|
||||
// Never called — server.go intercepts event file reads
|
||||
<-ctx.Done()
|
||||
return nil, base, nil
|
||||
}),
|
||||
),
|
||||
virtfs.FileNode("generate", 0666,
|
||||
virtfs.Doc("One-shot LLM generation"),
|
||||
|
|
|
|||
|
|
@ -18,7 +18,8 @@ import (
|
|||
// wireAgentEvents sets the state change callback on an agent.
|
||||
func wireAgentEvents(sessID string, ag *agent.Agent) {
|
||||
ag.SetOnStateChange(func(agentID, state string) {
|
||||
session.PublishEvent("session."+sessID+".agent."+agentID+".state", state)
|
||||
topic := "session." + sessID + ".agent." + agentID + ".state"
|
||||
session.PublishEvent(topic, state)
|
||||
})
|
||||
}
|
||||
|
||||
|
|
|
|||
|
|
@ -4,59 +4,77 @@ import (
|
|||
"context"
|
||||
"strings"
|
||||
"sync"
|
||||
|
||||
"github.com/simonfxr/pubsub"
|
||||
)
|
||||
|
||||
var eventBus = pubsub.NewBus()
|
||||
|
||||
// Event represents a published event with topic and payload.
|
||||
type Event struct {
|
||||
Topic string
|
||||
Payload string
|
||||
}
|
||||
|
||||
// PublishEvent publishes an event to a specific topic.
|
||||
func PublishEvent(topic, payload string) {
|
||||
ev := Event{Topic: topic, Payload: payload}
|
||||
eventBus.Publish(topic, ev)
|
||||
|
||||
// Publish to hierarchical wildcards
|
||||
parts := strings.Split(topic, ".")
|
||||
for i := len(parts) - 1; i >= 1; i-- {
|
||||
wildcard := strings.Join(parts[:i], ".") + ".*"
|
||||
eventBus.Publish(wildcard, ev)
|
||||
}
|
||||
eventBus.Publish("*", ev)
|
||||
// eventHub manages event subscriptions with a simple fan-out pattern.
|
||||
type eventHub struct {
|
||||
mu sync.RWMutex
|
||||
subs map[*eventSub]struct{}
|
||||
}
|
||||
|
||||
// SubscribeEvents returns a channel that receives all events matching the topic.
|
||||
// Events are never dropped — publishing blocks if the channel is full.
|
||||
func SubscribeEvents(ctx context.Context, topic string) <-chan Event {
|
||||
// Use a callback instead of SubscribeChan to avoid non-blocking TrySend
|
||||
ch := make(chan Event, 256)
|
||||
sub := eventBus.Subscribe(topic, func(ev Event) {
|
||||
type eventSub struct {
|
||||
ch chan Event
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
}
|
||||
|
||||
var hub = &eventHub{
|
||||
subs: make(map[*eventSub]struct{}),
|
||||
}
|
||||
|
||||
// PublishEvent publishes an event to all subscribers.
|
||||
func PublishEvent(topic, payload string) {
|
||||
ev := Event{Topic: topic, Payload: payload}
|
||||
|
||||
hub.mu.RLock()
|
||||
defer hub.mu.RUnlock()
|
||||
|
||||
for sub := range hub.subs {
|
||||
select {
|
||||
case ch <- ev:
|
||||
case <-ctx.Done():
|
||||
case sub.ch <- ev:
|
||||
case <-sub.ctx.Done():
|
||||
// Subscriber cancelled, skip
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
// SubscribeEvents returns a channel that receives all events.
|
||||
// The channel is closed when ctx is cancelled.
|
||||
func SubscribeEvents(ctx context.Context) <-chan Event {
|
||||
subCtx, cancel := context.WithCancel(ctx)
|
||||
sub := &eventSub{
|
||||
ch: make(chan Event, 256),
|
||||
ctx: subCtx,
|
||||
cancel: cancel,
|
||||
}
|
||||
|
||||
hub.mu.Lock()
|
||||
hub.subs[sub] = struct{}{}
|
||||
hub.mu.Unlock()
|
||||
|
||||
go func() {
|
||||
<-ctx.Done()
|
||||
eventBus.Unsubscribe(sub)
|
||||
close(ch)
|
||||
<-subCtx.Done()
|
||||
hub.mu.Lock()
|
||||
delete(hub.subs, sub)
|
||||
hub.mu.Unlock()
|
||||
close(sub.ch)
|
||||
}()
|
||||
return ch
|
||||
|
||||
return sub.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, 256)
|
||||
allEvents := SubscribeEvents(ctx, "*")
|
||||
allEvents := SubscribeEvents(ctx)
|
||||
go func() {
|
||||
defer close(ch)
|
||||
for ev := range allEvents {
|
||||
|
|
@ -100,45 +118,3 @@ func MatchTopic(filter, topic string) bool {
|
|||
}
|
||||
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{}) {
|
||||
var mu sync.Mutex
|
||||
var pending []byte
|
||||
var counter int64
|
||||
sig := make(chan struct{})
|
||||
|
||||
go func() {
|
||||
for ev := range ch {
|
||||
line := ev.Topic
|
||||
if ev.Payload != "" {
|
||||
line += " " + ev.Payload
|
||||
}
|
||||
d := []byte(line + "\n")
|
||||
mu.Lock()
|
||||
pending = d
|
||||
counter++
|
||||
close(sig)
|
||||
sig = make(chan struct{})
|
||||
mu.Unlock()
|
||||
}
|
||||
}()
|
||||
|
||||
readFn = func(base string) ([]byte, string, error) {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
if len(pending) == 0 {
|
||||
return nil, base, nil // no data, will block on signal
|
||||
}
|
||||
out := pending
|
||||
pending = nil
|
||||
return out, "", nil
|
||||
}
|
||||
signalFn = func() <-chan struct{} {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
return sig
|
||||
}
|
||||
return
|
||||
}
|
||||
|
|
|
|||
|
|
@ -46,75 +46,116 @@ func TestSetCwd_AbsolutePathUnchanged(t *testing.T) {
|
|||
}
|
||||
|
||||
func TestEventSubscription(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
|
||||
ch := SubscribeEvents(ctx, "*")
|
||||
|
||||
// Give subscription time to register
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
|
||||
// Publish events
|
||||
for i := 0; i < 5; i++ {
|
||||
PublishEvent(fmt.Sprintf("test.event.%d", i), fmt.Sprintf("payload%d", i))
|
||||
}
|
||||
|
||||
// Collect events
|
||||
timeout := time.After(1 * time.Second)
|
||||
var received []Event
|
||||
for {
|
||||
select {
|
||||
case ev := <-ch:
|
||||
received = append(received, ev)
|
||||
if len(received) >= 5 {
|
||||
goto done
|
||||
}
|
||||
case <-timeout:
|
||||
goto done
|
||||
}
|
||||
}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
|
||||
ch := SubscribeEvents(ctx)
|
||||
|
||||
// Give subscription time to register
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
|
||||
// Publish events
|
||||
for i := 0; i < 5; i++ {
|
||||
PublishEvent(fmt.Sprintf("test.event.%d", i), fmt.Sprintf("payload%d", i))
|
||||
}
|
||||
|
||||
// Collect events
|
||||
timeout := time.After(1 * time.Second)
|
||||
var received []Event
|
||||
for {
|
||||
select {
|
||||
case ev := <-ch:
|
||||
received = append(received, ev)
|
||||
if len(received) >= 5 {
|
||||
goto done
|
||||
}
|
||||
case <-timeout:
|
||||
goto done
|
||||
}
|
||||
}
|
||||
done:
|
||||
if len(received) < 5 {
|
||||
t.Errorf("Expected 5 events, got %d", len(received))
|
||||
for i, ev := range received {
|
||||
t.Logf(" [%d] topic=%s payload=%s", i, ev.Topic, ev.Payload)
|
||||
}
|
||||
}
|
||||
if len(received) < 5 {
|
||||
t.Errorf("Expected 5 events, got %d", len(received))
|
||||
for i, ev := range received {
|
||||
t.Logf(" [%d] topic=%s payload=%s", i, ev.Topic, ev.Payload)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestEventSubscriptionFiltered(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
|
||||
ch := SubscribeEventsFiltered(ctx, "*")
|
||||
|
||||
// Give subscription time to register
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
|
||||
// Publish events
|
||||
for i := 0; i < 5; i++ {
|
||||
PublishEvent(fmt.Sprintf("test.event.%d", i), fmt.Sprintf("payload%d", i))
|
||||
}
|
||||
|
||||
// Collect events
|
||||
timeout := time.After(1 * time.Second)
|
||||
var received []Event
|
||||
for {
|
||||
select {
|
||||
case ev := <-ch:
|
||||
received = append(received, ev)
|
||||
if len(received) >= 5 {
|
||||
goto done
|
||||
}
|
||||
case <-timeout:
|
||||
goto done
|
||||
}
|
||||
}
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
|
||||
ch := SubscribeEventsFiltered(ctx, "*")
|
||||
|
||||
// Give subscription time to register
|
||||
time.Sleep(10 * time.Millisecond)
|
||||
|
||||
// Publish events
|
||||
for i := 0; i < 5; i++ {
|
||||
PublishEvent(fmt.Sprintf("test.event.%d", i), fmt.Sprintf("payload%d", i))
|
||||
}
|
||||
|
||||
// Collect events
|
||||
timeout := time.After(1 * time.Second)
|
||||
var received []Event
|
||||
for {
|
||||
select {
|
||||
case ev := <-ch:
|
||||
received = append(received, ev)
|
||||
if len(received) >= 5 {
|
||||
goto done
|
||||
}
|
||||
case <-timeout:
|
||||
goto done
|
||||
}
|
||||
}
|
||||
done:
|
||||
if len(received) < 5 {
|
||||
t.Errorf("Expected 5 events, got %d", len(received))
|
||||
for i, ev := range received {
|
||||
t.Logf(" [%d] topic=%s payload=%s", i, ev.Topic, ev.Payload)
|
||||
}
|
||||
}
|
||||
if len(received) < 5 {
|
||||
t.Errorf("Expected 5 events, got %d", len(received))
|
||||
for i, ev := range received {
|
||||
t.Logf(" [%d] topic=%s payload=%s", i, ev.Topic, ev.Payload)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestEventSubscriptionRapidFire(t *testing.T) {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
defer cancel()
|
||||
|
||||
ch := SubscribeEventsFiltered(ctx, "*")
|
||||
|
||||
// Give subscription time to register
|
||||
time.Sleep(50 * time.Millisecond)
|
||||
|
||||
// Rapid-fire publish (simulating state changes during a turn)
|
||||
states := []string{"thinking", "calling: shell", "thinking", "calling: file_read", "thinking", "idle"}
|
||||
for _, state := range states {
|
||||
PublishEvent("session.abc.agent.xyz.state", state)
|
||||
}
|
||||
|
||||
// Collect events
|
||||
timeout := time.After(2 * time.Second)
|
||||
var received []Event
|
||||
for {
|
||||
select {
|
||||
case ev := <-ch:
|
||||
received = append(received, ev)
|
||||
t.Logf("Received: topic=%s payload=%s", ev.Topic, ev.Payload)
|
||||
if len(received) >= len(states) {
|
||||
goto done
|
||||
}
|
||||
case <-timeout:
|
||||
goto done
|
||||
}
|
||||
}
|
||||
done:
|
||||
if len(received) != len(states) {
|
||||
t.Errorf("Expected %d events, got %d", len(states), len(received))
|
||||
}
|
||||
for i, ev := range received {
|
||||
if i < len(states) && ev.Payload != states[i] {
|
||||
t.Errorf("Event %d: expected payload %q, got %q", i, states[i], ev.Payload)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
|
|
|||
3
go.mod
3
go.mod
|
|
@ -6,7 +6,6 @@ require (
|
|||
9fans.net/go v0.0.7
|
||||
github.com/JohannesKaufmann/html-to-markdown v1.6.0
|
||||
github.com/godbus/dbus/v5 v5.1.0
|
||||
github.com/simonfxr/pubsub v0.0.5
|
||||
github.com/tree-sitter-grammars/tree-sitter-yaml v0.7.2
|
||||
github.com/tree-sitter/go-tree-sitter v0.25.0
|
||||
github.com/tree-sitter/tree-sitter-c v0.24.2
|
||||
|
|
@ -24,7 +23,7 @@ require (
|
|||
ollie/virtfs v0.0.0
|
||||
)
|
||||
|
||||
require github.com/yalue/onnxruntime_go v1.35.0 // indirect
|
||||
require github.com/yalue/onnxruntime_go v1.35.0
|
||||
|
||||
require (
|
||||
github.com/PuerkitoBio/goquery v1.9.2 // indirect
|
||||
|
|
|
|||
4
go.sum
4
go.sum
|
|
@ -29,12 +29,9 @@ github.com/sebdah/goldie/v2 v2.5.3/go.mod h1:oZ9fp0+se1eapSRjfYbsV/0Hqhbuu3bJVvK
|
|||
github.com/sergi/go-diff v1.0.0/go.mod h1:0CfEIISq7TuYL3j771MWULgwwjU+GofnZX9QAmXWZgo=
|
||||
github.com/sergi/go-diff v1.3.1 h1:xkr+Oxo4BOQKmkn/B9eMK0g5Kg/983T9DqqPHwYqD+8=
|
||||
github.com/sergi/go-diff v1.3.1/go.mod h1:aMJSSKb2lpPvRNec0+w3fl7LP9IOFzdc9Pa4NFbPK1I=
|
||||
github.com/simonfxr/pubsub v0.0.5 h1:DJfvFoglqGvwJriIOC5NI5um34n2YX8KAU5+7jv768w=
|
||||
github.com/simonfxr/pubsub v0.0.5/go.mod h1:bQ+B2NEEHZ08VY/0xVFGaGL1KMJ7cl3yLaXe8xvyC24=
|
||||
github.com/stretchr/objx v0.1.0/go.mod h1:HFkY916IF+rwdDfMAkV7OtwuqBVzrE8GR6GFx+wExME=
|
||||
github.com/stretchr/testify v1.3.0/go.mod h1:M5WIy9Dh21IEIfnGCwXGc5bZfKNJtfHm1UVUgZn+9EI=
|
||||
github.com/stretchr/testify v1.4.0/go.mod h1:j7eGeouHqKxXV5pUuKE4zz7dFj8WfuZ+81PSLYec5m4=
|
||||
github.com/stretchr/testify v1.6.1/go.mod h1:6Fq8oRcR53rry900zMqJjRRixrwX3KX962/h/Wwjteg=
|
||||
github.com/stretchr/testify v1.10.0 h1:Xv5erBjTwe/5IxqUQTdXv5kgmIvbHo3QQyRwhJsOfJA=
|
||||
github.com/stretchr/testify v1.10.0/go.mod h1:r2ic/lqez/lEtzL7wO/rwa5dbSLXVDPFyf8C91i36aY=
|
||||
github.com/tree-sitter-grammars/tree-sitter-markdown v0.3.2 h1:hQhxY9y3XYZHA2JOnhhvKUHtMBBn0vOl7iF4lbBxKzE=
|
||||
|
|
@ -157,6 +154,5 @@ gopkg.in/check.v1 v1.0.0-20190902080502-41f04d3bba15/go.mod h1:Co6ibVJAznAaIkqp8
|
|||
gopkg.in/yaml.v2 v2.2.2/go.mod h1:hI93XBmqTisBFMUTm0b8Fm+jr3Dg1NNxqwp+5A1VGuI=
|
||||
gopkg.in/yaml.v2 v2.4.0 h1:D8xgwECY7CYvx+Y2n4sBz93Jn9JRvxdiyyo8CTfuKaY=
|
||||
gopkg.in/yaml.v2 v2.4.0/go.mod h1:RDklbk79AGWmwhnvt/jBztapEOGDOx6ZbXqjP6csGnQ=
|
||||
gopkg.in/yaml.v3 v3.0.0-20200313102051-9f266ea9e77c/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
|
||||
gopkg.in/yaml.v3 v3.0.1 h1:fxVm/GzAzEWqLHuvctI91KS9hhNmmWOoWu0XTYJS7CA=
|
||||
gopkg.in/yaml.v3 v3.0.1/go.mod h1:K4uyk7z7BCEPqu6E+C64Yfv1cQ7kz7rIZviUmN+EgEM=
|
||||
|
|
|
|||
Loading…
Reference in New Issue