Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -142,6 +142,8 @@ kurl graphql https://api.example.com/graphql --generate-query User
```

### 3. 📡 Server-Sent Events (SSE) Streamer (`kurl sse`)

Individual stream lines must be smaller than 1 MiB (including the line ending). Larger lines return an error instead of growing the scanner without a bound.
Stream live event feeds (logs, AI completions, notifications) with real-time timestamps and event-type colorization:
```bash
# Stream all SSE events
Expand Down
53 changes: 53 additions & 0 deletions internal/sse/content_type_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,53 @@
package sse

import (
"context"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"testing"
)

func TestRunSSEValidatesContentTypeBeforeOpeningLog(t *testing.T) {
for _, tc := range []struct {
name, contentType string
valid bool
}{
{"event stream", "text/event-stream", true},
{"charset parameter", "text/event-stream; charset=utf-8", true},
{"JSON", "application/json", false},
{"HTML", "text/html", false},
{"missing", "", false},
{"malformed", "text/event-stream; broken", false},
} {
t.Run(tc.name, func(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header()["Content-Type"] = []string{tc.contentType}
_, _ = w.Write([]byte("data: hello\n\n"))
}))
defer server.Close()
output := filepath.Join(t.TempDir(), "events.log")
if err := os.WriteFile(output, []byte("existing log"), 0600); err != nil {
t.Fatal(err)
}
err := RunSSE(context.Background(), Options{URL: server.URL, OutputFile: output, NoColor: true})
if tc.valid {
if err != nil {
t.Fatal(err)
}
return
}
if err == nil {
t.Error("expected invalid content type error")
}
data, readErr := os.ReadFile(output)
if readErr != nil {
t.Fatal(readErr)
}
if string(data) != "existing log" {
t.Errorf("invalid response overwrote existing log: %q", data)
}
})
}
}
30 changes: 30 additions & 0 deletions internal/sse/event_data_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
package sse

import (
"reflect"
"strings"
"testing"
)

func TestDataDispatchAndPersistentEventID(t *testing.T) {
for _, tc := range []struct {
name, input string
want []Event
}{
{"empty data", "data:\n\n", []Event{{Data: ""}}},
{"leading empty data", "data:\ndata: second\n\n", []Event{{Data: "\nsecond"}}},
{"id-only block", "id: 7\n\ndata: hello\n\ndata: again\n\n", []Event{{ID: "7", Data: "hello"}, {ID: "7", Data: "again"}}},
{"event-only block", "event: stale\n\ndata: hello\n\n", []Event{{Data: "hello"}}},
{"null id ignored", "id: 7\ndata: first\n\nid: bad\x00id\ndata: next\n\n", []Event{{ID: "7", Data: "first"}, {ID: "7", Data: "next"}}},
{"empty id resets", "id: 7\ndata: first\n\nid:\ndata: next\n\n", []Event{{ID: "7", Data: "first"}, {Data: "next"}}},
{"no dispatch at eof", "data: unfinished", nil},
} {
t.Run(tc.name, func(t *testing.T) {
var got []Event
err := ParseStream(strings.NewReader(tc.input), func(e Event) { got = append(got, e) })
if err != nil || !reflect.DeepEqual(got, tc.want) {
t.Fatalf("got %#v, %v; want %#v", got, err, tc.want)
}
})
}
}
48 changes: 48 additions & 0 deletions internal/sse/filter_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,48 @@
package sse

import (
"context"
"net/http"
"net/http/httptest"
"os"
"path/filepath"
"strings"
"testing"
)

func TestRunSSEFiltersEffectiveEventType(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/event-stream")
_, _ = w.Write([]byte("data: default-data\n\nevent: update\ndata: lower-data\n\nevent: Update\ndata: upper-data\n\nevent:\ndata: empty-type-data\n\n"))
}))
defer server.Close()
for _, tc := range []struct {
filter string
want, absent []string
}{
{"update", []string{"lower-data"}, []string{"default-data", "upper-data", "empty-type-data"}},
{"message", []string{"default-data", "empty-type-data"}, []string{"lower-data", "upper-data"}},
{"", []string{"default-data", "empty-type-data", "lower-data", "upper-data"}, nil},
} {
t.Run(tc.filter, func(t *testing.T) {
output := filepath.Join(t.TempDir(), "events.log")
if err := RunSSE(context.Background(), Options{URL: server.URL, FilterType: tc.filter, OutputFile: output, NoColor: true}); err != nil {
t.Fatal(err)
}
data, err := os.ReadFile(output)
if err != nil {
t.Fatal(err)
}
for _, want := range tc.want {
if !strings.Contains(string(data), want) {
t.Errorf("missing %q in %s", want, data)
}
}
for _, absent := range tc.absent {
if strings.Contains(string(data), absent) {
t.Errorf("unexpected %q in %s", absent, data)
}
}
})
}
}
39 changes: 39 additions & 0 deletions internal/sse/headers_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,39 @@
package sse

import (
"context"
"net/http"
"net/http/httptest"
"testing"
)

