mirror of
https://github.com/zeromicro/go-zero.git
synced 2026-05-14 02:10:00 +08:00
feat: support http stream response (#4055)
This commit is contained in:
@@ -5,6 +5,7 @@ import (
|
|||||||
"encoding/json"
|
"encoding/json"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"io"
|
||||||
"net/http"
|
"net/http"
|
||||||
"sync"
|
"sync"
|
||||||
|
|
||||||
@@ -88,6 +89,26 @@ func SetOkHandler(handler func(context.Context, any) any) {
|
|||||||
okHandler = handler
|
okHandler = handler
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// Stream writes data into w with streaming mode.
|
||||||
|
// The ctx is used to control the streaming loop, typically use r.Context().
|
||||||
|
// The fn is called repeatedly until it returns false.
|
||||||
|
func Stream(ctx context.Context, w http.ResponseWriter, fn func(w io.Writer) bool) {
|
||||||
|
for {
|
||||||
|
select {
|
||||||
|
case <-ctx.Done():
|
||||||
|
return
|
||||||
|
default:
|
||||||
|
hasMore := fn(w)
|
||||||
|
if flusher, ok := w.(http.Flusher); ok {
|
||||||
|
flusher.Flush()
|
||||||
|
}
|
||||||
|
if !hasMore {
|
||||||
|
return
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// WriteJson writes v as json string into w with code.
|
// WriteJson writes v as json string into w with code.
|
||||||
func WriteJson(w http.ResponseWriter, code int, v any) {
|
func WriteJson(w http.ResponseWriter, code int, v any) {
|
||||||
if err := doWriteJson(w, code, v); err != nil {
|
if err := doWriteJson(w, code, v); err != nil {
|
||||||
|
|||||||
@@ -1,10 +1,13 @@
|
|||||||
package httpx
|
package httpx
|
||||||
|
|
||||||
import (
|
import (
|
||||||
|
"bytes"
|
||||||
"context"
|
"context"
|
||||||
"errors"
|
"errors"
|
||||||
"fmt"
|
"fmt"
|
||||||
|
"io"
|
||||||
"net/http"
|
"net/http"
|
||||||
|
"net/http/httptest"
|
||||||
"strings"
|
"strings"
|
||||||
"testing"
|
"testing"
|
||||||
|
|
||||||
@@ -239,6 +242,60 @@ func TestWriteJsonMarshalFailed(t *testing.T) {
|
|||||||
assert.Equal(t, http.StatusInternalServerError, w.code)
|
assert.Equal(t, http.StatusInternalServerError, w.code)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestStream(t *testing.T) {
|
||||||
|
t.Run("regular case", func(t *testing.T) {
|
||||||
|
channel := make(chan string)
|
||||||
|
go func() {
|
||||||
|
defer close(channel)
|
||||||
|
for index := 0; index < 5; index++ {
|
||||||
|
channel <- fmt.Sprintf("%d", index)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
w := httptest.NewRecorder()
|
||||||
|
Stream(context.Background(), w, func(w io.Writer) bool {
|
||||||
|
output, ok := <-channel
|
||||||
|
if !ok {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
outputBytes := bytes.NewBufferString(output)
|
||||||
|
_, err := w.Write(append(outputBytes.Bytes(), []byte("\n")...))
|
||||||
|
return err == nil
|
||||||
|
})
|
||||||
|
|
||||||
|
assert.Equal(t, http.StatusOK, w.Code)
|
||||||
|
assert.Equal(t, "0\n1\n2\n3\n4\n", w.Body.String())
|
||||||
|
})
|
||||||
|
|
||||||
|
t.Run("context done", func(t *testing.T) {
|
||||||
|
channel := make(chan string)
|
||||||
|
go func() {
|
||||||
|
defer close(channel)
|
||||||
|
for index := 0; index < 5; index++ {
|
||||||
|
channel <- fmt.Sprintf("num: %d", index)
|
||||||
|
}
|
||||||
|
}()
|
||||||
|
|
||||||
|
w := httptest.NewRecorder()
|
||||||
|
ctx, cancel := context.WithCancel(context.Background())
|
||||||
|
cancel()
|
||||||
|
Stream(ctx, w, func(w io.Writer) bool {
|
||||||
|
output, ok := <-channel
|
||||||
|
if !ok {
|
||||||
|
return false
|
||||||
|
}
|
||||||
|
|
||||||
|
outputBytes := bytes.NewBufferString(output)
|
||||||
|
_, err := w.Write(append(outputBytes.Bytes(), []byte("\n")...))
|
||||||
|
return err == nil
|
||||||
|
})
|
||||||
|
|
||||||
|
assert.Equal(t, http.StatusOK, w.Code)
|
||||||
|
assert.Equal(t, "", w.Body.String())
|
||||||
|
})
|
||||||
|
}
|
||||||
|
|
||||||
type tracedResponseWriter struct {
|
type tracedResponseWriter struct {
|
||||||
headers map[string][]string
|
headers map[string][]string
|
||||||
builder strings.Builder
|
builder strings.Builder
|
||||||
|
|||||||
Reference in New Issue
Block a user