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

virtual-kubelet Provider的exec API返回WebSocket,如何实现RunInContainer?

实现Virtual-Kubelet Provider的RunInContainer方法(对接WebSocket版Exec API)

原代码依赖SPDY协议与K8S Exec API交互,若目标Provider的Exec接口返回WebSocket,需调整为WebSocket协议通信,核心实现逻辑如下:

核心思路

  1. 转换请求协议:将HTTP/SPDY的URL替换为WebSocket协议(ws://或wss://),并将Exec所需参数(命令、终端配置等)转为URL查询参数。
  2. 建立WebSocket连接:基于K8S客户端配置的认证信息(Token、CA证书等),配置WebSocket拨号选项并建立连接。
  3. 多路数据流转发: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,提问作者贾永鹏

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.11 20:50:35