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 minInterval window, they are not sent instantly. Instead, the broadcaster stores the most recent scope in a pending field 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:

  1. Sync Engine Emission: After each successful sync pass, internal/sync/engine.go calls broadcaster.Emit(scope) to announce new data.
  2. HTTP Handler Subscription: The /watch handler in internal/server/server.go creates an SSEStream, subscribes to the broadcaster (ch, unsub := broadcaster.Subscribe()), and enters a forwarding loop.
  3. 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 (default 10s) into a single notification.
  • The sync engine calls broadcaster.Emit() after each pass, while the /watch handler forwards coalesced events to clients via SSEStream.Send().
  • Set events_coalesce_interval to 0 to 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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →