Files
NixMsg/test/chaos/helper.go

193 lines
4.9 KiB
Go

// Package chaos 提供基于 toxiproxy 的弱网故障注入辅助。
// 上游地址由调用方传入,不写死端口。
package chaos
import (
"bytes"
"encoding/json"
"fmt"
"io"
"net/http"
"strings"
"time"
)
// Client 调用 toxiproxy HTTP API。
type Client struct {
BaseURL string
HTTPClient *http.Client
}
// NewClient 创建客户端。baseURL 形如 http://127.0.0.1:8474。
func NewClient(baseURL string) *Client {
return &Client{
BaseURL: strings.TrimRight(baseURL, "/"),
HTTPClient: &http.Client{
Timeout: 10 * time.Second,
},
}
}
// Proxy 描述一个 toxiproxy 代理。
type Proxy struct {
Name string `json:"name"`
Listen string `json:"listen"`
Upstream string `json:"upstream"`
Enabled bool `json:"enabled"`
}
// Toxic 描述一条故障规则。
type Toxic struct {
Name string `json:"name"`
Type string `json:"type"`
Stream string `json:"stream,omitempty"`
Toxicity float32 `json:"toxicity,omitempty"`
Attributes map[string]any `json:"attributes,omitempty"`
}
// CreateProxy 创建或覆盖同名代理。listen 可用 host:0 让 toxiproxy 分配端口。
func (c *Client) CreateProxy(name, listen, upstream string) (*Proxy, error) {
if name == "" {
return nil, fmt.Errorf("proxy name required")
}
if listen == "" {
return nil, fmt.Errorf("listen required")
}
if upstream == "" {
return nil, fmt.Errorf("upstream required")
}
body := Proxy{
Name: name,
Listen: listen,
Upstream: upstream,
Enabled: true,
}
var out Proxy
if err := c.doJSON(http.MethodPost, "/proxies", body, &out); err != nil {
return nil, err
}
return &out, nil
}
// DeleteProxy 删除代理;不存在时忽略。
func (c *Client) DeleteProxy(name string) error {
req, err := http.NewRequest(http.MethodDelete, c.BaseURL+"/proxies/"+name, nil)
if err != nil {
return err
}
resp, err := c.HTTPClient.Do(req)
if err != nil {
return err
}
defer func() { _ = resp.Body.Close() }()
if resp.StatusCode == http.StatusNotFound || resp.StatusCode == http.StatusNoContent || resp.StatusCode == http.StatusOK {
return nil
}
b, _ := io.ReadAll(resp.Body)
return fmt.Errorf("delete proxy: %s: %s", resp.Status, strings.TrimSpace(string(b)))
}
// GetProxy 读取代理(含实际 listen 地址)。
func (c *Client) GetProxy(name string) (*Proxy, error) {
var out Proxy
if err := c.doJSON(http.MethodGet, "/proxies/"+name, nil, &out); err != nil {
return nil, err
}
return &out, nil
}
// AddLatency 注入下行/上行延迟(毫秒)。stream 为空时默认 downstream。
func (c *Client) AddLatency(proxyName, toxicName string, latencyMs, jitterMs int, stream string) (*Toxic, error) {
if stream == "" {
stream = "downstream"
}
t := Toxic{
Name: toxicName,
Type: "latency",
Stream: stream,
Toxicity: 1,
Attributes: map[string]any{
"latency": latencyMs,
"jitter": jitterMs,
},
}
return c.addToxic(proxyName, t)
}
// AddResetPeer 在连接上注入 TCP RST(断开)。timeoutMs 为触发前等待。
func (c *Client) AddResetPeer(proxyName, toxicName string, timeoutMs int, stream string) (*Toxic, error) {
if stream == "" {
stream = "downstream"
}
t := Toxic{
Name: toxicName,
Type: "reset_peer",
Stream: stream,
Toxicity: 1,
Attributes: map[string]any{
"timeout": timeoutMs,
},
}
return c.addToxic(proxyName, t)
}
// RemoveToxic 删除一条 toxic。
func (c *Client) RemoveToxic(proxyName, toxicName string) error {
req, err := http.NewRequest(http.MethodDelete, c.BaseURL+"/proxies/"+proxyName+"/toxics/"+toxicName, nil)
if err != nil {
return err
}
resp, err := c.HTTPClient.Do(req)
if err != nil {
return err
}
defer func() { _ = resp.Body.Close() }()
if resp.StatusCode == http.StatusNotFound || resp.StatusCode == http.StatusNoContent || resp.StatusCode == http.StatusOK {
return nil
}
b, _ := io.ReadAll(resp.Body)
return fmt.Errorf("remove toxic: %s: %s", resp.Status, strings.TrimSpace(string(b)))
}
func (c *Client) addToxic(proxyName string, t Toxic) (*Toxic, error) {
var out Toxic
if err := c.doJSON(http.MethodPost, "/proxies/"+proxyName+"/toxics", t, &out); err != nil {
return nil, err
}
return &out, nil
}
func (c *Client) doJSON(method, path string, in any, out any) error {
var body io.Reader
if in != nil {
b, err := json.Marshal(in)
if err != nil {
return err
}
body = bytes.NewReader(b)
}
req, err := http.NewRequest(method, c.BaseURL+path, body)
if err != nil {
return err
}
if in != nil {
req.Header.Set("Content-Type", "application/json")
}
resp, err := c.HTTPClient.Do(req)
if err != nil {
return err
}
defer func() { _ = resp.Body.Close() }()
respBody, err := io.ReadAll(resp.Body)
if err != nil {
return err
}
if resp.StatusCode < 200 || resp.StatusCode >= 300 {
return fmt.Errorf("%s %s: %s: %s", method, path, resp.Status, strings.TrimSpace(string(respBody)))
}
if out == nil || len(respBody) == 0 {
return nil
}
return json.Unmarshal(respBody, out)
}