SSE Real-Time Update Mechanism in agentsview: Implementation Deep Dive
agentsview pushes live session changes to browsers using Server-Sent Events (SSE) centered around a reusable SSEStream helper that applies bounded write timeouts and integrates with the Huma framework to prevent stalled clients from blocking the server.
The agentsview repository implements a robust SSE real-time update mechanism to stream database changes and session events directly to connected clients. This architecture isolates low-level streaming protocols from HTTP routing logic, enabling both per-session updates and global broadcasts while maintaining server stability through strict write deadlines and configurable heartbeat intervals.
Core SSE Stream Abstraction
The foundation of agentsview’s real-time capability resides in internal/server/sse.go, which defines a self-contained SSEStream type that manages the HTTP response lifecycle and event formatting.
The SSEStream Type and Initialization
The SSEStream struct holds the http.ResponseWriter and its http.Flusher interface. The NewSSEStream constructor configures the mandatory SSE headers and immediately flushes them to the client:
Content-Type: text/event-streamCache-Control: no-cacheConnection: keep-alive
This initialization ensures the connection remains open for subsequent event streaming rather than closing after the initial response.
Write Deadline Protection
To prevent a stalled client from starving server resources, the Send method applies a bounded write deadline using sseWriteTimeout = 3s. If the client cannot accept data within this window, the write fails fast and the handler can terminate the connection. The ForceWriteDeadlineNow helper is utilized during server shutdown to unblock any pending writes immediately.
The Send method writes events in plain-text format following the SSE specification:
event: <name>\ndata: <payload>\n\n
For structured data, SendJSON marshals values to JSON before forwarding them to Send, ensuring type-safe transmission of complex objects like timing metadata.
Bridging Huma Framework and SSE
Since agentsview uses the Huma framework for HTTP routing, internal/server/huma_routes.go provides a newHumaSSEStream function that bridges Huma’s context abstraction to the raw SSEStream.
This helper extracts the underlying http.ResponseWriter via ctx.BodyWriter(), validates that the writer supports the http.Flusher interface, applies the standard SSE headers, performs an initial flush, and returns a configured *SSEStream. This pattern allows Huma route handlers to leverage the streaming abstraction without breaking framework conventions.
Per-Session Real-Time Updates
Session-specific monitoring is implemented in internal/server/huma_routes_sessions.go through the humaWatchSession function. This endpoint establishes a dedicated SSE stream for tracking changes to a specific session ID.
Session Monitoring and Heartbeats
Upon receiving a watch request, the handler initializes an SSE stream via newHumaSSEStream and invokes sessionMonitor from internal/server/events.go. The monitor creates a Go channel that emits signals whenever the session’s database state changes, utilizing sessionwatch.New(...).Events(ctx, sessionID) to poll the SQLite backend for modifications.
To prevent proxy timeouts and keep the connection alive, a heartbeat ticker runs at an interval calculated as sessionwatch.PollInterval * sessionwatch.HeartbeatTicks. Each tick sends a "heartbeat" event, ensuring the browser maintains an active connection even during periods of low database activity.
Event Types and JSON Payloads
When the watcher detects a change, the server emits two distinct events:
session_updated– A plain-text event indicating the session state has changedsession.timing– A JSON-encoded event containing detailed timing information, transmitted viaSendJSON
This dual-event approach allows lightweight clients to react to changes immediately while enabling rich clients to consume detailed metadata without additional API calls.
Global Event Broadcasting
Beyond individual session monitoring, agentsview supports global event distribution through the humaEvents function in internal/server/huma_routes.go. When the server’s broadcaster component is active, this endpoint opens an SSE stream, subscribes to the internal broadcaster via s.broadcaster.Subscribe(), and forwards every broadcasted message to connected clients.
The same heartbeat mechanism used in session watching prevents idle connections from timing out. Server code can publish messages to all connected clients using s.broadcaster.Publish("custom", msg), enabling real-time notifications across the entire user base.
Database Change Detection
The actual change detection logic resides in internal/server/events.go. The sessionMonitor function acts as a thin wrapper around sessionwatch.Watcher, which polls the SQLite database for modifications to the specified session ID. This polling-based approach avoids the complexity of database triggers or external message queues while still providing sub-second update latency for most operations.
Implementation Examples
Raw HTTP Handler with SSEStream
For custom endpoints outside the Huma router, you can initialize the stream directly:
func watchHandler(w http.ResponseWriter, r *http.Request) {
// Initialise the stream
stream, err := server.NewSSEStream(w)
if err != nil {
http.Error(w, "streaming not supported", http.StatusInternalServerError)
return
}
// Simple heartbeat loop
ticker := time.NewTicker(30 * time.Second)
defer ticker.Stop()
for {
select {
case <-r.Context().Done():
return // client disconnected
case <-ticker.C:
if !stream.Send("heartbeat", time.Now().UTC().String()) {
return // write failed
}
}
}
}
Consuming the Session Watch Endpoint
The Huma-registered route at /api/v1/sessions/:id/watch can be consumed via standard EventSource:
const evtSource = new EventSource("/api/v1/sessions/12345/watch");
evtSource.addEventListener("session_updated", e => {
console.log("session changed:", e.data);
});
evtSource.addEventListener("session.timing", e => {
const timing = JSON.parse(e.data);
console.log("timing info:", timing);
});
evtSource.addEventListener("heartbeat", e => {
console.log("ping:", e.data);
});
Broadcasting Server-Side Events
To push messages to all clients connected to the global events endpoint:
func (s *Server) broadcastMessage(msg string) {
if s.broadcaster == nil {
return
}
// All connected `/api/v1/events` clients receive this.
s.broadcaster.Publish("custom", msg)
}
Summary
internal/server/sse.goprovides the coreSSEStreamtype withNewSSEStream,Send, andSendJSONmethods, enforcing a 3-second write timeout to prevent blocking.newHumaSSEStreamininternal/server/huma_routes.gobridges the Huma framework to raw SSE streams by extractinghttp.ResponseWriterfromctx.BodyWriter().humaWatchSessionininternal/server/huma_routes_sessions.goimplements per-session monitoring with heartbeats and emits"session_updated"and"session.timing"events.sessionMonitorininternal/server/events.gowrapssessionwatch.Watcherto poll the SQLite database for changes.humaEventsenables global broadcasting vias.broadcaster.Subscribe()andPublish(), supporting real-time notifications across all connected clients.
Frequently Asked Questions
How does agentsview prevent SSE connections from blocking the server?
According to the source code in internal/server/sse.go, every write operation enforces a 3-second deadline via sseWriteTimeout. The Send method applies this deadline to the underlying connection, causing stalled clients to fail fast rather than holding goroutines indefinitely. Additionally, ForceWriteDeadlineNow unblocks pending writes during graceful shutdown.
What is the difference between per-session and global SSE endpoints?
The per-session endpoint (humaWatchSession) creates a dedicated sessionMonitor that polls the database for changes specific to a single session ID and emits "session_updated" events. In contrast, the global endpoint (humaEvents) subscribes to a central broadcaster (s.broadcaster.Subscribe()) and forwards messages published via s.broadcaster.Publish() to all connected clients regardless of session.
How does agentsview detect database changes for SSE transmission?
Change detection relies on sessionMonitor in internal/server/events.go, which wraps sessionwatch.New(...).Events(ctx, sessionID). This watcher polls the SQLite database at intervals defined by sessionwatch.PollInterval, emitting a signal on the returned channel whenever the session’s state differs from the previous check.
Why does agentsview send heartbeat events via SSE?
The heartbeat mechanism prevents proxy servers and browser connection timeouts from terminating idle streams. As implemented in humaWatchSession, a ticker runs at sessionwatch.PollInterval * sessionwatch.HeartbeatTicks to emit periodic "heartbeat" events, ensuring the TCP connection remains active even when no database changes occur for extended periods.
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 →