virtual-kubelet Provider的exec API返回WebSocket,如何实现RunInContainer?
实现Virtual-Kubelet Provider的RunInContainer方法(对接WebSocket版Exec API)
原代码依赖SPDY协议与K8S Exec API交互,若目标Provider的Exec接口返回WebSocket,需调整为WebSocket协议通信,核心实现逻辑如下:
核心思路
- 转换请求协议:将HTTP/SPDY的URL替换为WebSocket协议(
ws://或wss://),并将Exec所需参数(命令、终端配置等)转为URL查询参数。 - 建立WebSocket连接:基于K8S客户端配置的认证信息(Token、CA证书等),配置WebSocket拨号选项并建立连接。
- 多路数据流转发:K8S Exec WebSocket通过首字节通道标识区分不同数据流(0=stdin、1=stdout、2=stderr、3=终端resize),需将
api.AttachIO的输入输出与WebSocket消息对应转发,同时处理终端尺寸变更事件。
完整实现代码
import ( "context" "encoding/json" "errors" "fmt" "io" "net/http" "net/url" "github.com/gorilla/websocket" v1 "k8s.io/api/core/v1" "k8s.io/client-go/rest" "k8s.io/klog/v2" "github.com/virtual-kubelet/virtual-kubelet/node/api" ) func (p *MopProvider) RunInContainer(ctx context.Context, namespace, name, container string, cmd []string, attach api.AttachIO) error { klog.Infof("receive ExecInContainer request for container %q", container) // 1. 构建WebSocket请求参数与URL execOpts := v1.PodExecOptions{ Container: container, Command: cmd, Stdin: true, Stdout: true, Stderr: true, TTY: attach.TTY(), } // 将Exec选项转为URL查询参数 query := url.Values{} if execOpts.Stdin { query.Set("stdin", "true") } if execOpts.Stdout { query.Set("stdout", "true") } if execOpts.Stderr { query.Set("stderr", "true") } if execOpts.TTY { query.Set("tty", "true") } for _, c := range execOpts.Command { query.Add("command", c) } query.Set("container", execOpts.Container) // 替换协议为WebSocket,拼接完整请求URL baseURL, err := url.Parse(p.K8SClient.CoreV1().RESTClient().BaseURL().String()) if err != nil { return fmt.Errorf("parse base URL failed: %w", err) } baseURL.Scheme = "wss" baseURL.Path = fmt.Sprintf("/api/v1/namespaces/%s/pods/%s/exec", namespace, name) baseURL.RawQuery = query.Encode() // 2. 配置WebSocket拨号认证信息 config := p.K8SClient.GetConfig() dialer := websocket.DefaultDialer // 加载CA证书用于服务端校验 if config.CAData != nil { certPool, err := rest.CertPoolFromBytes(config.CAData) if err != nil { return fmt.Errorf("create cert pool failed: %w", err) } dialer.TLSClientConfig.RootCAs = certPool } else if config.CAFile != "" { tlsConfig, err := rest.LoadTLSConfig(config) if err != nil { return fmt.Errorf("load TLS config failed: %w", err) } dialer.TLSClientConfig.RootCAs = tlsConfig.RootCAs } // 设置Bearer Token认证头 header := http.Header{} if config.BearerToken != "" { header.Set("Authorization", fmt.Sprintf("Bearer %s", config.BearerToken)) } // 建立WebSocket连接 conn, _, err := dialer.DialContext(ctx, baseURL.String(), header) if err != nil { return fmt.Errorf("dial websocket failed: %w", err) } defer conn.Close() // 3. 启动数据流转发协程 errChan := make(chan error, 3) // 转发stdin到WebSocket(通道0) go func() { if !execOpts.Stdin { errChan <- nil return } buf := make([]byte, 1024) for { n, err := attach.Stdin().Read(buf) if err != nil { if errors.Is(err, io.EOF) { errChan <- nil return } errChan <- fmt.Errorf("read stdin failed: %w", err) return } // 首字节为通道标识0,后续为输入数据 msg := append([]byte{0}, buf[:n]...) if err := conn.WriteMessage(websocket.BinaryMessage, msg); err != nil { errChan <- fmt.Errorf("write stdin to websocket failed: %w", err) return } } }() // 转发WebSocket的stdout/stderr到attach(通道1、2) go func() { for { _, msg, err := conn.ReadMessage() if err != nil { if websocket.IsCloseError(err, websocket.CloseNormalClosure, websocket.CloseGoingAway) { errChan <- nil return } errChan <- fmt.Errorf("read websocket message failed: %w", err) return } if len(msg) == 0 { continue } channel := msg[0] data := msg[1:] switch channel { case 1: // 标准输出 if _, err := attach.Stdout().Write(data); err != nil { errChan <- fmt.Errorf("write stdout failed: %w", err) return } case 2: // 标准错误 if _, err := attach.Stderr().Write(data); err != nil { errChan <- fmt.Errorf("write stderr failed: %w", err) return } } } }() // 处理终端resize事件(发送到通道3) go func() { if !execOpts.TTY { errChan <- nil return } resizeCh := attach.Resize() for size := range resizeCh { resizeMsg, err := json.Marshal(map[string]uint16{ "width": uint16(size.Width), "height": uint16(size.Height), }) if err != nil { errChan <- fmt.Errorf("marshal resize message failed: %w", err) return } msg := append([]byte{3}, resizeMsg...) if err := conn.WriteMessage(websocket.BinaryMessage, msg); err != nil { errChan <- fmt.Errorf("write resize to websocket failed: %w", err) return } } errChan <- nil }() // 等待协程完成或出错 for i := 0; i < 3; i++ { if err := <-errChan; err != nil { return err } } return nil }
关键细节说明
- 通道标识规则:K8S Exec WebSocket约定首字节为数据流标识,0对应标准输入、1对应标准输出、2对应标准错误、3对应终端尺寸变更。
- 认证处理:必须携带K8S的Token与CA证书信息,否则无法通过API鉴权。
- 资源清理:通过
defer conn.Close()确保WebSocket连接在函数退出时关闭,避免资源泄漏。
内容的提问来源于stack exchange,提问作者贾永鹏
相关产品推荐
相关产品推荐