func TestRunSSEHonorsHostHeader(t *testing.T) {
got := make(chan string, 1)
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
got <- r.Host
w.Header().Set("Content-Type", "text/event-stream")
}))
defer srv.Close()
if err := RunSSE(context.Background(), Options{URL: srv.URL, Headers: []string{"hOsT: events.example"}, NoColor: true}); err != nil {
t.Fatal(err)
}
if host := <-got; host != "events.example" {
t.Fatalf("Host=%q", host)
}
}
func TestRunSSERejectsMalformedHeaderBeforeRequest(t *testing.T) {
called := make(chan struct{}, 1)
srv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
called <- struct{}{}
w.Header().Set("Content-Type", "text/event-stream")
}))
defer srv.Close()
if err := RunSSE(context.Background(), Options{URL: srv.URL, Headers: []string{"MissingColon"}, NoColor: true}); err == nil {
t.Error("malformed header accepted")
}
select {
case <-called:
t.Error("request sent despite malformed header")
default:
}
}
28 changes: 28 additions & 0 deletions internal/sse/large_event_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,28 @@
package sse

import (
"strings"
"testing"
)

func TestParseStreamLargeDataLine(t *testing.T) {
data := strings.Repeat("x", 128*1024)
var events []Event
err := ParseStream(strings.NewReader("data: "+data+"\n\n"), func(event Event) { events = append(events, event) })
if err != nil {
t.Fatal(err)
}
if len(events) != 1 {
t.Fatalf("got %d events, want 1", len(events))
}
if events[0].Data != data {
t.Fatal("large event data was not preserved")
}
}

func TestParseStreamRejectsOversizedLine(t *testing.T) {
err := ParseStream(strings.NewReader("data: "+strings.Repeat("x", 1024*1024)+"\n\n"), func(Event) { t.Error("oversized event was dispatched") })
if err == nil {
t.Fatal("expected bounded scanner to reject oversized line")
}
}
23 changes: 23 additions & 0 deletions internal/sse/line_endings_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
package sse

import (
"strings"
"testing"
"testing/iotest"
)

func TestEventStreamLineEndingsAndBOM(t *testing.T) {
for _, ending := range []string{"\n", "\r\n", "\r"} {
for _, bom := range []string{"", "\ufeff"} {
t.Run(bom+ending, func(t *testing.T) {
// One-byte reads cover CRLF and UTF-8 BOM split across network chunks.
input := bom + strings.Join([]string{"data: first", "", "data: second", "", ""}, ending)
var got []Event
err := ParseStream(iotest.OneByteReader(strings.NewReader(input)), func(e Event) { got = append(got, e) })
if err != nil || len(got) != 2 || got[0].Data != "first" || got[1].Data != "second" {
t.Fatalf("got %+v, %v", got, err)
}
})
}
}
}
43 changes: 43 additions & 0 deletions internal/sse/output_failure_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,43 @@
package sse

import (
"context"
"errors"
"net/http"
"net/http/httptest"
"testing"
"time"
)

type failingEventWriter struct {
calls int
failure error
}

func (w *failingEventWriter) Write(p []byte) (int, error) {
w.calls++
if w.calls > 1 {
return 0, w.failure
}
return len(p), nil
}

func TestSSEStopsWhenEventOutputFails(t *testing.T) {
server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
w.Header().Set("Content-Type", "text/event-stream")
_, _ = w.Write([]byte("data: hello\n\n"))
w.(http.Flusher).Flush()
<-r.Context().Done()
}))
defer server.Close()
ctx, cancel := context.WithTimeout(context.Background(), time.Second)
defer cancel()
failure := errors.New("output unavailable")
writer := &failingEventWriter{failure: failure}
if err := runSSE(ctx, Options{URL: server.URL, NoColor: true}, writer); !errors.Is(err, failure) {
t.Fatalf("got %v, want output failure", err)
}
if ctx.Err() != nil {
t.Fatal("waited for context timeout instead of stopping on write failure")
}
}
30 changes: 30 additions & 0 deletions internal/sse/retry_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,30 @@
package sse

import (
"strings"
"testing"
"time"
)

func TestRetryRequiresNonnegativeIntegerMilliseconds(t *testing.T) {
for _, value := range []string{"-1", "+1", "1.5", "1s", " 10", "", "999999999999999999999999"} {
t.Run(value, func(t *testing.T) {
var got []Event
err := ParseStream(strings.NewReader("retry: 1000\nretry: "+value+"\ndata: test\n\n"), func(e Event) { got = append(got, e) })
if err != nil || len(got) != 1 || got[0].Retry != time.Second {
t.Fatalf("got %+v, %v", got, err)
}
})
}
}

func TestRetryPersistsAcrossEventBoundaries(t *testing.T) {
var got []Event
err := ParseStream(strings.NewReader("retry: 1000\n\ndata: first\n\ndata: second\n\nretry: 0\ndata: third\n\n"), func(e Event) { got = append(got, e) })
if err != nil || len(got) != 3 {
t.Fatalf("got %+v, %v", got, err)
}
if got[0].Retry != time.Second || got[1].Retry != time.Second || got[2].Retry != 0 {
t.Fatalf("retry state lost: %+v", got)
}
}
Loading
Loading