Replace event queue with in-memory pub-sub event bus
Replaces the TTL-based event queue (internal/events) with a proper in-memory pub-sub bus (internal/eventbus) backed by simonfxr/pubsub. The agent/events 9P endpoint is now a streaming subscriber: clients block on reads and receive one JSON-encoded event per read until they disconnect. Each open(2) creates a per-connection subscription; clunk and connection teardown cancel it automatically. - internal/eventbus/bus.go: new package wrapping simonfxr/pubsub; provides Publish(agent, type, data) and Subscribe() -> (chan, cancel) - internal/events/events.go: deleted - internal/session/manager.go: SetEventQueue -> SetEventBus; all publish sites updated to use eventbus constants - internal/p9/server.go: fid gains eventCh/eventUnsub fields; Topen subscribes, Tread blocks on channel, Tclunk/serve-defer cancel - cmd/anvilsrv/main.go: wired to SetEventBus Co-Authored-By: Claude Sonnet 4.6 <noreply@anthropic.com>
This commit is contained in:
parent
4b80a279f8
commit
407b3732af
|
|
@ -205,8 +205,8 @@ func start(daemonize bool) {
|
|||
}
|
||||
defer srv.Close()
|
||||
|
||||
// Wire up event queue to session manager
|
||||
mgr.SetEventQueue(srv.Events())
|
||||
// Wire up event bus to session manager
|
||||
mgr.SetEventBus(srv.Events())
|
||||
|
||||
log.Println("anvilsrv started successfully")
|
||||
log.Printf("9P server listening on %s", srv.SocketPath())
|
||||
|
|
|
|||
1
go.mod
1
go.mod
|
|
@ -125,6 +125,7 @@ require (
|
|||
github.com/sergi/go-diff v1.3.1 // indirect
|
||||
github.com/shopspring/decimal v1.4.0 // indirect
|
||||
github.com/silvasur/buzhash v0.0.0-20160816060738-9bdec3dec7c6 // indirect
|
||||
github.com/simonfxr/pubsub v0.0.5 // indirect
|
||||
github.com/sirupsen/logrus v1.9.3 // indirect
|
||||
github.com/skratchdot/open-golang v0.0.0-20200116055534-eef842397966 // indirect
|
||||
github.com/sony/gobreaker v0.5.0 // indirect
|
||||
|
|
|
|||
2
go.sum
2
go.sum
|
|
@ -707,6 +707,8 @@ github.com/shopspring/decimal v1.4.0 h1:bxl37RwXBklmTi0C79JfXCEBD1cqqHt0bbgBAGFp
|
|||
github.com/shopspring/decimal v1.4.0/go.mod h1:gawqmDU56v4yIKSwfBSFip1HdCCXN8/+DMd9qYNcwME=
|
||||
github.com/silvasur/buzhash v0.0.0-20160816060738-9bdec3dec7c6 h1:31fhvQj+O9qDqMxUgQDOCQA5RV1iIFMzYPhBUyzg2p0=
|
||||
github.com/silvasur/buzhash v0.0.0-20160816060738-9bdec3dec7c6/go.mod h1:jk5gVE20+MCoyJ2TFiiMrbWPyaH4t9T5F3HwVdthB2w=
|
||||
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/sirupsen/logrus v1.4.1/go.mod h1:ni0Sbl8bgC9z8RoU9G6nDWqqs/fq4eDPysMBDgk/93Q=
|
||||
github.com/sirupsen/logrus v1.4.2/go.mod h1:tLMulIdttU9McNUspp0xgXVQah82FyeX6MwdIuYE2rE=
|
||||
github.com/sirupsen/logrus v1.9.0/go.mod h1:naHLuLoDiP4jHNo9R0sCBMtWGeIprob74mVsIT4qYEQ=
|
||||
|
|
|
|||
|
|
@ -0,0 +1,76 @@
|
|||
// Package eventbus implements an in-memory pub-sub event bus for anvillm events.
|
||||
// It wraps github.com/simonfxr/pubsub to provide typed event publishing and
|
||||
// per-subscriber streaming channels.
|
||||
package eventbus
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
ps "github.com/simonfxr/pubsub"
|
||||
)
|
||||
|
||||
// Event type constants.
|
||||
const (
|
||||
EventStateChange = "StateChange"
|
||||
EventUserRecv = "UserRecv"
|
||||
EventUserSend = "UserSend"
|
||||
EventBotRecv = "BotRecv"
|
||||
EventBotSend = "BotSend"
|
||||
)
|
||||
|
||||
// allTopic is the single topic used for all events.
|
||||
const allTopic = "events"
|
||||
|
||||
// Event is the structure for all published events.
|
||||
type Event struct {
|
||||
ID string `json:"id"`
|
||||
TS int64 `json:"ts"`
|
||||
Agent string `json:"agent"`
|
||||
Type string `json:"type"`
|
||||
Data any `json:"data"`
|
||||
}
|
||||
|
||||
// Bus is an in-memory pub-sub event bus.
|
||||
// It is safe for concurrent use from multiple goroutines.
|
||||
type Bus struct {
|
||||
bus *ps.Bus
|
||||
}
|
||||
|
||||
// New creates a new Bus.
|
||||
func New() *Bus {
|
||||
return &Bus{bus: ps.NewBus()}
|
||||
}
|
||||
|
||||
// Publish emits an event to all current subscribers.
|
||||
// It is non-blocking; slow subscribers will have events dropped.
|
||||
func (b *Bus) Publish(agent, eventType string, data any) {
|
||||
e := &Event{
|
||||
ID: uuid.New().String(),
|
||||
TS: time.Now().Unix(),
|
||||
Agent: agent,
|
||||
Type: eventType,
|
||||
Data: data,
|
||||
}
|
||||
b.bus.Publish(allTopic, e)
|
||||
}
|
||||
|
||||
// Subscribe returns a read channel that receives *Event values and a cancel
|
||||
// function. The channel has a buffer of 64 events; events are dropped when
|
||||
// the buffer is full (slow consumer). Calling cancel removes the subscription
|
||||
// and closes the channel.
|
||||
func (b *Bus) Subscribe() (<-chan *Event, func()) {
|
||||
ch := make(chan *Event, 64)
|
||||
sub := b.bus.SubscribeChan(allTopic, ch, ps.CloseOnUnsubscribe)
|
||||
cancel := func() {
|
||||
b.bus.Unsubscribe(sub)
|
||||
}
|
||||
return ch, cancel
|
||||
}
|
||||
|
||||
// MarshalEvent encodes an event as a JSON line (with trailing newline).
|
||||
func MarshalEvent(e *Event) []byte {
|
||||
data, _ := json.Marshal(e)
|
||||
return append(data, '\n')
|
||||
}
|
||||
|
|
@ -1,129 +0,0 @@
|
|||
// Package events implements a simple event queue with TTL and explicit ack.
|
||||
package events
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/google/uuid"
|
||||
)
|
||||
|
||||
const DefaultTTL = 2 * time.Minute
|
||||
|
||||
const (
|
||||
EventStateChange = "StateChange"
|
||||
EventUserRecv = "UserRecv"
|
||||
EventUserSend = "UserSend"
|
||||
EventBotRecv = "BotRecv"
|
||||
EventBotSend = "BotSend"
|
||||
)
|
||||
|
||||
type Event struct {
|
||||
ID string `json:"id"`
|
||||
TS int64 `json:"ts"`
|
||||
Agent string `json:"agent"`
|
||||
Type string `json:"type"`
|
||||
Data any `json:"data"`
|
||||
expiry time.Time
|
||||
}
|
||||
|
||||
type Queue struct {
|
||||
mu sync.Mutex
|
||||
events []*Event
|
||||
ttl time.Duration
|
||||
}
|
||||
|
||||
func NewQueue() *Queue {
|
||||
q := &Queue{ttl: DefaultTTL}
|
||||
go q.expireLoop()
|
||||
return q
|
||||
}
|
||||
|
||||
// Push adds an event to the queue.
|
||||
func (q *Queue) Push(agent, eventType string, data any) string {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
|
||||
e := &Event{
|
||||
ID: uuid.New().String(),
|
||||
TS: time.Now().Unix(),
|
||||
Agent: agent,
|
||||
Type: eventType,
|
||||
Data: data,
|
||||
expiry: time.Now().Add(q.ttl),
|
||||
}
|
||||
q.events = append(q.events, e)
|
||||
return e.ID
|
||||
}
|
||||
|
||||
// Read returns all pending events as JSON lines.
|
||||
func (q *Queue) Read() []byte {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
|
||||
if len(q.events) == 0 {
|
||||
return nil
|
||||
}
|
||||
|
||||
var out []byte
|
||||
for _, e := range q.events {
|
||||
line, _ := json.Marshal(e)
|
||||
out = append(out, line...)
|
||||
out = append(out, '\n')
|
||||
}
|
||||
return out
|
||||
}
|
||||
|
||||
// AckRequest is the JSON payload for acknowledging events.
|
||||
type AckRequest struct {
|
||||
Msg string `json:"msg"`
|
||||
IDs []string `json:"ids"`
|
||||
}
|
||||
|
||||
// Ack removes events with the given IDs.
|
||||
func (q *Queue) Ack(data []byte) error {
|
||||
var req AckRequest
|
||||
if err := json.Unmarshal(data, &req); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
|
||||
idSet := make(map[string]bool, len(req.IDs))
|
||||
for _, id := range req.IDs {
|
||||
idSet[id] = true
|
||||
}
|
||||
|
||||
filtered := q.events[:0]
|
||||
for _, e := range q.events {
|
||||
if !idSet[e.ID] {
|
||||
filtered = append(filtered, e)
|
||||
}
|
||||
}
|
||||
q.events = filtered
|
||||
return nil
|
||||
}
|
||||
|
||||
func (q *Queue) expireLoop() {
|
||||
ticker := time.NewTicker(30 * time.Second)
|
||||
defer ticker.Stop()
|
||||
for range ticker.C {
|
||||
q.expire()
|
||||
}
|
||||
}
|
||||
|
||||
func (q *Queue) expire() {
|
||||
q.mu.Lock()
|
||||
defer q.mu.Unlock()
|
||||
|
||||
now := time.Now()
|
||||
filtered := q.events[:0]
|
||||
for _, e := range q.events {
|
||||
if e.expiry.After(now) {
|
||||
filtered = append(filtered, e)
|
||||
}
|
||||
}
|
||||
q.events = filtered
|
||||
}
|
||||
|
|
@ -4,7 +4,7 @@ package p9
|
|||
import (
|
||||
"anvillm/internal/backend"
|
||||
"anvillm/internal/backend/tmux"
|
||||
"anvillm/internal/events"
|
||||
"anvillm/internal/eventbus"
|
||||
"anvillm/internal/mailbox"
|
||||
"anvillm/internal/session"
|
||||
"context"
|
||||
|
|
@ -113,7 +113,7 @@ type Server struct {
|
|||
mgr *session.Manager
|
||||
listener net.Listener
|
||||
socketPath string
|
||||
events *events.Queue
|
||||
events *eventbus.Bus
|
||||
beads *BeadsFS
|
||||
OnAliasChange func(backend.Session) // Called when session alias changes
|
||||
mu sync.RWMutex
|
||||
|
|
@ -126,10 +126,13 @@ type connState struct {
|
|||
}
|
||||
|
||||
type fid struct {
|
||||
qid plan9.Qid
|
||||
path string
|
||||
mode uint8
|
||||
offset int64
|
||||
qid plan9.Qid
|
||||
path string
|
||||
mode uint8
|
||||
offset int64
|
||||
// For streaming /events endpoint
|
||||
eventCh <-chan *eventbus.Event
|
||||
eventUnsub func()
|
||||
}
|
||||
|
||||
// NewServer creates and starts the 9P server.
|
||||
|
|
@ -166,7 +169,7 @@ func NewServer(mgr *session.Manager, beadsStore bd.Storage) (*Server, error) {
|
|||
mgr: mgr,
|
||||
listener: listener,
|
||||
socketPath: sockPath,
|
||||
events: events.NewQueue(),
|
||||
events: eventbus.New(),
|
||||
beads: beadsFS,
|
||||
}
|
||||
go s.acceptLoop()
|
||||
|
|
@ -178,8 +181,8 @@ func (s *Server) SocketPath() string {
|
|||
return s.socketPath
|
||||
}
|
||||
|
||||
// Events returns the event queue for pushing events.
|
||||
func (s *Server) Events() *events.Queue {
|
||||
// Events returns the event bus for publishing events.
|
||||
func (s *Server) Events() *eventbus.Bus {
|
||||
return s.events
|
||||
}
|
||||
|
||||
|
|
@ -200,6 +203,18 @@ func (s *Server) serve(conn net.Conn) {
|
|||
defer conn.Close()
|
||||
cs := &connState{fids: make(map[uint32]*fid)}
|
||||
|
||||
// Cancel any outstanding event subscriptions when the connection drops.
|
||||
defer func() {
|
||||
cs.mu.Lock()
|
||||
for _, f := range cs.fids {
|
||||
if f.eventUnsub != nil {
|
||||
f.eventUnsub()
|
||||
f.eventUnsub = nil
|
||||
}
|
||||
}
|
||||
cs.mu.Unlock()
|
||||
}()
|
||||
|
||||
for {
|
||||
fc, err := plan9.ReadFcall(conn)
|
||||
if err != nil {
|
||||
|
|
@ -241,6 +256,11 @@ func (s *Server) handle(cs *connState, fc *plan9.Fcall) *plan9.Fcall {
|
|||
return s.stat(cs, fc)
|
||||
case plan9.Tclunk:
|
||||
cs.mu.Lock()
|
||||
if f, ok := cs.fids[fc.Fid]; ok && f.eventUnsub != nil {
|
||||
// Cancel event subscription before dropping the fid.
|
||||
f.eventUnsub()
|
||||
f.eventUnsub = nil
|
||||
}
|
||||
delete(cs.fids, fc.Fid)
|
||||
cs.mu.Unlock()
|
||||
return &plan9.Fcall{Type: plan9.Rclunk, Tag: fc.Tag}
|
||||
|
|
@ -411,6 +431,14 @@ func (s *Server) open(cs *connState, fc *plan9.Fcall) *plan9.Fcall {
|
|||
}
|
||||
f.mode = fc.Mode
|
||||
f.offset = 0
|
||||
|
||||
// For the /events streaming endpoint, subscribe to the event bus.
|
||||
if f.path == "/events" {
|
||||
ch, cancel := s.events.Subscribe()
|
||||
f.eventCh = ch
|
||||
f.eventUnsub = cancel
|
||||
}
|
||||
|
||||
return &plan9.Fcall{Type: plan9.Ropen, Tag: fc.Tag, Qid: f.qid}
|
||||
}
|
||||
|
||||
|
|
@ -426,6 +454,17 @@ func (s *Server) read(cs *connState, fc *plan9.Fcall) *plan9.Fcall {
|
|||
return errFcall(fc, "bad fid")
|
||||
}
|
||||
|
||||
// Streaming /events: block until next event arrives (or channel closed).
|
||||
if f.eventCh != nil {
|
||||
e, ok := <-f.eventCh
|
||||
if !ok {
|
||||
// Channel closed (subscription cancelled); signal EOF.
|
||||
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: 0}
|
||||
}
|
||||
data := eventbus.MarshalEvent(e)
|
||||
return &plan9.Fcall{Type: plan9.Rread, Tag: fc.Tag, Count: uint32(len(data)), Data: data}
|
||||
}
|
||||
|
||||
var data []byte
|
||||
|
||||
if f.qid.Type&QTDir != 0 {
|
||||
|
|
@ -465,14 +504,6 @@ func (s *Server) write(cs *connState, fc *plan9.Fcall) *plan9.Fcall {
|
|||
return errFcall(fc, "beads not initialized")
|
||||
}
|
||||
|
||||
// /events - ack events
|
||||
if f.path == "/events" {
|
||||
if err := s.events.Ack(fc.Data); err != nil {
|
||||
return errFcall(fc, err.Error())
|
||||
}
|
||||
return &plan9.Fcall{Type: plan9.Rwrite, Tag: fc.Tag, Count: uint32(len(fc.Data))}
|
||||
}
|
||||
|
||||
// /ctl - create new session
|
||||
if f.path == "/ctl" {
|
||||
args := strings.Fields(input)
|
||||
|
|
@ -1069,10 +1100,6 @@ func (s *Server) readFile(path string) string {
|
|||
return strings.Join(lines, "\n") + "\n"
|
||||
}
|
||||
|
||||
if path == "/events" {
|
||||
return string(s.events.Read())
|
||||
}
|
||||
|
||||
if path == "/status" {
|
||||
var lines []string
|
||||
mailMgr := s.mgr.GetMailManager()
|
||||
|
|
|
|||
|
|
@ -4,7 +4,7 @@ package session
|
|||
import (
|
||||
"anvillm/internal/backend"
|
||||
"anvillm/internal/backend/tmux"
|
||||
"anvillm/internal/events"
|
||||
"anvillm/internal/eventbus"
|
||||
"anvillm/internal/mailbox"
|
||||
"context"
|
||||
"fmt"
|
||||
|
|
@ -19,7 +19,7 @@ type Manager struct {
|
|||
backends map[string]backend.Backend
|
||||
sessions map[string]backend.Session
|
||||
mailManager *mailbox.Manager
|
||||
eventQueue *events.Queue
|
||||
eventBus *eventbus.Bus
|
||||
OnStateChange func(sessionID, oldState, newState string)
|
||||
mu sync.RWMutex
|
||||
stopCh chan struct{}
|
||||
|
|
@ -29,46 +29,46 @@ type Manager struct {
|
|||
// NewManager creates a session manager with the given backends
|
||||
func NewManager(backends map[string]backend.Backend) *Manager {
|
||||
mailMgr := mailbox.NewManager()
|
||||
|
||||
|
||||
m := &Manager{
|
||||
backends: backends,
|
||||
sessions: make(map[string]backend.Session),
|
||||
mailManager: mailMgr,
|
||||
eventQueue: nil, // Set via SetEventQueue
|
||||
eventBus: nil, // Set via SetEventBus
|
||||
stopCh: make(chan struct{}),
|
||||
}
|
||||
|
||||
|
||||
// Set session getter for alias lookup
|
||||
mailMgr.SetSessionGetter(m)
|
||||
|
||||
|
||||
// Set state change callback
|
||||
m.OnStateChange = func(sessionID, oldState, newState string) {
|
||||
if m.eventQueue != nil {
|
||||
m.eventQueue.Push(sessionID, events.EventStateChange, map[string]string{
|
||||
if m.eventBus != nil {
|
||||
m.eventBus.Publish(sessionID, eventbus.EventStateChange, map[string]string{
|
||||
"old_state": oldState,
|
||||
"new_state": newState,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
// Wire up mailbox event callbacks
|
||||
mailMgr.SetEventCallbacks(
|
||||
func(senderID string, msg *mailbox.Message) {
|
||||
if m.eventQueue != nil {
|
||||
eventType := events.EventBotSend
|
||||
if m.eventBus != nil {
|
||||
evType := eventbus.EventBotSend
|
||||
if senderID == "user" {
|
||||
eventType = events.EventUserSend
|
||||
evType = eventbus.EventUserSend
|
||||
}
|
||||
m.eventQueue.Push(senderID, eventType, msg)
|
||||
m.eventBus.Publish(senderID, evType, msg)
|
||||
}
|
||||
},
|
||||
func(receiverID string, msg *mailbox.Message) {
|
||||
if m.eventQueue != nil {
|
||||
eventType := events.EventBotRecv
|
||||
if m.eventBus != nil {
|
||||
evType := eventbus.EventBotRecv
|
||||
if receiverID == "user" {
|
||||
eventType = events.EventUserRecv
|
||||
evType = eventbus.EventUserRecv
|
||||
}
|
||||
m.eventQueue.Push(receiverID, eventType, msg)
|
||||
m.eventBus.Publish(receiverID, evType, msg)
|
||||
}
|
||||
},
|
||||
)
|
||||
|
|
@ -276,9 +276,9 @@ func (m *Manager) GetMailManager() *mailbox.Manager {
|
|||
return m.mailManager
|
||||
}
|
||||
|
||||
// SetEventQueue sets the event queue for emitting events
|
||||
func (m *Manager) SetEventQueue(eq *events.Queue) {
|
||||
// SetEventBus sets the event bus for emitting events.
|
||||
func (m *Manager) SetEventBus(bus *eventbus.Bus) {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
m.eventQueue = eq
|
||||
m.eventBus = bus
|
||||
}
|
||||
|
|
|
|||
Reference in New Issue