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

407 lines
16 KiB
Go
Raw Permalink 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.

// proxy.go — 「把 gin 請求轉發到 local agent」的共用邏輯。
//
// 大量 device / camera / media / model load-to-device endpoint 都會走同一條路徑:
// 1. 從 UserContext 拿到當前使用者
// 2. 透過 SessionStore / ProxyClient 找到該使用者的 active session token
// 3. 用 Forwarder.ForwardHTTP 代理請求body / headers / path 原樣送)
// 4. 把 response 原樣寫回 gin.ResponseWriter支援 streaming
//
// 把這段抽成 handler 產生器,讓 devices.go / camera.go 等只需宣告路徑即可。
package api
import (
"context"
"errors"
"fmt"
"io"
"net/http"
"strings"
"sync"
"time"
"github.com/gin-gonic/gin"
"visiona-backend/internal/session"
)
// defaultProxyRequestTimeout 是「非 streaming 型」proxy 請求的整體 timeout。
//
// 對 streaming 端點MJPEG / SSE不套用此 timeout — 我們靠 gin 的 ctx 取消機制
// 在 browser 關閉時順帶關 conn。300s 對 scan / flash可能很慢夠寬鬆。
const defaultProxyRequestTimeout = 300 * time.Second
// proxyOptions 控制 proxy handler 的細部行為。
type proxyOptions struct {
// streaming 若為 true 代表 response body 可能是長連線MJPEG / SSE
// 這種情況下我們不套 timeout、並對 gin.Writer.Flush 啟用 chunk 推送。
streaming bool
// rewritePath 可選:若非空,就把請求 path 改寫成這個值再送到 local agent。
// 雛形大多不需要api-server 的路徑與 local agent 的路徑一致)。
rewritePath string
}
// newProxyHandler 產生一個 gin.HandlerFunc會把當前請求透過 Forwarder 轉發到
// local agent由 UserContext 對應的 active session 決定)。
//
// 用法:
//
// g.GET("/devices", newProxyHandler(deps, proxyOptions{}))
// g.GET("/camera/stream", newProxyHandler(deps, proxyOptions{streaming: true}))
func newProxyHandler(deps Deps, opts proxyOptions) gin.HandlerFunc {
return func(c *gin.Context) {
// 1. 檢查必要依賴
if deps.Forwarder == nil || deps.SessionStore == nil {
WriteNotImplemented(c, "forwarder/session store not configured")
return
}
// 2. 找當前使用者的 active session token
// Phase 0.7 security fix C1 (見 .autoflow/05-implementation/review/phase-0.7-security-audit.md)
// 移除 demo-user fallbackapiGroup 下所有 handler 都被 AuthMiddleware 保護,
// 拿不到 UserContext 代表 middleware 設定錯誤,回 500 比 silent fallback 安全。
uc, ok := UserContextFrom(c)
if !ok || uc.UserID == "" {
WriteError(c, http.StatusInternalServerError, ErrCodeInternalError,
"missing user context (auth middleware misconfigured?)", nil)
return
}
userID := uc.UserID
token, err := pickActiveSessionToken(c.Request.Context(), deps.SessionStore, userID, deps.Logger)
if err != nil {
writeTunnelError(c, err)
return
}
// 3. 決定 rewrite path可選
outPath := c.Request.URL.Path
if opts.rewritePath != "" {
outPath = opts.rewritePath
}
if c.Request.URL.RawQuery != "" {
outPath += "?" + c.Request.URL.RawQuery
}
// 4. 組出「打給 local agent」的 http.Request
ctx := c.Request.Context()
if !opts.streaming {
// 對非 streaming 端點加個總 timeout避免 local agent hang 住
var cancel context.CancelFunc
ctx, cancel = context.WithTimeout(ctx, defaultProxyRequestTimeout)
defer cancel()
}
outReq, err := http.NewRequestWithContext(ctx, c.Request.Method, outPath, c.Request.Body)
if err != nil {
WriteError(c, http.StatusInternalServerError, ErrCodeInternalError,
"proxy: build upstream request: "+err.Error(), nil)
return
}
// 複製 headers過濾掉 hop-by-hopForwarder 不會動,但避免重複)
copyProxyRequestHeaders(c.Request.Header, outReq.Header)
// Content-Length 要保留
if cl := c.Request.ContentLength; cl > 0 {
outReq.ContentLength = cl
}
// 5. 呼叫 Forwarder
resp, err := deps.Forwarder.ForwardHTTP(ctx, token, outReq)
if err != nil {
writeTunnelError(c, err)
return
}
defer resp.Body.Close()
// 6. 把 response 寫回 gin.Writer
writeProxyResponse(c, resp, opts.streaming)
}
}
// newWebSocketProxyHandler 產生一個 gin.HandlerFunc把 browser 的 WebSocket 連線
// 經 tunnel 轉發到 local agent推論結果 WS `inference:<deviceId>` 走這條回顯)。
//
// 認證:本 handler 掛在 apiGroupAuthMiddleware走 same-origin cookie 認證。
// 瀏覽器的 `new WebSocket(...)` 對 same-origin 會自動帶 visiona_session cookie
// AuthMiddleware 驗 cookie → 放行。**刻意不支援 token-in-URL**security 定案:
// long-lived session token 進 URL 會落 access log → 帳號接管風險Critical
//
// 授權S2 / IDOR與所有 /api/devices/:id/* proxy 路徑同一套 posture —
// pickActiveSessionToken 只挑「當前 user 自己的」active session tokenrequest 只會
// 被送到該 user 自己的 local agent。deviceId 只是透傳給 local agent 的 path 參數,
// 攻擊者帶別人的 deviceId 也只會打到自己的 tunnel打不到別人的 agent。多 user /
// 多 device 的 strict deviceId ↔ session 綁定屬 Phase 1見 pickActiveSessionToken 註解)。
//
// 流程(照 POC edge-ai-platform relay/server.go:233-296 proxyWebSocket 移植後半段):
// 1. 找 user 的 active session token
// 2. Forwarder.ForwardWebSocketOpenStream → 寫 upgrade req → 讀 101
// 3. Hijack browser 連線 → 回寫 101 → 與 tunnel conn 雙向 io.Copy
//
// WS payload 契約:透明轉發 raw bytesapi-server 不解析 / 不包 envelope
// local agent 直接推 raw driver.InferenceResult JSON前端自行解析
func newWebSocketProxyHandler(deps Deps) gin.HandlerFunc {
return func(c *gin.Context) {
if deps.Forwarder == nil || deps.SessionStore == nil {
WriteNotImplemented(c, "forwarder/session store not configured")
return
}
uc, ok := UserContextFrom(c)
if !ok || uc.UserID == "" {
WriteError(c, http.StatusInternalServerError, ErrCodeInternalError,
"missing user context (auth middleware misconfigured?)", nil)
return
}
token, err := pickActiveSessionToken(c.Request.Context(), deps.SessionStore, uc.UserID, deps.Logger)
if err != nil {
writeTunnelError(c, err)
return
}
// 建出「打給 local agent」的 upgrade request沿用原 path + query + WS header。
outPath := c.Request.URL.Path
if c.Request.URL.RawQuery != "" {
outPath += "?" + c.Request.URL.RawQuery
}
// WS 是長連線,不套 defaultProxyRequestTimeoutS1streaming 不設 timeout
outReq, err := http.NewRequestWithContext(c.Request.Context(), c.Request.Method, outPath, nil)
if err != nil {
WriteError(c, http.StatusInternalServerError, ErrCodeInternalError,
"ws proxy: build upstream request: "+err.Error(), nil)
return
}
// WS upgrade 依賴 Upgrade / Connection / Sec-WebSocket-* header必須完整帶上。
// 不走 copyProxyRequestHeaders那會剝 Connection/Upgrade 等 hop-by-hop
copyWebSocketUpgradeHeaders(c.Request.Header, outReq.Header)
wsConn, err := deps.Forwarder.ForwardWebSocket(c.Request.Context(), token, outReq)
if err != nil {
// local agent 拒絕 upgrade非 101→ 轉發原狀態碼給 browser。
if resp, ok := session.AsWSUpgradeError(err); ok {
defer resp.Body.Close()
writeProxyResponse(c, resp, false)
return
}
writeTunnelError(c, err)
return
}
defer wsConn.Conn.Close()
// Hijack browser 連線 → 回寫 101 → 雙向 pipe照 POC proxyWebSocket 後半段)。
hijacker, ok := c.Writer.(http.Hijacker)
if !ok {
WriteError(c, http.StatusInternalServerError, ErrCodeInternalError,
"ws proxy: hijacking not supported", nil)
return
}
clientConn, clientBuf, err := hijacker.Hijack()
if err != nil {
logOrDefault(deps.Logger).Warn("ws proxy: hijack failed",
"error", err, "request_id", RequestIDFrom(c))
return
}
defer clientConn.Close()
// 把 local agent 的 101 response 原樣寫回 browser完成 browser 端 WS 握手。
if err := wsConn.Resp.Write(clientBuf); err != nil {
logOrDefault(deps.Logger).Warn("ws proxy: write 101 to browser failed",
"error", err, "request_id", RequestIDFrom(c))
return
}
if err := clientBuf.Flush(); err != nil {
logOrDefault(deps.Logger).Warn("ws proxy: flush 101 to browser failed",
"error", err, "request_id", RequestIDFrom(c))
return
}
// 雙向 byte pipebrowser ↔ tunnel。任一方向結束就關掉另一邊。
var wg sync.WaitGroup
wg.Add(2)
go func() {
defer wg.Done()
_, _ = io.Copy(wsConn.Conn, clientConn)
_ = wsConn.Conn.Close()
}()
go func() {
defer wg.Done()
_, _ = io.Copy(clientConn, wsConn.Conn)
_ = clientConn.Close()
}()
wg.Wait()
}
}
// copyWebSocketUpgradeHeaders 複製 WS upgrade 所需的 header 到 upstream request。
//
// 與 copyProxyRequestHeaders 的差別:**保留** Connection / UpgradeWS 握手必需,
// 雖是 hop-by-hop並帶上 Sec-WebSocket-*。同樣剝掉 Authorization / Origin
// Origin 剝除理由同 copyProxyRequestHeaders避免 local agent CORS 對 tunnel 中繼
// 請求誤判)。
func copyWebSocketUpgradeHeaders(src, dst http.Header) {
for name, values := range src {
if strings.EqualFold(name, "Authorization") || strings.EqualFold(name, "Origin") {
continue
}
// 其餘(含 Connection / Upgrade / Sec-WebSocket-Key / -Version / -Protocol /
// -Extensions全部保留WS 握手需要。
for _, v := range values {
dst.Add(name, v)
}
}
}
// pickActiveSessionToken 找出當前使用者在雲端的 active session token。
//
// 雛形邏輯(單一 user + 單一 agent走 Store.List過濾 userID 對得上的第一筆。
// OIDC 模式下 userID 是 Member Center 簽出的 OIDC subUUIDtunnel session 在
// pairing exchange 時被綁到同個 sub因此能對上。
//
// 多 user / 多 device 階段Phase 1需要 store.ListByUser(userID) 原生介面,
// 見 session.Store TODO。
//
// Phase 0.7 security audit M2 (見 .autoflow/05-implementation/review/phase-0.7-security-audit.md)
// **保留寬鬆比對待人工介入修復**
// - 完整修法是「s.UserID != "" && s.UserID == userID」strict equality
// - 但 prototype 的 relay.NewLocalHandle (internal/relay/local_handle.go:31)
// 在 tunnel handshake 時不查 SessionTokenStore所以 Summary.UserID 永遠為空
// - 改 strict 會讓所有 e2e proxy 鏈路全斷TestE2E_FullFlow_PairingToForward 等)
// - 正解需 relay 端在 HandleTunnelConnect 時拿 token 查 SessionTokenStore
// 取得 user_id 並寫入 LocalHandle.summary.UserID屬 Phase 1 follow-up
//
// 暫保留寬鬆比對C1/M1 handler-side strict UserContext 已優先處理 — 任何 request
// 進入此函式時 userID 必非空handler 在前面已 abort 500所以唯一仍寬鬆的條件是
// s.UserID == ""relay-side 尚未 backfill
//
// logger 參數保留給未來觀測list 失敗時 log warn目前尚未使用測試傳 nil 即可。
func pickActiveSessionToken(ctx context.Context, store session.Store, userID string, _ any) (string, error) {
listCtx, cancel := context.WithTimeout(ctx, 2*time.Second)
defer cancel()
summaries, err := store.List(listCtx)
if err != nil {
return "", fmt.Errorf("proxy: list sessions: %w", err)
}
if len(summaries) == 0 {
return "", session.ErrSessionNotFound
}
for _, s := range summaries {
// 寬鬆比對handler 已確保 userID 非空C1 strict mode
// 暫接受 s.UserID == "" 直到 relay 端 backfill UserIDM2 待人工介入)。
if s.UserID == "" || s.UserID == userID {
return s.Token, nil
}
}
return "", session.ErrSessionNotFound
}
// writeTunnelError 把 forwarder / store 的錯誤映射到統一的 API 錯誤格式。
//
// - ErrSessionNotFound / ErrSessionClosed → 502 TUNNEL_DISCONNECTED
// - 其他 → 502 TUNNEL_ERROR本質上是 local agent 不可達)
func writeTunnelError(c *gin.Context, err error) {
if errors.Is(err, session.ErrSessionNotFound) || errors.Is(err, session.ErrSessionClosed) {
WriteError(c, http.StatusBadGateway, ErrCodeTunnelDisconnect,
"local agent 未連線或 tunnel 斷開", nil)
return
}
WriteError(c, http.StatusBadGateway, ErrCodeTunnelError,
"tunnel error: "+err.Error(), nil)
}
// copyProxyRequestHeaders 把 src 的 headers 複製到 dst但略過 hop-by-hop。
//
// 對齊 RFC 7230 §6.1 hop-by-hop headers
//
// Connection, Keep-Alive, Proxy-Authenticate, Proxy-Authorization,
// TE, Trailers, Transfer-Encoding, Upgrade
//
// 這些由 Forwarder / underlying conn 自動處理,不該 blind copy。
//
// **此函式只用在 tunnel proxy 轉發路徑**newProxyHandler → local agent 中繼),
// 不影響 api-server 自身面向瀏覽器的 CORSMiddleware那是另一套 gin-contrib/cors
func copyProxyRequestHeaders(src, dst http.Header) {
for name, values := range src {
if isHopByHopHeader(name) {
continue
}
// Authorization header 雛形不必送local agent 沒有對應的 auth 系統);
// 但保留其他 custom headerX-From-Api 等 test fixture 會用)
if strings.EqualFold(name, "Authorization") {
continue
}
// Origin header 必須剝掉。原因:
// 雲端瀏覽器打「開始推論」等 state-changing POST 時會自動帶上 stage 網域的
// Origin如 http://stage-9527.innovedus.com。此 Origin 經 nginx → api-server
// → tunnel 一路原樣透傳到 local agentlocal agent 的 CORSMiddleware 只把
// 127.0.0.1/localhost/::1 列白名單,非白名單 + state-changing 方法會直接 403。
// 但 tunnel 中繼請求並非瀏覽器對 local agent 的 cross-origin 請求,不該被
// local agent 的 CORS 邏輯攔。剝掉 Origin 後 local agent 走「無 Origin →
// same-origin 放行」的快速路徑(見 local-agent middleware.go CORSMiddleware
//
// 安全性:這只發生在 api-server → local agent 這段 server-to-server 中繼,
// 瀏覽器 → api-server 那段仍由 api-server 自己的 CORSMiddleware 把關,不受影響。
if strings.EqualFold(name, "Origin") {
continue
}
for _, v := range values {
dst.Add(name, v)
}
}
}
// isHopByHopHeader 回報 header 名稱是否為 hop-by-hop。
func isHopByHopHeader(name string) bool {
switch strings.ToLower(name) {
case "connection", "keep-alive", "proxy-authenticate", "proxy-authorization",
"te", "trailers", "transfer-encoding", "upgrade":
return true
}
return false
}
// writeProxyResponse 把 upstream response 原樣寫回 gin。
//
// 支援 streaming若 streaming=true 且 response 有 Flusher每次 Read 後立即 Flush。
// 這讓 MJPEG / SSE 的 frame 能即時抵達 browser。
func writeProxyResponse(c *gin.Context, resp *http.Response, streaming bool) {
// 複製 headers略過 hop-by-hop
for name, values := range resp.Header {
if isHopByHopHeader(name) {
continue
}
for _, v := range values {
c.Writer.Header().Add(name, v)
}
}
c.Writer.WriteHeader(resp.StatusCode)
if !streaming {
// 非 streaming一口氣 copy 完
_, _ = io.Copy(c.Writer, resp.Body)
return
}
// Streaming邊讀邊 flush。buffer 大小 8KB平衡延遲與 syscall 次數。
buf := make([]byte, 8*1024)
flusher, _ := c.Writer.(http.Flusher)
for {
n, rerr := resp.Body.Read(buf)
if n > 0 {
if _, werr := c.Writer.Write(buf[:n]); werr != nil {
// browser 斷線 → 停止conn 會在 resp.Body.Close 時關掉 upstream
return
}
if flusher != nil {
flusher.Flush()
}
}
if rerr != nil {
// io.EOF 或連線結束都是正常
return
}
}
}