convert streams to SSE #336
This commit is contained in:
@@ -155,7 +155,7 @@ func (s *public) infoStreamCtrl(w http.ResponseWriter, r *http.Request) {
|
||||
fn := func() steamEventFn {
|
||||
lastTS := time.Time{}
|
||||
lastCount := 0
|
||||
return func() (data []byte, upd bool, err error) {
|
||||
return func() (event string, data []byte, upd bool, err error) {
|
||||
key := cache.NewKey(locator.SiteID).ID(URLKey(r)).Scopes(locator.SiteID, locator.URL)
|
||||
data, err = s.cache.Get(key, func() ([]byte, error) {
|
||||
info, e := s.dataService.Info(locator, s.readOnlyAge)
|
||||
@@ -172,10 +172,10 @@ func (s *public) infoStreamCtrl(w http.ResponseWriter, r *http.Request) {
|
||||
return encodeJSONWithHTML(info)
|
||||
})
|
||||
if err != nil {
|
||||
return data, false, err
|
||||
return "info", data, false, err
|
||||
}
|
||||
|
||||
return data, upd, nil
|
||||
return "info", data, upd, nil
|
||||
}
|
||||
}
|
||||
|
||||
@@ -234,7 +234,7 @@ func (s *public) lastCommentsStreamCtrl(w http.ResponseWriter, r *http.Request)
|
||||
|
||||
fn := func() steamEventFn {
|
||||
sinceTime := time.Now()
|
||||
return func() (data []byte, upd bool, err error) {
|
||||
return func() (event string, data []byte, upd bool, err error) {
|
||||
key := cache.NewKey(siteID).ID(URLKey(r)).Scopes(lastCommentsScope)
|
||||
data, err = s.cache.Get(key, func() ([]byte, error) {
|
||||
comments, e := s.dataService.Last(siteID, 1, sinceTime, rest.GetUserOrEmpty(r))
|
||||
@@ -248,7 +248,7 @@ func (s *public) lastCommentsStreamCtrl(w http.ResponseWriter, r *http.Request)
|
||||
sinceTime = time.Now()
|
||||
return encodeJSONWithHTML(comments)
|
||||
})
|
||||
return data, upd, err
|
||||
return "last", data, upd, err
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -569,10 +569,11 @@ func TestRest_InfoStream(t *testing.T) {
|
||||
assert.Equal(t, 200, code)
|
||||
wg.Wait()
|
||||
|
||||
t.Logf(string(body))
|
||||
recs := strings.Split(strings.TrimSuffix(string(body), "\n"), "\n")
|
||||
require.Equal(t, 10, len(recs), "10 records")
|
||||
assert.True(t, strings.Contains(recs[0], `"count":2`), recs[0])
|
||||
assert.True(t, strings.Contains(recs[9], `"count":11`), recs[9])
|
||||
require.Equal(t, 10*3, len(recs), "10 records. each 2 lines +1 emty line")
|
||||
assert.True(t, strings.Contains(recs[0+1], `"count":2`), recs[0])
|
||||
assert.True(t, strings.Contains(recs[9*3+1], `"count":11`), recs[9])
|
||||
|
||||
_, code = get(t, ts.URL+"/api/v1/stream/info?site=radio-t&url=https://radio-t.com/blah123")
|
||||
assert.Equal(t, 500, code)
|
||||
@@ -659,9 +660,9 @@ func TestRest_InfoStreamCancel(t *testing.T) {
|
||||
wg.Wait()
|
||||
|
||||
recs := strings.Split(strings.TrimSuffix(string(body), "\n"), "\n")
|
||||
require.Equal(t, 2, len(recs), "should have 2 records")
|
||||
assert.True(t, strings.Contains(recs[0], `"count":2`), recs[0])
|
||||
assert.True(t, strings.Contains(recs[1], `"count":3`), recs[1])
|
||||
require.Equal(t, 2*3, len(recs), "should have 2 events")
|
||||
assert.True(t, strings.Contains(recs[0*3+1], `"count":2`), recs[0])
|
||||
assert.True(t, strings.Contains(recs[1*3+1], `"count":3`), recs[1])
|
||||
}
|
||||
|
||||
func TestRest_Robots(t *testing.T) {
|
||||
@@ -706,11 +707,13 @@ func TestRest_LastCommentsStream(t *testing.T) {
|
||||
assert.Equal(t, 200, r.StatusCode)
|
||||
|
||||
wg.Wait()
|
||||
t.Logf("headers: %+v", r.Header)
|
||||
assert.Equal(t, "text/event-stream", r.Header.Get("content-type"))
|
||||
|
||||
recs := strings.Split(strings.TrimSuffix(string(body), "\n"), "\n")
|
||||
require.Equal(t, 9, len(recs), "9 records")
|
||||
require.Equal(t, 9*3, len(recs), "9 events")
|
||||
t.Logf("%v", recs)
|
||||
assert.True(t, strings.Contains(recs[0], `test 123`), recs[0])
|
||||
assert.True(t, strings.Contains(recs[1], `test 123`), recs[1])
|
||||
}
|
||||
|
||||
func TestRest_LastCommentsStreamTimeout(t *testing.T) {
|
||||
@@ -765,8 +768,8 @@ func TestRest_LastCommentsStreamCancel(t *testing.T) {
|
||||
wg.Wait()
|
||||
|
||||
recs := strings.Split(strings.TrimSuffix(string(body), "\n"), "\n")
|
||||
require.Equal(t, 2, len(recs), "2 records")
|
||||
assert.True(t, strings.Contains(recs[0], `test 123`), recs[0])
|
||||
require.Equal(t, 2*3, len(recs), "2 events")
|
||||
assert.True(t, strings.Contains(recs[0+1], `test 123`), recs[0+1])
|
||||
}
|
||||
|
||||
func TestRest_LastCommentsStreamTooMany(t *testing.T) {
|
||||
|
||||
@@ -2,6 +2,7 @@ package api
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"sync/atomic"
|
||||
@@ -13,18 +14,18 @@ import (
|
||||
|
||||
// Streamer creates endless stream of \n separated json records send to remote client
|
||||
type Streamer struct {
|
||||
TimeOut time.Duration
|
||||
Refresh time.Duration
|
||||
MaxActive int32
|
||||
|
||||
TimeOut time.Duration
|
||||
Refresh time.Duration
|
||||
MaxActive int32
|
||||
activeCount int32
|
||||
}
|
||||
|
||||
type steamEventFn func() (data []byte, upd bool, err error)
|
||||
type steamEventFn func() (event string, data []byte, upd bool, err error)
|
||||
|
||||
type steamEventResp struct {
|
||||
data []byte
|
||||
err error
|
||||
data []byte
|
||||
event string
|
||||
err error
|
||||
}
|
||||
|
||||
// Activate starts blocking function streaming update created by eventFn to ResponseWriter
|
||||
@@ -39,6 +40,10 @@ func (s *Streamer) Activate(ctx context.Context, eventFn func() steamEventFn, w
|
||||
return errors.New("too many streams")
|
||||
}
|
||||
|
||||
if ww, ok := w.(http.ResponseWriter); ok {
|
||||
ww.Header().Set("Content-Type", "text/event-stream")
|
||||
}
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-ctx.Done(): // request closed by remote client
|
||||
@@ -54,7 +59,10 @@ func (s *Streamer) Activate(ctx context.Context, eventFn func() steamEventFn, w
|
||||
if resp.err != nil {
|
||||
return resp.err
|
||||
}
|
||||
if _, e := w.Write(resp.data); e != nil {
|
||||
|
||||
// make server-sent event record
|
||||
// see https://developer.mozilla.org/en-US/docs/Web/API/Server-sent_events/Using_server-sent_events
|
||||
if _, e := fmt.Fprintf(w, "event: %s\ndata: %s\n", resp.event, string(resp.data)); e != nil {
|
||||
return errors.Wrap(e, "send to stream failed")
|
||||
}
|
||||
if fw, okFlush := w.(http.Flusher); okFlush {
|
||||
@@ -78,13 +86,13 @@ func (s *Streamer) eventsCh(ctx context.Context, fn steamEventFn) <-chan steamEv
|
||||
case <-ctx.Done(): // request closed by remote client
|
||||
return
|
||||
case <-tick.C:
|
||||
resp, upd, err := fn()
|
||||
event, resp, upd, err := fn()
|
||||
if err != nil {
|
||||
ch <- steamEventResp{data: nil, err: errors.Wrap(err, "can't get stream data")}
|
||||
ch <- steamEventResp{event: event, data: nil, err: errors.Wrap(err, "can't get stream data")}
|
||||
return
|
||||
}
|
||||
if upd {
|
||||
ch <- steamEventResp{data: resp, err: nil}
|
||||
ch <- steamEventResp{event: event, data: resp, err: nil}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -19,19 +19,19 @@ func TestStream_Timeout(t *testing.T) {
|
||||
|
||||
eventFn := func() steamEventFn {
|
||||
n := 0
|
||||
return func() (data []byte, upd bool, err error) {
|
||||
return func() (event string, data []byte, upd bool, err error) {
|
||||
n++
|
||||
if n%2 == 0 || n > 10 {
|
||||
return nil, false, nil
|
||||
return "test", nil, false, nil
|
||||
}
|
||||
return []byte(fmt.Sprintf("some data %d\n", n)), true, nil
|
||||
return "test", []byte(fmt.Sprintf("some data %d\n", n)), true, nil
|
||||
}
|
||||
}
|
||||
|
||||
buf := bytes.Buffer{}
|
||||
err := s.Activate(context.Background(), eventFn, &buf)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, "some data 1\nsome data 3\nsome data 5\nsome data 7\nsome data 9\n", buf.String())
|
||||
assert.Equal(t, "event: test\ndata: some data 1\n\nevent: test\ndata: some data 3\n\nevent: test\ndata: some data 5\n\nevent: test\ndata: some data 7\n\nevent: test\ndata: some data 9\n\n", buf.String())
|
||||
}
|
||||
|
||||
func TestStream_Cancel(t *testing.T) {
|
||||
@@ -43,12 +43,12 @@ func TestStream_Cancel(t *testing.T) {
|
||||
|
||||
eventFn := func() steamEventFn {
|
||||
n := 0
|
||||
return func() (data []byte, upd bool, err error) {
|
||||
return func() (event string, data []byte, upd bool, err error) {
|
||||
n++
|
||||
if n%2 == 0 {
|
||||
return nil, false, nil
|
||||
return "test", nil, false, nil
|
||||
}
|
||||
return []byte(fmt.Sprintf("some data %d\n", n)), true, nil
|
||||
return "test", []byte(fmt.Sprintf("some data %d\n", n)), true, nil
|
||||
}
|
||||
}
|
||||
|
||||
@@ -57,5 +57,5 @@ func TestStream_Cancel(t *testing.T) {
|
||||
defer cancel()
|
||||
err := s.Activate(ctx, eventFn, &buf)
|
||||
assert.NoError(t, err)
|
||||
assert.Equal(t, "some data 1\nsome data 3\nsome data 5\nsome data 7\nsome data 9\n", buf.String())
|
||||
assert.Equal(t, "event: test\ndata: some data 1\n\nevent: test\ndata: some data 3\n\nevent: test\ndata: some data 5\n\nevent: test\ndata: some data 7\n\nevent: test\ndata: some data 9\n\n", buf.String())
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user