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

Go HTTP CONNECT代理转发WebSocket时遇Broken pipe等错误求助

HTTP CONNECT代理WebSocket测试报错问题

我用Go实现了一个简单的HTTP CONNECT代理,核心逻辑是在隧道关闭前,在客户端与目标服务器之间双向盲转发数据包。用curl通过代理访问外部网站时一切正常,但测试自定义的WebSocket客户端和服务端时,代理日志频繁出现连接重置或管道破裂的错误。

代理代码

type ProxyHttpServer struct {
    Addr string
}

func (proxy *ProxyHttpServer) ServeHTTP(w http.ResponseWriter, r *http.Request) {
    if r.Method == http.MethodConnect {
        proxy.handleConnect(w, r)
    }
}

func (proxy *ProxyHttpServer) handleConnect(w http.ResponseWriter, r *http.Request) {
    // 获取底层连接控制权
    hij, ok := w.(http.Hijacker)
    if !ok {
        http.Error(w, "Hijacking not supported", http.StatusInternalServerError)
        return
    }

    // 劫持连接
    clientConn, _, err := hij.Hijack() // 与客户端的连接
    if err != nil {
        http.Error(w, err.Error(), http.StatusInternalServerError)
        return
    }
    defer clientConn.Close()

    // 建立与目标服务器的TCP连接
    upstreamConn, err := net.Dial("tcp", r.Host) // 与目标服务器的连接
    if err != nil {
        http.Error(w, err.Error(), http.StatusServiceUnavailable)
        return
    }
    defer upstreamConn.Close()

    // TCP连接建立完成,向客户端返回200响应
    clientConn.Write([]byte("HTTP/1.1 200 Connection established\r\n\r\n"))

    // 启动双向数据转发
    proxy.relayData(clientConn, upstreamConn)
}

func (proxy *ProxyHttpServer) relayData(responseStream io.ReadWriteCloser, outboundStream io.ReadWriteCloser) {
    wg := sync.WaitGroup{}
    wg.Add(2)

    // 从目标服务器转发到客户端
    go func() {
        defer wg.Done()
        proxy.transfer(responseStream, outboundStream, "outboundStream -> responseStream")
    }()

    // 从客户端转发到目标服务器
    go func() {
        defer wg.Done()
        proxy.transfer(outboundStream, responseStream, "responseStream -> outboundStream")
    }()

    wg.Wait()
}

func (proxy *ProxyHttpServer) transfer(destination io.WriteCloser, source io.ReadCloser, logMessage string) {
    // 从源读取数据并写入目标
    n, copyErr := io.Copy(destination, source)
    if copyErr != nil {
        log.Printf("failed to copy data [%s]: %v", logMessage, copyErr)
        return
    }
    log.Printf("successfully copied %v bytes [%s]", n, logMessage)
}

func main() {
    proxy := ProxyHttpServer{
        Addr: ":8887",
    }

    log.Printf("Listening on port %s", proxy.Addr)

    err := http.ListenAndServe(":8887", &proxy)
    if err != nil {
        log.Fatal(err)
    }
}

错误日志

运行WebSocket客户端传输数据时,代理输出以下错误:

2024/10/07 13:28:58 Listening on port :8887
2024/10/07 13:29:08 successfully copied 684 bytes [responseStream -> outboundStream]
2024/10/07 13:29:08 failed to copy data [outboundStream -> responseStream]: writeto tcp 127.0.0.1:33088->127.0.0.1:8080: readfrom tcp 127.0.0.1:8887->127.0.0.1:60268: splice: connection reset by peer
2024/10/07 13:29:16 successfully copied 684 bytes [responseStream -> outboundStream]
2024/10/07 13:29:16 failed to copy data [outboundStream -> responseStream]: writeto tcp 127.0.0.1:51142->127.0.0.1:8080: readfrom tcp 127.0.0.1:8887->127.0.0.1:46066: splice: broken pipe
2024/10/07 13:29:21 successfully copied 1343 bytes [outboundStream -> responseStream]
2024/10/07 13:29:21 failed to copy data [responseStream -> outboundStream]: writeto tcp 127.0.0.1:8887->127.0.0.1:42628: readfrom tcp 127.0.0.1:52260->127.0.0.1:8080: splice: connection reset by peer
2024/10/07 13:29:28 successfully copied 684 bytes [responseStream -> outboundStream]
2024/10/07 13:29:28 failed to copy data [outboundStream -> responseStream]: writeto tcp 127.0.0.1:52268->127.0.0.1:8080: readfrom tcp 127.0.0.1:8887->127.0.0.1:42632: splice: broken pipe

WebSocket服务端代码(server.go)

func main() {
    http.HandleFunc("/ws", func(w http.ResponseWriter, r *http.Request) {
        conn, err := websocket.Accept(w, r, &websocket.AcceptOptions{
            InsecureSkipVerify: true, // 测试时关闭验证
        })
        if err != nil {
            log.Println("Failed to accept connection:", err)
            return
        }
        defer conn.Close(websocket.StatusInternalError, "Internal server error")

        ctx := context.Background()

        var msg string
        err = wsjson.Read(ctx, conn, &msg)
        if err != nil {
            log.Println("Failed to read message:", err)
            return
        }

        log.Println("Received message:", msg)

        err = wsjson.Write(ctx, conn, "Hello, client!")
        if err != nil {
            log.Println("Failed to write message:", err)
            return
        }

        conn.Close(websocket.StatusNormalClosure, "Done")
    })

    // 注册信号处理,优雅关闭服务
    go func() {
        sigint := make(chan os.Signal, 1)
        signal.Notify(sigint, os.Interrupt, syscall.SIGTERM)
        <-sigint

        log.Println("Shutting down server...")
        os.Exit(0)
    }()

    certFile := "./cert.pem"   // SSL证书路径
    keyFile := "./private.key" // SSL私钥路径

    log.Println("Server started on :8080")
    err := http.ListenAndServeTLS(":8080", certFile, keyFile, nil)
    if err != nil {
        log.Fatal("ListenAndServe:", err)
    }
}

