package httpapi // Real-connection SSE test. httptest.NewRecorder buffers and never flushes, // so this uses httptest.NewServer + a streaming client to verify that: // - the raw serveSSE handler (not the generated 501 stub) serves the route, // - an event created *after* the client connects is delivered in real time // (i.e. flushed before the connection closes), // - the SSE `data:` payload is the canonical gen.Event shape. // Guarded by OIKOS_TEST_DATABASE_URL. import ( "bufio" "bytes" "context" "encoding/json" "net/http" "net/http/httptest" "strings" "testing" "time" ) func TestSSEStreamRealtimeDelivery(t *testing.T) { srv := httptest.NewServer(newTestHandler(t, devConfig())) defer srv.Close() ctx, cancel := context.WithTimeout(context.Background(), 10*time.Second) defer cancel() // Connect to the stream. req, _ := http.NewRequestWithContext(ctx, "GET", srv.URL+"/api/v1/events/stream", nil) req.Header.Set("Authorization", "Bearer "+testAuthToken) resp, err := http.DefaultClient.Do(req) if err != nil { t.Fatalf("connect stream: %v", err) } defer resp.Body.Close() if resp.StatusCode != 200 { t.Fatalf("stream status = %d, want 200", resp.StatusCode) } if ct := resp.Header.Get("Content-Type"); !strings.HasPrefix(ct, "text/event-stream") { t.Fatalf("content-type = %q, want text/event-stream", ct) } // Read SSE frames in a goroutine. dataCh := make(chan map[string]any, 4) go func() { sc := bufio.NewScanner(resp.Body) for sc.Scan() { line := sc.Text() if strings.HasPrefix(line, "data: ") { var m map[string]any if json.Unmarshal([]byte(strings.TrimPrefix(line, "data: ")), &m) == nil { dataCh <- m } } } }() // Give the subscriber a moment to register, then trigger an event by // POSTing to the SAME live server (same DB → NOTIFY the listener sees). time.Sleep(300 * time.Millisecond) payload, _ := json.Marshal(map[string]any{"slug": "service:sse-rt", "type": "service", "name": "sse-rt"}) createReq, _ := http.NewRequestWithContext(ctx, "POST", srv.URL+"/api/v1/entities", bytes.NewReader(payload)) createReq.Header.Set("Content-Type", "application/json") createReq.Header.Set("Authorization", "Bearer "+testAuthToken) cResp, err := http.DefaultClient.Do(createReq) if err != nil { t.Fatalf("trigger create: %v", err) } cResp.Body.Close() if cResp.StatusCode != 201 { t.Fatalf("trigger create status = %d, want 201", cResp.StatusCode) } // The event must arrive in real time (well before the 10s ctx deadline), // proving the handler flushes rather than buffering until close. select { case ev := <-dataCh: if ev["type"] != "entity.created" { t.Errorf("event type = %v, want entity.created", ev["type"]) } // canonical shape: snake_case + decoded data object if _, ok := ev["entity_id"]; !ok { t.Errorf("missing snake_case entity_id: %v", ev) } if d, ok := ev["data"].(map[string]any); !ok || d["slug"] != "service:sse-rt" { t.Errorf("data not a decoded object with slug: %v", ev["data"]) } case <-time.After(3 * time.Second): t.Fatal("SSE event not delivered within 3s (flushing broken?)") } }