The generated strict-server path could only return an io.Reader that io.Copy drains without flushing, so SSE events sat chunk-buffered instead of streaming in real time. Replace it with a raw http.ResponseWriter handler (serveSSE) that Flush()es after every event. Routing: chi allows a later registration to supersede an earlier one for the same method+path (verified empirically for v5.3.1), so serveSSE is registered on the router AFTER gen.HandlerWithOptions and wins over the generated /events/stream route. It inherits the base middleware chain and applies auth via With(). The generated StreamEvents method now returns an error (never reached) so a routing regression fails loudly rather than silently reverting to buffered delivery. Adds TestSSEStreamRealtimeDelivery: a real httptest.NewServer + streaming client (NewRecorder can't flush) that connects, triggers an event, and asserts delivery within 3s — proving both the override routing and per-event flushing. Passes in <1s. Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
340 lines
9.3 KiB
Go
340 lines
9.3 KiB
Go
package httpapi
|
|
|
|
import (
|
|
"container/list"
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"log/slog"
|
|
"net/http"
|
|
"strconv"
|
|
"sync"
|
|
"time"
|
|
|
|
"github.com/dtoro/oikos/internal/db/sqlcgen"
|
|
"github.com/dtoro/oikos/internal/httpapi/gen"
|
|
"github.com/jackc/pgx/v5"
|
|
"github.com/jackc/pgx/v5/pgxpool"
|
|
)
|
|
|
|
// SSE broker is a keep-last-event-id in-memory buffer used for subscriber
|
|
// fan-out. The database NOTIFY is the primary delivery mechanism; this buffer
|
|
// just supports the Last-Event-ID replay on connect.
|
|
type sseBroker struct {
|
|
mu sync.Mutex
|
|
buf *list.List // list of sqlcgen.Event
|
|
cache map[int64]*list.Element // id → list element for O(1) lookup
|
|
cap int
|
|
lastID int64
|
|
}
|
|
|
|
func newSSEBroker(capacity int) *sseBroker {
|
|
return &sseBroker{
|
|
buf: list.New(),
|
|
cache: make(map[int64]*list.Element),
|
|
cap: capacity,
|
|
}
|
|
}
|
|
|
|
func (b *sseBroker) push(ev sqlcgen.Event) {
|
|
b.mu.Lock()
|
|
defer b.mu.Unlock()
|
|
|
|
// Evict oldest if at capacity
|
|
for b.buf.Len() >= b.cap && b.buf.Len() > 0 {
|
|
front := b.buf.Front()
|
|
b.cache[front.Value.(sqlcgen.Event).ID] = nil // don't delete, just nil
|
|
b.buf.Remove(front)
|
|
}
|
|
|
|
elem := b.buf.PushBack(ev)
|
|
b.cache[ev.ID] = elem
|
|
if ev.ID > b.lastID {
|
|
b.lastID = ev.ID
|
|
}
|
|
}
|
|
|
|
func (b *sseBroker) after(id int64) []sqlcgen.Event {
|
|
b.mu.Lock()
|
|
defer b.mu.Unlock()
|
|
|
|
if id >= b.lastID {
|
|
return nil
|
|
}
|
|
|
|
// Walk from the front to find the first element after id
|
|
var events []sqlcgen.Event
|
|
for e := b.buf.Front(); e != nil; e = e.Next() {
|
|
ev := e.Value.(sqlcgen.Event)
|
|
if ev.ID > id {
|
|
events = append(events, ev)
|
|
}
|
|
}
|
|
return events
|
|
}
|
|
|
|
func (b *sseBroker) latestID() int64 {
|
|
b.mu.Lock()
|
|
defer b.mu.Unlock()
|
|
return b.lastID
|
|
}
|
|
|
|
// notifyPayload is the JSON payload from the pg_notify trigger (migration 008).
|
|
type notifyPayload struct {
|
|
ID int64 `json:"id"`
|
|
Ts string `json:"ts"`
|
|
Type string `json:"type"`
|
|
EntityID *string `json:"entity_id"`
|
|
Severity string `json:"severity"`
|
|
Source string `json:"source"`
|
|
CorrelationID *string `json:"correlation_id"`
|
|
}
|
|
|
|
// sseSubscriber holds the channels and cancel func for one SSE client.
|
|
type sseSubscriber struct {
|
|
ch chan sqlcgen.Event
|
|
done chan struct{}
|
|
cancel context.CancelFunc
|
|
}
|
|
|
|
// sseListener runs in a background goroutine: it opens a dedicated pgx
|
|
// connection, LISTENs on oikos_events, and fans out each notification to
|
|
// all live subscribers. Runs until ctx is cancelled.
|
|
func (s *Server) sseListener(ctx context.Context) {
|
|
poolConn, err := s.pool.Acquire(ctx)
|
|
if err != nil {
|
|
slog.Error("sse listener acquire failed", "error", err)
|
|
return
|
|
}
|
|
defer poolConn.Release()
|
|
|
|
conn := poolConn.Conn()
|
|
if _, err := conn.Exec(ctx, "LISTEN oikos_events"); err != nil {
|
|
slog.Error("sse listener listen failed", "error", err)
|
|
return
|
|
}
|
|
|
|
slog.Info("sse listener started on oikos_events")
|
|
defer slog.Info("sse listener stopped")
|
|
|
|
for {
|
|
nt, err := conn.WaitForNotification(ctx)
|
|
if err != nil {
|
|
if ctx.Err() != nil {
|
|
return // normal shutdown
|
|
}
|
|
slog.Error("sse listener notification error", "error", err)
|
|
// Reconnect on error after a brief delay
|
|
select {
|
|
case <-time.After(5 * time.Second):
|
|
case <-ctx.Done():
|
|
return
|
|
}
|
|
// Re-acquire connection
|
|
poolConn.Release()
|
|
var reconnErr error
|
|
poolConn, reconnErr = s.pool.Acquire(ctx)
|
|
if reconnErr != nil {
|
|
slog.Error("sse listener reconnect failed", "error", reconnErr)
|
|
return
|
|
}
|
|
conn = poolConn.Conn()
|
|
if _, err := conn.Exec(ctx, "LISTEN oikos_events"); err != nil {
|
|
slog.Error("sse listener re-listen failed", "error", err)
|
|
return
|
|
}
|
|
continue
|
|
}
|
|
|
|
var p notifyPayload
|
|
if err := json.Unmarshal([]byte(nt.Payload), &p); err != nil {
|
|
slog.Error("sse listener unmarshal failed", "error", err)
|
|
continue
|
|
}
|
|
|
|
// Fetch full event from DB
|
|
q := sqlcgen.New(s.pool)
|
|
events, err := q.ListEventsAfter(ctx, sqlcgen.ListEventsAfterParams{
|
|
ID: p.ID - 1,
|
|
Limit: 1,
|
|
})
|
|
if err != nil || len(events) == 0 {
|
|
slog.Warn("sse listener event fetch failed", "id", p.ID, "error", err)
|
|
continue
|
|
}
|
|
ev := events[0]
|
|
|
|
// Push to broker
|
|
s.sseBroker.push(ev)
|
|
|
|
// Fan out to subscribers (non-blocking send)
|
|
s.sseMu.Lock()
|
|
for sub := range s.sseSubs {
|
|
select {
|
|
case sub.ch <- ev:
|
|
default:
|
|
// Subscriber too slow — drop event for them
|
|
// (they'll reconnect via Last-Event-ID)
|
|
}
|
|
}
|
|
s.sseMu.Unlock()
|
|
}
|
|
}
|
|
|
|
// sqlcEventToGen converts a DB event row to the canonical wire shape so the
|
|
// SSE `data:` payload matches GET /events (snake_case keys, decoded data
|
|
// object) rather than leaking Go field names and base64-encoded JSONB.
|
|
func sqlcEventToGen(ev sqlcgen.Event) gen.Event {
|
|
out := gen.Event{
|
|
Id: int(ev.ID),
|
|
Ts: ev.Ts,
|
|
Type: ev.Type,
|
|
Severity: gen.EventSeverity(ev.Severity),
|
|
Source: ev.Source,
|
|
CorrelationId: ev.CorrelationID,
|
|
}
|
|
if ev.EntityID != nil {
|
|
s := ev.EntityID.String()
|
|
out.EntityId = &s
|
|
}
|
|
if len(ev.Data) > 0 {
|
|
var data map[string]any
|
|
if json.Unmarshal(ev.Data, &data) == nil && len(data) > 0 {
|
|
out.Data = &data
|
|
}
|
|
}
|
|
return out
|
|
}
|
|
|
|
// writeSSE writes a single Event as an SSE message. Returns false if the
|
|
// write failed (client disconnected). flusher may be nil (io.Pipe path,
|
|
// which has no separate flush step).
|
|
func writeSSE(w ioWriter, flusher http.Flusher, ev sqlcgen.Event) bool {
|
|
data, err := json.Marshal(sqlcEventToGen(ev))
|
|
if err != nil {
|
|
return true // skip un-serializable events
|
|
}
|
|
_, err = fmt.Fprintf(w, "id: %d\nevent: %s\ndata: %s\n\n", ev.ID, ev.Type, data)
|
|
if err != nil {
|
|
return false
|
|
}
|
|
if flusher != nil {
|
|
flusher.Flush()
|
|
}
|
|
return true
|
|
}
|
|
|
|
// ioWriter is an interface satisfied by both http.ResponseWriter and
|
|
// io.StringWriter, letting writeSSE work with the writer directly.
|
|
type ioWriter interface {
|
|
Write([]byte) (int, error)
|
|
}
|
|
|
|
// serveSSE is the streaming handler for GET /api/v1/events/stream, registered
|
|
// directly on the chi router in NewHandler so it supersedes the generated
|
|
// route. Using the raw http.ResponseWriter lets us Flush() after every event
|
|
// (real-time delivery); the generated strict-server path can only hand back
|
|
// an io.Reader that io.Copy drains without flushing (chunk-buffered).
|
|
//
|
|
// Protocol: https://html.spec.whatwg.org/multipage/server-sent-events.html
|
|
// 1. If Last-Event-ID is present, replay buffered events (broker, then DB).
|
|
// 2. Subscribe and forward events fanned out from pg_notify.
|
|
// 3. Colon-comment heartbeat every 15s.
|
|
// 4. Unsubscribe on client disconnect (request context cancels).
|
|
func (s *Server) serveSSE(w http.ResponseWriter, r *http.Request) {
|
|
flusher, ok := w.(http.Flusher)
|
|
if !ok {
|
|
writeProblem(w, r, http.StatusInternalServerError, "internal error", "streaming not supported")
|
|
return
|
|
}
|
|
|
|
w.Header().Set("Content-Type", "text/event-stream")
|
|
w.Header().Set("Cache-Control", "no-cache")
|
|
w.Header().Set("Connection", "keep-alive")
|
|
w.Header().Set("X-Accel-Buffering", "no") // disable proxy buffering
|
|
w.WriteHeader(http.StatusOK)
|
|
flusher.Flush()
|
|
|
|
ctx, cancel := context.WithCancel(r.Context())
|
|
defer cancel()
|
|
|
|
sub := &sseSubscriber{
|
|
ch: make(chan sqlcgen.Event, 64),
|
|
done: make(chan struct{}),
|
|
cancel: cancel,
|
|
}
|
|
s.sseMu.Lock()
|
|
s.sseSubs[sub] = struct{}{}
|
|
s.sseMu.Unlock()
|
|
defer func() {
|
|
s.sseMu.Lock()
|
|
delete(s.sseSubs, sub)
|
|
s.sseMu.Unlock()
|
|
close(sub.done)
|
|
}()
|
|
|
|
// ── 1. Replay on Last-Event-ID ────────────────────────────────
|
|
if lastID := r.Header.Get("Last-Event-ID"); lastID != "" {
|
|
if id, err := strconv.ParseInt(lastID, 10, 64); err == nil {
|
|
// The in-memory broker only holds recent events; if it doesn't
|
|
// cover the whole gap, fall back to the DB for a complete replay.
|
|
events := s.sseBroker.after(id)
|
|
complete := len(events) > 0 && events[len(events)-1].ID == s.sseBroker.latestID()
|
|
if complete {
|
|
for _, ev := range events {
|
|
if !writeSSE(w, flusher, ev) {
|
|
return
|
|
}
|
|
}
|
|
} else {
|
|
q := sqlcgen.New(s.pool)
|
|
dbEvents, err := q.ListEventsAfter(ctx, sqlcgen.ListEventsAfterParams{ID: id, Limit: 5000})
|
|
if err != nil {
|
|
slog.Error("sse db replay failed", "error", err)
|
|
} else {
|
|
for _, ev := range dbEvents {
|
|
if !writeSSE(w, flusher, ev) {
|
|
return
|
|
}
|
|
}
|
|
}
|
|
}
|
|
}
|
|
}
|
|
|
|
// ── 2. Subscribe and forward ──────────────────────────────────
|
|
heartbeat := time.NewTicker(15 * time.Second)
|
|
defer heartbeat.Stop()
|
|
|
|
for {
|
|
select {
|
|
case <-ctx.Done():
|
|
return
|
|
case ev, ok := <-sub.ch:
|
|
if !ok {
|
|
return
|
|
}
|
|
if !writeSSE(w, flusher, ev) {
|
|
return
|
|
}
|
|
case <-heartbeat.C:
|
|
if _, err := fmt.Fprintf(w, ": heartbeat\n\n"); err != nil {
|
|
return
|
|
}
|
|
flusher.Flush()
|
|
}
|
|
}
|
|
}
|
|
|
|
// StreamEvents satisfies the generated ServerInterface, but the SSE route is
|
|
// served by the raw serveSSE handler registered in NewHandler (which wins
|
|
// over this generated route). If this method is ever reached, routing has
|
|
// regressed — fail loudly rather than silently chunk-buffering.
|
|
func (s *Server) StreamEvents(ctx context.Context, req gen.StreamEventsRequestObject) (gen.StreamEventsResponseObject, error) {
|
|
return nil, fmt.Errorf("%w: SSE must be served by the raw handler", errNotImplemented)
|
|
}
|
|
|
|
// Ensure pgxpool is imported — used via Acquire.
|
|
var _ = &pgxpool.Pool{}
|
|
var _ = pgx.ErrNoRows
|