How AgentsView Implements SSE Real-Time Updates with Event Coalescing
AgentsView combines an SSEStream wrapper with a Broadcaster hub that uses leading-edge and trailing-edge event coalescing to stream live updates efficiently, collapsing rapid bursts into single notifications while delivering isolated changes immediately.
The open-source AgentsView project (kenn-io/agentsview) keeps browsers and CLI clients synchronized with backend state changes using SSE real-time updates with event coalescing. This architecture prevents network saturation during intensive sync operations while ensuring critical updates reach users with minimal latency.
The SSE Architecture
The implementation centers on two tightly integrated components: a low-level stream wrapper that manages HTTP connections and a broadcaster hub that intelligently throttles event delivery.
SSEStream: Managing the HTTP Connection
The SSEStream struct in internal/server/sse.go handles the raw Server-Sent Events protocol. It wraps an http.ResponseWriter to write properly formatted SSE frames, enforces a bounded write deadline (sseWriteTimeout = 3s) to detect stalled clients, and provides a Send method for JSON payloads. The server instantiates this stream for any endpoint matching the /watch pattern.
The Broadcaster Hub
Located in internal/server/broadcaster.go, the Broadcaster acts as a publish-subscribe hub that sync-engine components use to announce data changes. Each SSE client subscribes to receive Event{Scope string} objects. The broadcaster implements leading-edge + trailing-edge rate limiting (event coalescing) to optimize network usage.
How Event Coalescing Works
Event coalescing prevents the server from spamming clients during high-frequency update periods, such as long synchronization runs. The mechanism operates on a configurable minInterval (default 10s, exposed as EventsCoalesceInterval in internal/config/config.go).
- Leading-edge behavior: When an event arrives after a quiet period exceeding
minInterval, the broadcaster sends it immediately. - Trailing-edge behavior: If subsequent events arrive within the
minIntervalwindow, they are not sent instantly. Instead, the broadcaster stores the most recent scope in apendingfield and schedules a single trailing broadcast to fire when the interval expires. - Coalescing effect: The trailing broadcast delivers only the latest scope, collapsing an entire burst of updates into one SSE packet.
This design guarantees low latency for isolated updates while protecting clients from flooding during heavy activity.
End-to-End Data Flow
The data pipeline connects the synchronization engine to the browser through three stages:
- Sync Engine Emission: After each successful sync pass,
internal/sync/engine.gocallsbroadcaster.Emit(scope)to announce new data. - HTTP Handler Subscription: The
/watchhandler ininternal/server/server.gocreates anSSEStream, subscribes to the broadcaster (ch, unsub := broadcaster.Subscribe()), and enters a forwarding loop. - Client Delivery: Each event received from the broadcaster is forwarded to the client via
sse.Send(event.Scope).
Creating the SSE Connection
The HTTP handler wires the stream to the broadcaster as shown in this simplified excerpt from internal/server/server.go:
func handleWatch(w http.ResponseWriter, r *http.Request) {
// Create the SSE connection with proper headers and timeout
stream, err := NewSSEStream(w)
if err != nil {
http.Error(w, err.Error(), http.StatusInternalServerError)
return
}
// Subscribe to the global broadcaster
events, unsubscribe := broadcaster.Subscribe()
defer unsubscribe()
// Forward each coalesced event to the client
for ev := range events {
if !stream.Send("update", ev.Scope) {
// Client disconnected; exit loop
return
}
}
}
The Coalescing Logic
The Broadcaster.Emit method in internal/server/broadcaster.go implements the leading-edge and trailing-edge logic:
func (b *Broadcaster) Emit(scope string) {
b.mu.Lock()
defer b.mu.Unlock()
now := time.Now()
// Immediate send if outside coalesce window or disabled
if b.minInterval == 0 || b.lastEmit.IsZero() ||
now.Sub(b.lastEmit) >= b.minInterval {
b.pending = nil
if b.timer != nil {
b.timer.Stop()
b.timer = nil
}
b.timerGen++
b.lastEmit = now
b.broadcastLocked(Event{Scope: scope})
return
}
// Inside window: store latest scope and schedule trailing flush
b.pending = &Event{Scope: scope}
if b.timer == nil {
gen := b.timerGen
wait := b.minInterval - now.Sub(b.lastEmit)
b.timer = time.AfterFunc(wait, func() {
b.flushTrailing(gen)
})
}
}
Configuration and Tuning
You can tune the coalescing behavior via the configuration file defined in internal/config/config.go. The EventsCoalesceInterval field accepts a duration string:
# config.toml
events_coalesce_interval = "10s"
Setting events_coalesce_interval = "0s" disables coalescing entirely, causing every Emit call to broadcast immediately. The sseWriteTimeout remains hardcoded at 3s in internal/server/sse.go to prevent hung connections from consuming resources.
Summary
- SSEStream (
internal/server/sse.go) manages the raw HTTP connection, write deadlines, and SSE framing. - Broadcaster (
internal/server/broadcaster.go) implements pub-sub with leading-edge immediate delivery and trailing-edge coalescing. - Event coalescing collapses bursts of updates within the configured
minInterval(default10s) into a single notification. - The sync engine calls
broadcaster.Emit()after each pass, while the/watchhandler forwards coalesced events to clients viaSSEStream.Send(). - Set
events_coalesce_intervalto0to disable coalescing for real-time granular updates.
Frequently Asked Questions
What is event coalescing in AgentsView?
Event coalescing is a rate-limiting strategy in the Broadcaster component that collapses multiple rapid update notifications into a single SSE packet. When updates occur faster than the configured events_coalesce_interval, only the most recent scope is retained and delivered after the interval expires, reducing network overhead.
How does the broadcaster decide when to send immediate versus coalesced updates?
The broadcaster uses leading-edge detection: if an event arrives after a quiet period exceeding minInterval, it sends immediately. If additional events arrive within that window, they trigger the trailing-edge behavior, where the latest event is stored and a timer schedules delivery for when the interval expires.
Can I disable event coalescing if I need every single update?
Yes. Set events_coalesce_interval = "0s" in your configuration file. When minInterval is zero, the Emit method in internal/server/broadcaster.go bypasses the timer logic and calls broadcastLocked immediately for every event, ensuring no updates are deferred or merged.
What happens if a client connection stalls during SSE streaming?
The SSEStream wrapper enforces a 3s write deadline (sseWriteTimeout) on the underlying http.ResponseWriter. If a write operation blocks longer than this timeout, the connection is terminated, allowing the handler to clean up the broadcaster subscription and free resources.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →