jim800121chen dc6ca211ae feat(ws): tunnel WS forward — 推論結果經 tunnel 推回雲端(後端塊2)
推論工作區後端塊2:實作 WS forward,讓 local agent 的推論結果 WS
(inference:<deviceId>)經 tunnel 推回雲端瀏覽器 canvas overlay。照
edge-ai-platform POC relay/server.go proxyWebSocket 移植(唯讀參考)。

- forwarder.go ForwardWebSocket:OpenStream→寫 upgrade→讀 101→WebSocketConn
  (新增 WebSocketConn/wsUpgradeError/AsWSUpgradeError)
- proxy.go newWebSocketProxyHandler:hijack browser→回寫 101→雙向 io.Copy
  cross-close(無 goroutine leak)+ copyWebSocketUpgradeHeaders(保留
  Sec-WebSocket-*、剝 Authorization/Origin)
- camera.go registerWebSocketRoutes:GET /ws/devices/:id/inference
- api.go wsAuthGroup(/ws + AuthMiddleware):same-origin cookie 認證,
  無 token-in-URL(security 定案)
- stubs.go:移除 WS inference 501 stub、更新過時 doc comment(Mi-1/2)

架構差異:POC 單 binary,visionA api-server(auth)+remote-proxy 雙 binary,
remote-proxy raw byte pipe 透明穿過 WS upgrade bytes。local agent 端未動。

認證/授權(security 定案 + S1/S2):same-origin cookie、剝 Auth/Origin、
WS 不套 300s timeout、走 pickActiveSessionToken(帶別人 deviceId 也只打到
自己 tunnel;多租戶嚴格綁定屬 Phase 1 M2 debt)。

Reviewer 0C/1M/3Mi 通過(修後)。Major-1 已修:收窄 all_endpoints_require_auth_test
的 /ws/ 白名單,讓 authed inference WS 納入「無 cookie 應 401」回歸檢查
(附守得住證明:移除 auth 測試即 FAIL)。+forwarder 5 test + camera_ws
端到端四層真連線雙向 pipe。build/vet/全回歸綠。

Co-Authored-By: Claude Opus 4.8 (1M context) <noreply@anthropic.com>
2026-07-09 04:56:54 +08:00

431 lines
16 KiB
Go
Raw Blame History

This file contains ambiguous Unicode characters

This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.

// forwarder.go — api-server → remote-proxy 的 raw forward client。
//
// 這是雛形雙 binary 架構下「api-server 把前端 HTTP 請求轉發到 local agent」
// 的核心元件。
//
// 整條路徑:
//
// browser ─HTTP─► api-server handler
// │
// │ Forwarder.ForwardHTTP / OpenStream
// ▼
// raw TCP dial remote-proxy: POST /internal/forward/raw?token=...
// │ (B3 Major-1 修復後新增的 hijack endpoint)
// ▼
// remote-proxy hijack 自己的連線 → yamux.OpenStream → 雙向 io.Copy
// │
// ▼
// local agent (yamux client) 把 stream 上的 HTTP request
// ▼
// 轉到本地 127.0.0.1:3721local-tool回 response
//
// 對齊 `.autoflow/04-architecture/api/api-internal.md` §POST /internal/forward/raw
// 與 `.autoflow/04-architecture/tunnel.md` §3.3。
package session
import (
"bufio"
"context"
"errors"
"fmt"
"io"
"log/slog"
"net"
"net/http"
"net/url"
"strings"
"time"
)
// defaultDialTimeout 是 raw TCP dial remote-proxy 的最大等待時間。
const defaultDialTimeout = 10 * time.Second
// defaultHandshakeTimeout 是讀取「HTTP/1.1 200 Connected」握手的最大等待時間。
const defaultHandshakeTimeout = 10 * time.Second
// Forwarder 把 api-server 的 HTTP 請求 forward 到 remote-proxy。
//
// 並發安全:本 struct 的方法不共享可變狀態,每個 OpenStream 走獨立 net.Conn
// 多個 goroutine 可同時呼叫。
type Forwarder struct {
// proxyHost 是從 baseURL 解析出來的 host:port供 net.Dial 用。
proxyHost string
// dialer 用於 raw TCP dial。獨立成欄位以利測試 / 未來換成 TLS dial。
dialer net.Dialer
logger *slog.Logger
}
// NewForwarder 從 baseURLhttp://localhost:3801建立 Forwarder。
//
// baseURL 必須是 http:// 或 https:// 開頭;其他 scheme 視為錯誤但延遲到
// 第一次呼叫時才回(保持建構簽章簡單)。
//
// **注意**:雛形 internal port 是純 HTTPnetwork policy 阻擋外部存取,見
// api-internal.md §安全。Phase 1 加 mTLS 時,本 Forwarder 需擴充支援 TLS。
func NewForwarder(baseURL string, logger *slog.Logger) *Forwarder {
if logger == nil {
logger = slog.Default()
}
host := parseHostFromBaseURL(baseURL)
return &Forwarder{
proxyHost: host,
dialer: net.Dialer{Timeout: defaultDialTimeout},
logger: logger,
}
}
// parseHostFromBaseURL 從 baseURL 取出 host:port失敗時回傳空字串
// (後續 OpenStream 會拒絕並回明確錯誤)。
func parseHostFromBaseURL(baseURL string) string {
if baseURL == "" {
return ""
}
u, err := url.Parse(baseURL)
if err != nil {
return ""
}
return u.Host
}
// OpenStream 對 remote-proxy 開一條 raw TCP 連線,完成 hijack 握手,並回傳
// 一條可以直接用 net.Conn 語意操作的連線(底層是 yamux stream
//
// 用法(典型 api-server handler
//
// conn, err := forwarder.OpenStream(ctx, sessionToken)
// if err != nil { ... }
// defer conn.Close()
//
// httpReq.Write(conn) // 送 HTTP request
// resp, _ := http.ReadResponse(bufio.NewReader(conn), httpReq)
// io.Copy(browserResponseWriter, resp.Body) // streaming friendly
//
// 失敗回傳的 error
// - ErrSessionNotFoundremote-proxy 在 hijack 前回 502 JSON
// - 其他 wrapped errordial / 握手 / 解析錯誤
//
// 注意caller 拿到 conn 後**必須自己負責 Close**;本函式內部不會 set deadline
// 因為 streaming 場景MJPEG / SSE需要無限長的存活時間。
func (f *Forwarder) OpenStream(ctx context.Context, sessionToken string) (net.Conn, error) {
if f.proxyHost == "" {
return nil, errors.New("session: forwarder has no proxy host (check VISIONA_PROXY_INTERNAL_URL)")
}
if sessionToken == "" {
return nil, errors.New("session: forwarder.OpenStream requires non-empty sessionToken")
}
// 1. raw TCP dial
conn, err := f.dialer.DialContext(ctx, "tcp", f.proxyHost)
if err != nil {
return nil, fmt.Errorf("session: dial remote-proxy %s: %w", f.proxyHost, err)
}
// 2. 寫 POST /internal/forward/raw?token=...
// 仿 dialRawForward 測試 helper 的格式(見 internal/relay/integration_raw_test.go
reqLine := fmt.Sprintf(
"POST /internal/forward/raw?token=%s HTTP/1.1\r\n"+
"Host: %s\r\n"+
"Content-Length: 0\r\n"+
"\r\n",
url.QueryEscape(sessionToken), f.proxyHost,
)
// 設一個短的握手 deadline避免 remote-proxy 假死時 hang 住。
if err := conn.SetWriteDeadline(time.Now().Add(defaultHandshakeTimeout)); err != nil {
_ = conn.Close()
return nil, fmt.Errorf("session: set write deadline: %w", err)
}
if _, err := conn.Write([]byte(reqLine)); err != nil {
_ = conn.Close()
return nil, fmt.Errorf("session: write forward request: %w", err)
}
// 3. 讀握手 — 預期 "HTTP/1.1 200 Connected\r\n\r\n"
if err := conn.SetReadDeadline(time.Now().Add(defaultHandshakeTimeout)); err != nil {
_ = conn.Close()
return nil, fmt.Errorf("session: set read deadline: %w", err)
}
reader := bufio.NewReader(conn)
statusLine, err := reader.ReadString('\n')
if err != nil {
_ = conn.Close()
return nil, fmt.Errorf("session: read handshake status: %w", err)
}
statusLine = strings.TrimRight(statusLine, "\r\n")
// 解析 status code
// 格式HTTP/1.1 200 Connected 或 HTTP/1.1 502 Bad Gateway
if !strings.HasPrefix(statusLine, "HTTP/1.1 200") {
// 非 200 → 把 body 讀出來幫 debug常見502 = TUNNEL_DISCONNECTED
bodyHint := drainAndPeek(reader)
_ = conn.Close()
// session 不存在的明確錯誤對應 ErrSessionNotFound
if strings.Contains(statusLine, "502") {
return nil, fmt.Errorf("%w: remote-proxy responded %q (body hint: %s)",
ErrSessionNotFound, statusLine, bodyHint)
}
return nil, fmt.Errorf("session: forward handshake failed: %q (body hint: %s)",
statusLine, bodyHint)
}
// 4. 把握手後的 header 讀完(一直讀到空行)
for {
line, err := reader.ReadString('\n')
if err != nil {
_ = conn.Close()
return nil, fmt.Errorf("session: read handshake headers: %w", err)
}
if line == "\r\n" || line == "\n" {
break
}
}
// 5. 清掉 deadline因為後續 streaming 場景不該再 timeout
if err := conn.SetDeadline(time.Time{}); err != nil {
_ = conn.Close()
return nil, fmt.Errorf("session: clear deadline: %w", err)
}
// 6. 如果 reader 裡還有預讀資料bufio.NewReader 可能讀超過一行),
// 回傳一個包裝 conn 把預讀的 byte 接回 stream。
// 這個情境在 raw forward 上理論上不會發生remote-proxy 在發出
// "200 Connected\r\n\r\n" 之後不會主動寫資料 — 它要等 caller 寫
// request 才會從 yamux stream 收 response但保險起見處理。
if buffered := reader.Buffered(); buffered > 0 {
peek, _ := reader.Peek(buffered)
f.logger.Warn("forwarder: unexpected bytes after handshake; wrapping conn",
"bytes", buffered)
return newPrefixConn(conn, append([]byte(nil), peek...)), nil
}
return conn, nil
}
// ForwardHTTP 是「給定 http.Request回傳 *http.Response」的高階 helper。
//
// 內部實作:
// 1. OpenStream 拿 raw TCP已 hijack連線
// 2. req.Write(conn) 把完整 HTTP request 寫進去
// 3. http.ReadResponse 讀出 response不消耗 body
//
// 重要response.Body **包住 conn 本身**(所以 caller 必須在用完後 Close
// response.Body這允許 streaming bodyMJPEG / SSE / chunked原樣轉發。
//
// req 的 URL.Host / Scheme 會被覆寫成 "127.0.0.1" / "http",因為 local agent
// 收到的是「打到自己 localhost」的請求caller 設定的 Host header 會被保留。
func (f *Forwarder) ForwardHTTP(ctx context.Context, sessionToken string, req *http.Request) (*http.Response, error) {
if req == nil {
return nil, errors.New("session: ForwardHTTP requires non-nil req")
}
conn, err := f.OpenStream(ctx, sessionToken)
if err != nil {
return nil, err
}
// 改寫 req 為「打給 local agent」格式
// - URL.Scheme = httpURL.Host = 127.0.0.1 → req.Write 才不會報錯
// - RequestURI 必須清空client 端不能設)
// - 不覆寫 req.Hostcaller 自行決定要不要保留 browser 的 Host
//
// 注意req 本身可能已被外部使用,這裡複製 URL 避免副作用。
outReq := req.Clone(ctx)
if outReq.URL == nil {
outReq.URL = &url.URL{}
}
outReq.URL.Scheme = "http"
outReq.URL.Host = "127.0.0.1"
outReq.RequestURI = ""
if outReq.Host == "" {
outReq.Host = "127.0.0.1"
}
// 把 request 寫到 conn
if err := outReq.Write(conn); err != nil {
_ = conn.Close()
return nil, fmt.Errorf("session: write request to forwarded conn: %w", err)
}
// 讀 response — 不可以 close conn因為 response.Body 還會用到
resp, err := http.ReadResponse(bufio.NewReader(conn), outReq)
if err != nil {
_ = conn.Close()
return nil, fmt.Errorf("session: read response from forwarded conn: %w", err)
}
// 把 conn 包進 response.Body 的 close chaincaller close body 時連 conn 一起關
resp.Body = &bodyWithConn{ReadCloser: resp.Body, conn: conn}
return resp, nil
}
// WebSocketConn 是 ForwardWebSocket 成功後回傳的結果。
//
// Conn 是「已完成 WS upgrade101 已讀掉)」的 tunnel 連線,語意上等同 net.Conn
// 後續讀寫的都是 WebSocket frame bytesapi-server 不解析 frame只做透明 byte pipe
// Resp 是 local agent 回的 101 Switching Protocols response含 Sec-WebSocket-Accept
// 等 headerhandler 需把它原樣寫回 browser 端 hijacked 連線,才能完成 browser↔agent
// 的 WS 握手。
//
// caller 必須負責 Conn.Close()。
type WebSocketConn struct {
Conn net.Conn
Resp *http.Response
}
// ForwardWebSocket 把一個 WebSocket upgrade 請求經 tunnel 轉發到 local agent
// 回傳「已升級101 已讀)」的 tunnel 連線 + local agent 的 101 response。
//
// 架構說明visionA 雙 binary vs POC 單 relay
//
// POCedge-ai-platform relay/server.go:233-296 proxyWebSocket是單 binaryrelay
// 直接持有 yamux session在同一個 handler 內完成「寫 upgrade → 讀 101 → hijack
// browser → 雙向 copy」。
//
// visionA 拆成 api-server面向瀏覽器 + auth+ remote-proxyrelay。api-server 透過
// OpenStream 對 remote-proxy 開一條 raw TCPremote-proxy 端 (internal_forward_raw.go)
// 對這條連線做**透明 byte pipe、不解析 HTTP**,所以 WS upgrade bytes 原樣穿過
// remote-proxy 到 local agentclient.go:370 handleWebSocket 已實作本地端 pipe
//
// 因此 visionA 的 WS forward 分工:
// - 本函式forwarder 層OpenStream → 寫 upgrade req → 讀 101 → 回 conn+resp
// - handler 層camera.go newWebSocketProxyHandlerhijack browser → 回寫 101 →
// 雙向 io.Copy照 POC proxyWebSocket 的後半段移植)
//
// 失敗回傳的 error
// - ErrSessionNotFound無 active sessionOpenStream 已映射)
// - 非 101local agent 拒絕 upgrade回傳含該 response 的 error 供 handler 轉發原狀態碼)
// - 其他 wrapped errordial / 寫 / 讀失敗
func (f *Forwarder) ForwardWebSocket(ctx context.Context, sessionToken string, req *http.Request) (*WebSocketConn, error) {
if req == nil {
return nil, errors.New("session: ForwardWebSocket requires non-nil req")
}
conn, err := f.OpenStream(ctx, sessionToken)
if err != nil {
return nil, err
}
// 改寫 req 為「打給 local agent」格式同 ForwardHTTP
outReq := req.Clone(ctx)
if outReq.URL == nil {
outReq.URL = &url.URL{}
}
outReq.URL.Scheme = "http"
outReq.URL.Host = "127.0.0.1"
outReq.RequestURI = ""
if outReq.Host == "" {
outReq.Host = "127.0.0.1"
}
// 寫 upgrade request 到 tunnel保留 Upgrade / Connection / Sec-WebSocket-* header —
// 這些是 hop-by-hop但 WS upgrade 依賴它們,所以 req.Write 會原樣送出)。
if err := outReq.Write(conn); err != nil {
_ = conn.Close()
return nil, fmt.Errorf("session: write ws upgrade request: %w", err)
}
// 讀 upgrade response。用 bufio.Reader 讀,讀完 101 header 後若有多讀的 byte
// WS frame 可能緊跟在 101 後面)要接回 conn避免丟失。
br := bufio.NewReader(conn)
resp, err := http.ReadResponse(br, outReq)
if err != nil {
_ = conn.Close()
return nil, fmt.Errorf("session: read ws upgrade response: %w", err)
}
if resp.StatusCode != http.StatusSwitchingProtocols {
// local agent 拒絕 upgrade非 101。把 conn 包進 response.Body 讓 handler
// 能讀出 body 轉發原狀態碼給 browser再由 handler close。
resp.Body = &bodyWithConn{ReadCloser: resp.Body, conn: conn}
return nil, &wsUpgradeError{Resp: resp}
}
// 101 成功。把 bufio 預讀但還沒被消費的 byte 接回 conn 開頭,回傳給 handler。
outConn := conn
if buffered := br.Buffered(); buffered > 0 {
peek, _ := br.Peek(buffered)
outConn = newPrefixConn(conn, append([]byte(nil), peek...))
}
return &WebSocketConn{Conn: outConn, Resp: resp}, nil
}
// wsUpgradeError 代表「tunnel 通了,但 local agent 回的不是 101」。
// 帶著原 response 讓 handler 能把原狀態碼 / body 轉發回 browser。
type wsUpgradeError struct {
Resp *http.Response
}
func (e *wsUpgradeError) Error() string {
return fmt.Sprintf("session: ws upgrade rejected by local agent: %s", e.Resp.Status)
}
// AsWSUpgradeError 若 err 是 ws upgrade 被拒(非 101回傳該 response 與 true。
// handler 用它把 local agent 的原狀態碼轉發給 browser。
func AsWSUpgradeError(err error) (*http.Response, bool) {
var e *wsUpgradeError
if errors.As(err, &e) {
return e.Resp, true
}
return nil, false
}
// ----------------------------------------------------------------------
// Helpers
// ----------------------------------------------------------------------
// drainAndPeek 嘗試讀少量 byte 給 error message 加上 context
// 不阻塞太久,最多 256 byte。
//
// 呼叫前提caller 必須已經對 underlying conn 設過 ReadDeadline這個函式只
// 在 OpenStream 握手失敗的 error path 被呼叫,該路徑已經 SetReadDeadline
// 到 defaultHandshakeTimeout所以 Read 不會 hang 住;若 deadline 已過,
// Read 會立刻回 0 + deadline error行為仍然是「不阻塞」。
func drainAndPeek(reader *bufio.Reader) string {
buf := make([]byte, 256)
n, _ := reader.Read(buf)
return strings.TrimSpace(string(buf[:n]))
}
// bodyWithConn 把 ReadCloser 與底層 net.Conn 綁在一起,
// caller close body 時順便關 conn避免 leak
type bodyWithConn struct {
io.ReadCloser
conn net.Conn
}
// Close 同時關閉 body 與底層 conn以最後一個非 nil 的 error 回傳。
func (b *bodyWithConn) Close() error {
bodyErr := b.ReadCloser.Close()
connErr := b.conn.Close()
if bodyErr != nil {
return bodyErr
}
return connErr
}
// prefixConn 把預讀的 byte 接回 net.Conn 開頭,供 caller 透明使用。
//
// 並發說明net.Conn 本身對單一 goroutine 讀 + 單一 goroutine 寫是安全的。
// prefixConn 只包裝 Readprefix 的讀取不會跨 goroutine 共享Read 慣例上
// 只由 reader goroutine 呼叫),所以這裡不需要額外的 mutex。
type prefixConn struct {
net.Conn
prefix []byte
}
func newPrefixConn(c net.Conn, prefix []byte) *prefixConn {
return &prefixConn{Conn: c, prefix: prefix}
}
func (p *prefixConn) Read(b []byte) (int, error) {
if len(p.prefix) > 0 {
n := copy(b, p.prefix)
p.prefix = p.prefix[n:]
return n, nil
}
return p.Conn.Read(b)
}