fix(service): 保持直接响应的实时流式写入(by AI)
This commit is contained in:
parent
548aecf156
commit
0a3fe19a56
@ -1,5 +1,10 @@
|
||||
# CHANGELOG - go/service
|
||||
|
||||
## v1.5.23 (2026-08-17)
|
||||
- **流式响应**: `Response.Write` 保持直接写入和实时 Flush,不再因服务注册了输出过滤器而被统一缓冲。
|
||||
- **显式过滤**: 新增 `Response.WriteFiltered`,静态文件等确实需要输出过滤器处理的内容可显式进入缓冲流程。
|
||||
- **低代码字节桥接**: 新增 `Response.WriteBytes`,可将动态运行时的字节数组无损写入 SSE 或二进制响应。
|
||||
|
||||
## v1.5.22 (2026-07-19)
|
||||
- **服务发现升级**: 升级 `discover` 至 `v1.5.8`,继承 Redis 类型化注册、心跳、注销和节点拉取流程。
|
||||
- **依赖对齐**: 随模块图同步升级 `cast`、`config`、`crypto`、`encoding`、`file`、`http`、`id`、`log`、`rand`、`safe`、`shell` 的 patch 版本。
|
||||
|
||||
@ -80,6 +80,8 @@ func main() {
|
||||
- **文档生成**: `service.MakeDocument()` 返回全量接口描述
|
||||
- **依赖注入**: `service.GetInjectT[T]()` 快速获取已注入的对象或组件
|
||||
|
||||
低代码需要转发二进制或 SSE 流时,使用 `response.WriteBytes(chunk)` 原样写入动态运行时提供的字节数组;文本输出继续使用 `response.WriteString(text)`。`Write`、`WriteBytes` 和 `WriteString` 都直接写入响应并保持流式行为,不会被全局输出过滤器强制缓冲;确实需要过滤器处理的内容应显式使用 `response.WriteFiltered(bytes)`。
|
||||
|
||||
## 配置指南 (ServiceConfig)
|
||||
|
||||
详细配置项可查阅 `config.go` 中的 `ServiceConfig` 结构。通过 `config.Load` 支持从 `env.yml` 或环境变量加载。
|
||||
|
||||
13
TEST.md
13
TEST.md
@ -1,5 +1,11 @@
|
||||
# Service Module Test Report
|
||||
|
||||
## v1.5.23 验证
|
||||
- `TestResponseWriteBytes` 验证动态字节数组写入时保持 UTF-8 内容不变。
|
||||
- `TestResponseWriteBypassesOutputFilterBuffer` 验证直接响应即使存在输出过滤器也保持实时写入。
|
||||
- `TestResponseWriteFilteredOptsIntoOutputFilterBuffer` 验证显式过滤写入仍可由静态 HTML 注入等后置过滤器处理。
|
||||
- `go test -v ./...` 与 `go test -bench=. ./...` 全部通过。
|
||||
|
||||
## v1.5.22 验证
|
||||
- `go list -m all`:模块图使用 `discover v1.5.8`、`redis v1.5.11` 与稳定版 `redigo v1.9.3`。
|
||||
- `go test -v ./...` 与 `go test -bench=. ./...` 全部通过。
|
||||
@ -9,9 +15,9 @@
|
||||
- `go test -v ./...` 与 `go test -bench=. ./...` 全部通过。
|
||||
|
||||
## 性能测试 (Benchmark)
|
||||
- 测试日期: 2026-07-19
|
||||
- 版本: v1.5.22
|
||||
- 指标: `BenchmarkRouting`: **2269 ns/op**
|
||||
- 测试日期: 2026-08-17
|
||||
- 版本: v1.5.23
|
||||
- 指标: `BenchmarkRouting`: **2312 ns/op**
|
||||
- 环境: Darwin / Apple M3 Max
|
||||
|
||||
## 单元测试覆盖 (Unit Test)
|
||||
@ -39,6 +45,7 @@
|
||||
- [x] Client Keys: Device-Id/Session-Id 应答头条件化(请求有则不应答)、静态文件/WebSocket 仅 Cookie
|
||||
- [x] Cookie 头智能过滤: 排除列表中 key 从 Cookie 内容中剔除,保留业务 Cookie
|
||||
- [x] Response Body: 200 响应和 dev 模式下 keepBody 捕获
|
||||
- [x] Response Dynamic Bytes: 低代码动态数组通过 `WriteBytes` 写入时保持 UTF-8 字节不变
|
||||
|
||||
## 基础设施对齐验证
|
||||
- [x] 成功集成 `apigo.cc/go/cast` 用于参数解析与类型强转。
|
||||
|
||||
@ -204,7 +204,12 @@ func (rh *RouteHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
|
||||
}
|
||||
|
||||
filter:
|
||||
// 7. 后置过滤器 (即使 response.changed 也要执行,比如静态文件的 HTML 注入)
|
||||
// Direct response writes are already on the wire and intentionally bypass
|
||||
// output filters. Static responses opt in through WriteFiltered.
|
||||
if response.changed && !response.filteredWrite && result == nil {
|
||||
return
|
||||
}
|
||||
// 7. 后置过滤器
|
||||
for _, filter := range ws.outFilters {
|
||||
newResult, done := filter(args, request, response, result, requestLogger)
|
||||
if newResult != nil {
|
||||
@ -223,7 +228,7 @@ filter:
|
||||
// 过滤器模式:所有内容都应该从 result 或 response.body 中写出
|
||||
if result != nil {
|
||||
outputResult(response, result)
|
||||
} else if response.changed {
|
||||
} else if response.filteredWrite {
|
||||
response.PhysicalWrite(response.body)
|
||||
}
|
||||
} else {
|
||||
|
||||
100
response.go
100
response.go
@ -4,8 +4,10 @@ import (
|
||||
"apigo.cc/go/cast"
|
||||
"apigo.cc/go/file"
|
||||
"apigo.cc/go/jsmod"
|
||||
"fmt"
|
||||
"io"
|
||||
"net/http"
|
||||
"reflect"
|
||||
)
|
||||
|
||||
// Response 封装 http.ResponseWriter
|
||||
@ -13,14 +15,15 @@ type Response struct {
|
||||
Id string
|
||||
Writer http.ResponseWriter `js:"-"`
|
||||
Code int
|
||||
body []byte `js:"-"`
|
||||
outLen int `js:"-"`
|
||||
changed bool `js:"-"`
|
||||
headerWritten bool `js:"-"`
|
||||
dontLog200 bool `js:"-"`
|
||||
dontLogArgs []string `js:"-"`
|
||||
ProxyHeader *http.Header `js:"-"`
|
||||
server *WebServer `js:"-"`
|
||||
body []byte `js:"-"`
|
||||
outLen int `js:"-"`
|
||||
changed bool `js:"-"`
|
||||
filteredWrite bool `js:"-"`
|
||||
headerWritten bool `js:"-"`
|
||||
dontLog200 bool `js:"-"`
|
||||
dontLogArgs []string `js:"-"`
|
||||
ProxyHeader *http.Header `js:"-"`
|
||||
server *WebServer `js:"-"`
|
||||
}
|
||||
|
||||
func (r *Response) SetCookie(cookie *Cookie) {
|
||||
@ -61,12 +64,6 @@ func (r *Response) Write(bytes []byte) (int, error) {
|
||||
r.changed = true
|
||||
r.outLen += len(bytes)
|
||||
|
||||
// 如果有输出过滤器,我们必须先缓冲,不能直接写入网线,否则会导致重复输出
|
||||
if r.server != nil && r.server.hasOutFilter {
|
||||
r.body = append(r.body, bytes...)
|
||||
return len(bytes), nil
|
||||
}
|
||||
|
||||
// 缓冲 body 用于日志记录
|
||||
r.keepBody(bytes)
|
||||
|
||||
@ -80,6 +77,81 @@ func (r *Response) Write(bytes []byte) (int, error) {
|
||||
return n, nil
|
||||
}
|
||||
|
||||
// WriteFiltered buffers bytes for registered output filters. Direct Write calls
|
||||
// remain streamable and bypass output filters.
|
||||
func (r *Response) WriteFiltered(bytes []byte) (int, error) {
|
||||
if r.server == nil || !r.server.hasOutFilter {
|
||||
return r.Write(bytes)
|
||||
}
|
||||
r.checkWriteHeader()
|
||||
r.changed = true
|
||||
r.filteredWrite = true
|
||||
r.outLen += len(bytes)
|
||||
r.body = append(r.body, bytes...)
|
||||
return len(bytes), nil
|
||||
}
|
||||
|
||||
// WriteBytes writes a byte array received through a dynamic runtime without
|
||||
// converting its bytes to text first. This keeps UTF-8 stream chunks intact.
|
||||
func (r *Response) WriteBytes(value any) (int, error) {
|
||||
bytes, err := responseBytes(value)
|
||||
if err != nil {
|
||||
return 0, jsmod.MakeError(err)
|
||||
}
|
||||
return r.Write(bytes)
|
||||
}
|
||||
|
||||
func responseBytes(value any) ([]byte, error) {
|
||||
if value == nil {
|
||||
return nil, nil
|
||||
}
|
||||
if bytes, ok := value.([]byte); ok {
|
||||
return bytes, nil
|
||||
}
|
||||
if text, ok := value.(string); ok {
|
||||
return []byte(text), nil
|
||||
}
|
||||
|
||||
rv := reflect.ValueOf(value)
|
||||
if rv.Kind() != reflect.Array && rv.Kind() != reflect.Slice {
|
||||
return nil, fmt.Errorf("response bytes must be a string or byte array, got %T", value)
|
||||
}
|
||||
bytes := make([]byte, rv.Len())
|
||||
for i := 0; i < rv.Len(); i++ {
|
||||
item := rv.Index(i)
|
||||
for item.IsValid() && (item.Kind() == reflect.Interface || item.Kind() == reflect.Pointer) {
|
||||
if item.IsNil() {
|
||||
return nil, fmt.Errorf("response byte at index %d is nil", i)
|
||||
}
|
||||
item = item.Elem()
|
||||
}
|
||||
if !item.IsValid() {
|
||||
return nil, fmt.Errorf("response byte at index %d is nil", i)
|
||||
}
|
||||
switch item.Kind() {
|
||||
case reflect.Uint, reflect.Uint8, reflect.Uint16, reflect.Uint32, reflect.Uint64:
|
||||
if item.Uint() > 255 {
|
||||
return nil, fmt.Errorf("response byte at index %d is out of range", i)
|
||||
}
|
||||
bytes[i] = byte(item.Uint())
|
||||
case reflect.Int, reflect.Int8, reflect.Int16, reflect.Int32, reflect.Int64:
|
||||
if item.Int() < 0 || item.Int() > 255 {
|
||||
return nil, fmt.Errorf("response byte at index %d is out of range", i)
|
||||
}
|
||||
bytes[i] = byte(item.Int())
|
||||
case reflect.Float32, reflect.Float64:
|
||||
n := item.Float()
|
||||
if n < 0 || n > 255 || n != float64(byte(n)) {
|
||||
return nil, fmt.Errorf("response byte at index %d is out of range", i)
|
||||
}
|
||||
bytes[i] = byte(n)
|
||||
default:
|
||||
return nil, fmt.Errorf("response byte at index %d has type %s", i, item.Kind())
|
||||
}
|
||||
}
|
||||
return bytes, nil
|
||||
}
|
||||
|
||||
// keepBody 缓冲数据用于日志记录,限制大小防止内存问题
|
||||
func (r *Response) keepBody(bytes []byte) {
|
||||
maxBuf := 200
|
||||
|
||||
44
response_test.go
Normal file
44
response_test.go
Normal file
@ -0,0 +1,44 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"net/http/httptest"
|
||||
"testing"
|
||||
)
|
||||
|
||||
func TestResponseWriteBytes(t *testing.T) {
|
||||
recorder := httptest.NewRecorder()
|
||||
response := NewResponse(recorder, nil)
|
||||
want := "data: 第一条\n\n"
|
||||
input := make([]any, len([]byte(want)))
|
||||
for i, value := range []byte(want) {
|
||||
input[i] = int64(value)
|
||||
}
|
||||
if _, err := response.WriteBytes(input); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if got := recorder.Body.String(); got != want {
|
||||
t.Fatalf("body = %q, want %q", got, want)
|
||||
}
|
||||
}
|
||||
|
||||
func TestResponseWriteBypassesOutputFilterBuffer(t *testing.T) {
|
||||
recorder := httptest.NewRecorder()
|
||||
response := NewResponse(recorder, &WebServer{hasOutFilter: true})
|
||||
if _, err := response.Write([]byte("stream")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if recorder.Body.String() != "stream" || response.filteredWrite {
|
||||
t.Fatalf("direct write was buffered: wire=%q buffer=%q", recorder.Body.String(), response.body)
|
||||
}
|
||||
}
|
||||
|
||||
func TestResponseWriteFilteredOptsIntoOutputFilterBuffer(t *testing.T) {
|
||||
recorder := httptest.NewRecorder()
|
||||
response := NewResponse(recorder, &WebServer{hasOutFilter: true})
|
||||
if _, err := response.WriteFiltered([]byte("html")); err != nil {
|
||||
t.Fatal(err)
|
||||
}
|
||||
if recorder.Body.Len() != 0 || !response.filteredWrite || string(response.body) != "html" {
|
||||
t.Fatalf("filtered write was not buffered: wire=%q buffer=%q", recorder.Body.String(), response.body)
|
||||
}
|
||||
}
|
||||
@ -116,7 +116,7 @@ func (ws *WebServer) processStatic(requestPath string, request *Request, respons
|
||||
if len(indexFiles) == 0 {
|
||||
indexFiles = []string{"index.html", "index.htm"}
|
||||
}
|
||||
|
||||
|
||||
for _, indexFile := range indexFiles {
|
||||
f := filepath.Join(filePath, indexFile)
|
||||
if i := file.GetFileInfo(f); i != nil && !i.IsDir {
|
||||
@ -159,6 +159,6 @@ func (ws *WebServer) processStatic(requestPath string, request *Request, respons
|
||||
return false
|
||||
}
|
||||
|
||||
_, _ = response.Write(data)
|
||||
_, _ = response.WriteFiltered(data)
|
||||
return true
|
||||
}
|
||||
|
||||
Loading…
x
Reference in New Issue
Block a user