From 458eb27d8c62b294105861be09359a2fcddc9e8a Mon Sep 17 00:00:00 2001 From: Umputun Date: Mon, 3 Jun 2019 16:44:46 -0500 Subject: [PATCH] add stream api for /last --- README.md | 6 +- backend/app/rest/api/rest.go | 1 + backend/app/rest/api/rest_public.go | 128 ++++++++++++++++------- backend/app/rest/api/rest_public_test.go | 95 ++++++++++++++++- 4 files changed, 188 insertions(+), 42 deletions(-) diff --git a/README.md b/README.md index 19b8b937..db47d40e 100644 --- a/README.md +++ b/README.md @@ -625,7 +625,11 @@ Sort can be `time`, `active` or `score`. Supported sort order with prefix -/+, i ``` * `GET /api/v1/info?site=site-idd&url=post-url` - returns `PostInfo` for site and url -* `GET /api/v1/stream/info?site=site-idd&url=post-url` - returns stream with `PostInfo` records ("\n" separated) for site and url` + +### Streaming API + +* `GET /api/v1/stream/info?site=site-idd&url=post-url` - returns stream with `PostInfo` records ("\n" separated) for the site and url` +* `GET /api/v1/stream/last?site=site-id` - returns updates stream with comments ("\n" separated) for the site` ### RSS feeds diff --git a/backend/app/rest/api/rest.go b/backend/app/rest/api/rest.go index 85955a4d..0d029ea2 100644 --- a/backend/app/rest/api/rest.go +++ b/backend/app/rest/api/rest.go @@ -244,6 +244,7 @@ func (s *Rest) routes() chi.Router { rstream.Use(tollbooth_chi.LimitHandler(tollbooth.NewLimiter(10, nil))) rstream.Use(authMiddleware.Trace, middleware.NoCache, logInfoWithBody) rstream.Get("/info", s.pubRest.infoStreamCtrl) + rstream.Get("/last", s.pubRest.lastCommentsStreamCtrl) }) // open routes, cached diff --git a/backend/app/rest/api/rest_public.go b/backend/app/rest/api/rest_public.go index 5e2f7a10..7be58849 100644 --- a/backend/app/rest/api/rest_public.go +++ b/backend/app/rest/api/rest_public.go @@ -155,10 +155,10 @@ func (s *public) infoStreamCtrl(w http.ResponseWriter, r *http.Request) { locator := store.Locator{SiteID: r.URL.Query().Get("site"), URL: r.URL.Query().Get("url")} log.Printf("[DEBUG] start stream for %+v, timeout=%v, refresh=%v", locator, s.streamTimeOut, s.streamRefresh) - key := cache.NewKey(locator.SiteID).ID(URLKey(r)).Scopes(locator.SiteID, locator.URL) lastTS := time.Time{} lastCount := 0 info := func() (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) if e != nil { @@ -196,7 +196,7 @@ func (s *public) infoStreamCtrl(w http.ResponseWriter, r *http.Request) { return } if _, e := w.Write(resp.data); e != nil { - log.Printf("[WARN] failed to send stream, %v", e) + log.Printf("[WARN] failed to send info stream, %v", e) return } if fw, okFlush := w.(http.Flusher); okFlush { @@ -206,42 +206,6 @@ func (s *public) infoStreamCtrl(w http.ResponseWriter, r *http.Request) { } } -type eventFn func() (data []byte, upd bool, err error) - -type eventResp struct { - data []byte - err error -} - -// populate updates to chan, break on context close -func (s *public) eventsCh(ctx context.Context, fn eventFn) <-chan eventResp { - ch := make(chan eventResp) - go func() { - tick := time.NewTicker(s.streamRefresh) - defer func() { - close(ch) - tick.Stop() - }() - for { - select { - case <-ctx.Done(): // request closed by remote client - log.Printf("[DEBUG] info stream closed by remote client, %v", ctx.Err()) - return - case <-tick.C: - resp, upd, err := fn() - if err != nil { - ch <- eventResp{data: nil, err: errors.Wrap(err, "can't get stream data")} - return - } - if upd { - ch <- eventResp{data: resp, err: nil} - } - } - } - }() - return ch -} - // GET /last/{limit}?site=siteID&since=unix_ts_msec - last comments for the siteID, across all posts, sorted by time, optionally // limited with "since" param func (s *public) lastCommentsCtrl(w http.ResponseWriter, r *http.Request) { @@ -285,6 +249,58 @@ func (s *public) lastCommentsCtrl(w http.ResponseWriter, r *http.Request) { } } +// GET /stream/last?site=siteID& - stream of last comments last comments for the siteID, across all posts +func (s *public) lastCommentsStreamCtrl(w http.ResponseWriter, r *http.Request) { + siteID := r.URL.Query().Get("site") + log.Printf("[DEBUG] get last comments stream for %s", siteID) + + sinceTime := time.Now() + + info := func() (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)) + if e != nil { + return nil, e + } + if len(comments) > 0 { + sinceTime = comments[0].Timestamp + upd = true + } + sinceTime = time.Now() + return encodeJSONWithHTML(comments) + }) + return data, upd, err + } + + updCh := s.eventsCh(r.Context(), info) + + for { + select { + case <-r.Context().Done(): // request closed by remote client + return + case <-time.After(s.streamTimeOut): // request closed by timeout + log.Printf("[DEBUG] last comments stream closed due to timeout") + return + case resp, ok := <-updCh: // new update + if !ok { // closed + return + } + if resp.err != nil { + rest.SendErrorJSON(w, r, http.StatusInternalServerError, resp.err, "can't get last comments", rest.ErrInternal) + return + } + if _, e := w.Write(resp.data); e != nil { + log.Printf("[WARN] failed to send last comments stream, %v", e) + return + } + if fw, okFlush := w.(http.Flusher); okFlush { + fw.Flush() + } + } + } +} + // GET /id/{id}?site=siteID&url=post-url - gets a comment by id func (s *public) commentByIDCtrl(w http.ResponseWriter, r *http.Request) { @@ -506,3 +522,39 @@ func (s *public) applyView(comments []store.Comment, view string) []store.Commen } return comments } + +type eventFn func() (data []byte, upd bool, err error) + +type eventResp struct { + data []byte + err error +} + +// populate updates to chan, break on context close +func (s *public) eventsCh(ctx context.Context, fn eventFn) <-chan eventResp { + ch := make(chan eventResp) + go func() { + tick := time.NewTicker(s.streamRefresh) + defer func() { + close(ch) + tick.Stop() + }() + for { + select { + case <-ctx.Done(): // request closed by remote client + log.Printf("[DEBUG] info stream closed by remote client, %v", ctx.Err()) + return + case <-tick.C: + resp, upd, err := fn() + if err != nil { + ch <- eventResp{data: nil, err: errors.Wrap(err, "can't get stream data")} + return + } + if upd { + ch <- eventResp{data: resp, err: nil} + } + } + } + }() + return ch +} diff --git a/backend/app/rest/api/rest_public_test.go b/backend/app/rest/api/rest_public_test.go index 0e3f0013..51504333 100644 --- a/backend/app/rest/api/rest_public_test.go +++ b/backend/app/rest/api/rest_public_test.go @@ -565,11 +565,9 @@ func TestRest_InfoStreamTimeout(t *testing.T) { postComment(t, ts.URL) st := time.Now() - body, code := get(t, ts.URL+"/api/v1/stream/info?site=radio-t&url=https://radio-t.com/blah1") + _, code := get(t, ts.URL+"/api/v1/stream/info?site=radio-t&url=https://radio-t.com/blah1") assert.Equal(t, 200, code) assert.True(t, time.Since(st) > time.Millisecond*450 && time.Since(st) < time.Millisecond*500, time.Since(st)) - recs := strings.Split(strings.TrimSuffix(string(body), "\n"), "\n") - require.True(t, len(recs) < 10, "not all for 10 streamed, only %d", len(recs)) } func TestRest_InfoStreamCancel(t *testing.T) { @@ -623,6 +621,97 @@ func TestRest_Robots(t *testing.T) { "Allow: /api/v1/list\nAllow: /api/v1/config\nAllow: /api/v1/img\nAllow: /api/v1/avatar\nAllow: /api/v1/picture\n", string(body)) } +func TestRest_LastCommentsStream(t *testing.T) { + ts, srv, teardown := startupT(t) + srv.pubRest.readOnlyAge = 10000000 // make sure we don't hit read-only + srv.pubRest.streamRefresh = 10 * time.Millisecond + srv.pubRest.streamTimeOut = 500 * time.Millisecond + + postComment(t, ts.URL) + + defer teardown() + wg := sync.WaitGroup{} + wg.Add(1) + go func() { + defer wg.Done() + for i := 1; i < 10; i++ { + time.Sleep(100 * time.Millisecond) + postComment(t, ts.URL) + } + }() + + client := http.Client{} + req, err := http.NewRequest("GET", ts.URL+"/api/v1/stream/last?site=radio-t", nil) + require.Nil(t, err) + r, err := client.Do(req) + require.Nil(t, err) + defer r.Body.Close() + body, err := ioutil.ReadAll(r.Body) + require.Nil(t, err) + assert.Equal(t, 200, r.StatusCode) + + wg.Wait() + + recs := strings.Split(strings.TrimSuffix(string(body), "\n"), "\n") + require.Equal(t, 9, len(recs), "9 records") + t.Logf("%v", recs) + assert.True(t, strings.Contains(recs[0], `test 123`), recs[0]) +} + +func TestRest_LastCommentsStreamTimeout(t *testing.T) { + ts, srv, teardown := startupT(t) + defer teardown() + srv.pubRest.readOnlyAge = 10000000 // make sure we don't hit read-only + srv.pubRest.streamRefresh = 10 * time.Millisecond + srv.pubRest.streamTimeOut = 450 * time.Millisecond + + postComment(t, ts.URL) + + st := time.Now() + _, code := get(t, ts.URL+"/api/v1/stream/last?site=radio-t") + assert.Equal(t, 200, code) + assert.True(t, time.Since(st) > time.Millisecond*450 && time.Since(st) < time.Millisecond*500, time.Since(st)) +} + +func TestRest_LastCommentsStreamCancel(t *testing.T) { + ts, srv, teardown := startupT(t) + srv.pubRest.readOnlyAge = 10000000 // make sure we don't hit read-only + srv.pubRest.streamRefresh = 10 * time.Millisecond + srv.pubRest.streamTimeOut = 500 * time.Millisecond + + postComment(t, ts.URL) + + defer teardown() + wg := sync.WaitGroup{} + wg.Add(1) + go func() { + defer wg.Done() + for i := 1; i < 10; i++ { + time.Sleep(100 * time.Millisecond) + postComment(t, ts.URL) + } + }() + + client := http.Client{} + req, err := http.NewRequest("GET", ts.URL+"/api/v1/stream/last?site=radio-t", nil) + require.Nil(t, err) + ctx, cancel := context.WithTimeout(context.Background(), 250*time.Millisecond) + defer cancel() + req = req.WithContext(ctx) + r, err := client.Do(req) + require.Nil(t, err) + defer r.Body.Close() + body, err := ioutil.ReadAll(r.Body) + require.EqualError(t, err, "context deadline exceeded") + assert.Equal(t, 200, r.StatusCode) + + 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]) +} + func postComment(t *testing.T, url string) { resp, e := post(t, url+"/api/v1/comment", `{"text": "test 123", "locator":{"url": "https://radio-t.com/blah1", "site": "radio-t"}}`)