batch ChatUpdated D-Bus signals with 50ms coalescing

Previously every streaming token fired a separate D-Bus signal,
flooding the session bus and back-pressuring the HTTP stream reader.
New chatEmitter buffers text and flushes at most every 50ms. Log
writes remain immediate so GetChat always returns current data.
flushSync() drains the buffer before StateChanged signals.
This commit is contained in:
Levi Neely 2026-07-16 06:03:06 +02:00
parent 99a8dbfe88
commit 7900514548
1 changed files with 74 additions and 12 deletions

86
main.go
View File

@ -11,6 +11,7 @@ import (
"strings" "strings"
"sync" "sync"
"syscall" "syscall"
"time"
"github.com/godbus/dbus/v5" "github.com/godbus/dbus/v5"
"github.com/godbus/dbus/v5/introspect" "github.com/godbus/dbus/v5/introspect"
@ -42,6 +43,71 @@ type managedSession struct {
peers map[string]bool peers map[string]bool
} }
// chatEmitter batches ChatUpdated D-Bus signals to avoid flooding the bus.
// Text is appended to the session log immediately, but the D-Bus signal is
// coalesced and emitted at most every 50ms.
type chatEmitter struct {
conn *dbus.Conn
sess *managedSession
mu sync.Mutex
buf string
offset int64
timer *time.Timer
running bool
}
func newChatEmitter(conn *dbus.Conn, sess *managedSession) *chatEmitter {
return &chatEmitter{conn: conn, sess: sess}
}
func (e *chatEmitter) append(text string) {
// Write to log immediately so GetChat always returns the latest.
e.sess.logMu.Lock()
if e.buf == "" {
e.offset = int64(len(e.sess.log))
}
e.sess.log = append(e.sess.log, []byte(text)...)
e.sess.logMu.Unlock()
e.mu.Lock()
e.buf += text
if !e.running {
e.running = true
e.timer = time.AfterFunc(50*time.Millisecond, e.flush)
}
e.mu.Unlock()
}
func (e *chatEmitter) flush() {
e.mu.Lock()
text := e.buf
offset := e.offset
e.buf = ""
e.running = false
e.mu.Unlock()
if text != "" {
e.conn.Emit(busPath, busIface+".ChatUpdated", e.sess.id, offset, text)
}
}
// flushSync drains any pending buffered text immediately (used on state transitions).
func (e *chatEmitter) flushSync() {
e.mu.Lock()
if e.timer != nil {
e.timer.Stop()
}
text := e.buf
offset := e.offset
e.buf = ""
e.running = false
e.mu.Unlock()
if text != "" {
e.conn.Emit(busPath, busIface+".ChatUpdated", e.sess.id, offset, text)
}
}
// SessionManager is the D-Bus exported object. // SessionManager is the D-Bus exported object.
type SessionManager struct { type SessionManager struct {
mu sync.RWMutex mu sync.RWMutex
@ -135,6 +201,7 @@ func (m *SessionManager) CreateSession(cwd, backendName, modelName, agentName st
} }
// Subscribe to events for the chat log + signal dispatch // Subscribe to events for the chat log + signal dispatch
emitter := newChatEmitter(m.conn, sess)
streamingRole := "" streamingRole := ""
core.Bus().Subscribe("event", func(ev agent.Event) { core.Bus().Subscribe("event", func(ev agent.Event) {
var text string var text string
@ -175,12 +242,7 @@ func (m *SessionManager) CreateSession(cwd, backendName, modelName, agentName st
return return
} }
sess.logMu.Lock() emitter.append(text)
offset := int64(len(sess.log))
sess.log = append(sess.log, []byte(text)...)
sess.logMu.Unlock()
m.conn.Emit(busPath, busIface+".ChatUpdated", sess.id, offset, text)
}) })
// Watch state changes // Watch state changes
@ -189,9 +251,11 @@ func (m *SessionManager) CreateSession(cwd, backendName, modelName, agentName st
for { for {
next, ok := core.WaitChange(ctx, agent.WatchState, current) next, ok := core.WaitChange(ctx, agent.WatchState, current)
if !ok { if !ok {
emitter.flushSync()
return return
} }
current = next current = next
emitter.flushSync()
m.conn.Emit(busPath, busIface+".StateChanged", sess.id, current) m.conn.Emit(busPath, busIface+".StateChanged", sess.id, current)
// Send notification when a turn completes // Send notification when a turn completes
@ -749,6 +813,7 @@ func (m *SessionManager) restoreSession(ps *agent.PersistedSession) error {
} }
// Subscribe to events for the chat log + signal dispatch // Subscribe to events for the chat log + signal dispatch
emitter := newChatEmitter(m.conn, sess)
streamingRole := "" streamingRole := ""
core.Bus().Subscribe("event", func(ev agent.Event) { core.Bus().Subscribe("event", func(ev agent.Event) {
var text string var text string
@ -789,12 +854,7 @@ func (m *SessionManager) restoreSession(ps *agent.PersistedSession) error {
return return
} }
sess.logMu.Lock() emitter.append(text)
offset := int64(len(sess.log))
sess.log = append(sess.log, []byte(text)...)
sess.logMu.Unlock()
m.conn.Emit(busPath, busIface+".ChatUpdated", sess.id, offset, text)
}) })
// Watch state changes // Watch state changes
@ -803,9 +863,11 @@ func (m *SessionManager) restoreSession(ps *agent.PersistedSession) error {
for { for {
next, ok := core.WaitChange(ctx, agent.WatchState, current) next, ok := core.WaitChange(ctx, agent.WatchState, current)
if !ok { if !ok {
emitter.flushSync()
return return
} }
current = next current = next
emitter.flushSync()
m.conn.Emit(busPath, busIface+".StateChanged", sess.id, current) m.conn.Emit(busPath, busIface+".StateChanged", sess.id, current)
if current == "idle" { if current == "idle" {