You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在单元测试中测试Server-Sent Event?无循环与time.Sleep的替代方案

单元测试读取SSE响应:避开循环与Sleep的方案

下面是几种不用空循环或time.Sleep就能读取Server-Sent Event(SSE)响应的实用方案:

1. 自定义带通知的ResponseRecorder

默认的http.ResponseRecorder不会在内容写入时触发通知,我们可以封装一个自定义Recorder,在数据写入时通过通道发送就绪信号:

import (
    "bytes"
    "net/http"
    "sync"
)

// NotifyRecorder 带写入通知的自定义响应记录器
type NotifyRecorder struct {
    http.ResponseWriter
    Body     bytes.Buffer
    doneChan chan struct{}
    once     sync.Once
}

func NewNotifyRecorder(w http.ResponseWriter) *NotifyRecorder {
    return &NotifyRecorder{
        ResponseWriter: w,
        doneChan:       make(chan struct{}),
    }
}

func (nr *NotifyRecorder) Write(b []byte) (int, error) {
    n, err := nr.Body.Write(b)
    // 第一次写入数据时关闭通道,发送就绪通知
    nr.once.Do(func() {
        close(nr.doneChan)
    })
    return n, err
}

// Done 返回就绪通知通道
func (nr *NotifyRecorder) Done() <-chan struct{} {
    return nr.doneChan
}

测试中使用这个Recorder:

func (f *fixture) readResponse(t *testing.T) []byte {
    t.Helper()

    timeout := time.After(100 * time.Millisecond)
    select {
    case <-f.w.Done(): // 等待数据写入就绪
        return f.w.Body.Bytes()
    case <-timeout:
        t.Error("未收到响应")
        return nil
    }
}

这个方案通过通道通知替代轮询,第一次写入数据就触发信号,彻底避免空循环和sleep的低效等待。

2. 使用sync.Cond条件变量

利用条件变量sync.Cond,让测试代码在响应写入完成时被主动唤醒,无需轮询:

首先在fixture结构体中添加条件变量:

import (
    "net/http/httptest"
    "sync"
)

type fixture struct {
    w    *httptest.ResponseRecorder
    mu   sync.Mutex
    cond *sync.Cond
}

func newFixture() *fixture {
    w := httptest.NewRecorder()
    f := &fixture{w: w}
    f.cond = sync.NewCond(&f.mu)
    return f
}

在SSE处理函数中,写入响应后唤醒等待的测试goroutine:

import "fmt"

func handleSSE(w http.ResponseWriter, r *http.Request, f *fixture) {
    w.Header().Set("Content-Type", "text/event-stream")
    w.Header().Set("Cache-Control", "no-cache")
    w.Header().Set("Connection", "keep-alive")

    // 写入SSE事件
    _, _ = fmt.Fprintf(w, "data: hello\n\n")
    
    // 刷新响应(确保数据写入缓冲区)
    if flusher, ok := w.(http.Flusher); ok {
        flusher.Flush()
    }

    // 唤醒等待的测试逻辑
    f.mu.Lock()
    f.cond.Signal()
    f.mu.Unlock()
}

修改readResponse函数:

func (f *fixture) readResponse(t *testing.T) []byte {
    t.Helper()

    f.mu.Lock()
    defer f.mu.Unlock()

    done := make(chan struct{})
    go func() {
        f.cond.Wait()
        close(done)
    }()

    select {
    case <-done:
        return f.w.Body.Bytes()
    case <-time.After(100 * time.Millisecond):
        t.Error("未收到响应")
        return nil
    }
}

条件变量让测试代码在特定条件(响应写入完成)满足时被主动唤醒,完全消除无效轮询。

3. 结合流式读取与通知机制

SSE是流式响应,我们可以先通过通知机制等待数据写入,再用bufio.Scanner逐行解析SSE事件:

import (
    "bufio"
    "bytes"
    "strings"
)

func (f *fixture) readResponse(t *testing.T) []byte {
    t.Helper()

    timeout := time.After(100 * time.Millisecond)
    select {
    case <-f.w.Done(): // 先等待数据写入就绪
        reader := bytes.NewReader(f.w.Body.Bytes())
        scanner := bufio.NewScanner(reader)
        
        // 扫描并提取SSE数据行
        for scanner.Scan() {
            line := scanner.Text()
            if strings.HasPrefix(line, "data: ") {
                return []byte(strings.TrimPrefix(line, "data: "))
            }
        }
        
        if err := scanner.Err(); err != nil {
            t.Error("读取响应失败:", err)
            return nil
        }
        t.Error("未找到有效的SSE事件")
        return nil
    case <-timeout:
        t.Error("未收到响应")
        return nil
    }
}

这个方案既保证了等待的高效性,又能精准提取SSE事件内容,适合需要验证具体事件数据的测试场景。

内容的提问来源于stack exchange,提问作者Ado Ren

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.24 15:22:47