WebSocket客户端代码(client.go)

func main() {
    proxyURL, err := url.Parse("http://127.0.0.1:8887")
    if err != nil {
        log.Fatal("Failed to parse proxy URL:", err)
    }

    httpTransport := &http.Transport{
        Proxy: http.ProxyURL(proxyURL),
        TLSClientConfig: &tls.Config{InsecureSkipVerify: true},
    }

    httpClient := &http.Client{
        Transport: httpTransport,
    }

    opts := &websocket.DialOptions{
        HTTPClient: httpClient,
    }

    ctx := context.Background()
    conn, err := Dial(ctx, "wss://localhost:8080/ws", opts)
    if err != nil {
        log.Fatal("%v", err)
    }
    defer conn.Close(websocket.StatusInternalError, "Internal server error")

    err = wsjson.Write(ctx, conn, "Hello, server!")
    if err != nil {
        log.Fatal("Failed to send message:", err)
    }

    var msg string
    err = wsjson.Read(ctx, conn, &msg)
    if err != nil {
        log.Fatal("Failed to read message:", err)
    }

    log.Println("Received message:", msg)

    conn.Close(websocket.StatusNormalClosure, "Done")
}

curl测试验证代理正常

用curl通过代理访问外部网站完全正常,执行命令:

curl -x http://127.0.0.1:8887  https://www.reddit.com/ -v

代理日志输出:

2024/10/07 13:59:41 Listening on port :8887
2024/10/07 14:00:14 successfully copied 883 bytes [responseStream -> outboundStream]
2024/10/07 14:00:14 successfully copied 27470 bytes [outboundStream -> responseStream]

问题分析与解决方法

问题本质

报错是因为WebSocket会话结束时,一端已关闭连接,但代理的转发协程仍在尝试读写:

  1. WebSocket服务端发送响应后立刻主动调用conn.Close()关闭连接,此时代理从服务端读取数据的协程会触发connection reset by peer。
  2. 客户端收到响应后也主动关闭连接,代理向客户端写入数据的协程会触发broken pipe。

curl测试正常是因为curl与目标服务器的连接关闭流程更优雅,会确保所有数据传输完成后再关闭连接。

解决方法

1. 忽略正常关闭的错误

在代理的transfer函数中,过滤掉连接重置、管道破裂这类正常关闭的错误:

import "syscall"

func (proxy *ProxyHttpServer) transfer(destination io.WriteCloser, source io.ReadCloser, logMessage string) {
    n, copyErr := io.Copy(destination, source)
    if copyErr != nil {
        // 判断错误类型,忽略正常关闭的错误
        if opErr, ok := copyErr.(*net.OpError); ok {
            if sysErr, ok := opErr.Err.(syscall.Errno); ok {
                if sysErr == syscall.ECONNRESET || sysErr == syscall.EPIPE {
                    log.Printf("Connection closed normally [%s]", logMessage)
                    return
                }
            }
        }
        log.Printf("failed to copy data [%s]: %v", logMessage, copyErr)
        return
    }
    log.Printf("successfully copied %v bytes [%s]", n, logMessage)
}
2. 优化WebSocket服务端/客户端的关闭逻辑

避免重复关闭连接,让会话自然结束:

  • 服务端:删除主动调用的conn.Close(websocket.StatusNormalClosure, "Done"),仅保留defer中的正常关闭:
    defer conn.Close(websocket.StatusNormalClosure, "Done")
    
  • 客户端:同样删除末尾的conn.Close(),依赖defer完成正常关闭。
3. 优化代理转发逻辑

当任意一端的转发协程结束时,主动关闭另一端的连接,避免无效读写:

import "context"

func (proxy *ProxyHttpServer) relayData(responseStream io.ReadWriteCloser, outboundStream io.ReadWriteCloser) {
    ctx, cancel := context.WithCancel(context.Background())
    defer cancel()
    // 确保两端连接最终关闭
    defer responseStream.Close()
    defer outboundStream.Close()

    wg := sync.WaitGroup{}
    wg.Add(2)

    // 转发目标服务器到客户端
    go func() {
        defer wg.Done()
        defer cancel() // 结束时通知另一个协程退出
        proxy.transfer(responseStream, outboundStream, "outboundStream -> responseStream")
    }()

    // 转发客户端到目标服务器
    go func() {
        defer wg.Done()
        defer cancel() // 结束时通知另一个协程退出
        proxy.transfer(outboundStream, responseStream, "responseStream -> outboundStream")
    }()

    // 等待任意协程结束或context被取消
    select {
    case <-ctx.Done():
        // 主动关闭两端连接,终止剩余转发
        responseStream.Close()
        outboundStream.Close()
    case <-func() chan struct{} {
        ch := make(chan struct{})
        go func() {
            wg.Wait()
            close(ch)
        }()
        return ch
    }():
    }
}

内容的提问来源于stack exchange,提问作者ray an

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 14:22:03