推論工作區後端塊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>
407 lines
16 KiB
Go
407 lines
16 KiB
Go
// 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 fallback:apiGroup 下所有 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-hop(Forwarder 不會動,但避免重複)
|
||
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 掛在 apiGroup(AuthMiddleware)下,走 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 token,request 只會
|
||
// 被送到該 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.ForwardWebSocket:OpenStream → 寫 upgrade req → 讀 101
|
||
// 3. Hijack browser 連線 → 回寫 101 → 與 tunnel conn 雙向 io.Copy
|
||
//
|
||
// WS payload 契約:透明轉發 raw bytes,api-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 是長連線,不套 defaultProxyRequestTimeout(S1:streaming 不設 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 pipe:browser ↔ 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 / Upgrade(WS 握手必需,
|
||
// 雖是 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 sub(UUID),tunnel 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 UserID(M2 待人工介入)。
|
||
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 header(X-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 agent,local 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
|
||
}
|
||
}
|
||
}
|