mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-10-06 18:06:36 +08:00
fix: preserve shared limiters with per-IP rules
This commit is contained in:
@@ -292,27 +292,13 @@ func (h *Handler) syncForwardServicesWithWarnings(forward *forwardRecord, method
|
|||||||
if user != nil && user.MaxConn > 0 {
|
if user != nil && user.MaxConn > 0 {
|
||||||
userMaxConn = user.MaxConn
|
userMaxConn = user.MaxConn
|
||||||
}
|
}
|
||||||
connLimiterConfig := buildConnLimiterConfig(forward, userMaxConn)
|
connLimiterConfigs := buildConnLimiterConfigs(forward, userMaxConn)
|
||||||
|
|
||||||
for _, fp := range ports {
|
for _, fp := range ports {
|
||||||
runtimeLimiters := forwardRuntimeLimiters{ConnLimiter: connLimiterConfig.Name}
|
runtimeLimiters := forwardRuntimeLimiters{ConnLimiter: joinLimiterNames(connLimiterConfigs)}
|
||||||
if ipSpeed != nil {
|
trafficLimiterNames := make([]string, 0, 2)
|
||||||
runtimeLimiters.TrafficLimiter = fmt.Sprintf("rule_traffic_limit_%d", forward.ID)
|
if limiterID != nil && speed != nil {
|
||||||
if err := h.ensureTrafficLimiterOnNode(fp.NodeID, runtimeLimiters.TrafficLimiter, speed, ipSpeed); err != nil {
|
totalLimiterName := strconv.FormatInt(*limiterID, 10)
|
||||||
// If the limiter push fails because the node is offline, skip it with a warning
|
|
||||||
if isNodeOfflineOrTimeoutError(err) {
|
|
||||||
node, _ := h.getNodeRecord(fp.NodeID)
|
|
||||||
nodeName := fmt.Sprintf("%d", fp.NodeID)
|
|
||||||
if node != nil && strings.TrimSpace(node.Name) != "" {
|
|
||||||
nodeName = strings.TrimSpace(node.Name)
|
|
||||||
}
|
|
||||||
warnings = append(warnings, fmt.Sprintf("节点 %s 不在线,已跳过下发", nodeName))
|
|
||||||
continue
|
|
||||||
}
|
|
||||||
return nil, err
|
|
||||||
}
|
|
||||||
} else if limiterID != nil && speed != nil {
|
|
||||||
runtimeLimiters.TrafficLimiter = strconv.FormatInt(*limiterID, 10)
|
|
||||||
if err := h.ensureLimiterOnNode(fp.NodeID, *limiterID, *speed); err != nil {
|
if err := h.ensureLimiterOnNode(fp.NodeID, *limiterID, *speed); err != nil {
|
||||||
// If the limiter push fails because the node is offline, skip it with a warning
|
// If the limiter push fails because the node is offline, skip it with a warning
|
||||||
if isNodeOfflineOrTimeoutError(err) {
|
if isNodeOfflineOrTimeoutError(err) {
|
||||||
@@ -326,9 +312,28 @@ func (h *Handler) syncForwardServicesWithWarnings(forward *forwardRecord, method
|
|||||||
}
|
}
|
||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
|
trafficLimiterNames = append(trafficLimiterNames, totalLimiterName)
|
||||||
}
|
}
|
||||||
|
if ipSpeed != nil {
|
||||||
|
ruleLimiterName := fmt.Sprintf("rule_traffic_limit_%d", forward.ID)
|
||||||
|
if err := h.ensureTrafficLimiterOnNode(fp.NodeID, ruleLimiterName, nil, ipSpeed); err != nil {
|
||||||
|
// If the limiter push fails because the node is offline, skip it with a warning
|
||||||
|
if isNodeOfflineOrTimeoutError(err) {
|
||||||
|
node, _ := h.getNodeRecord(fp.NodeID)
|
||||||
|
nodeName := fmt.Sprintf("%d", fp.NodeID)
|
||||||
|
if node != nil && strings.TrimSpace(node.Name) != "" {
|
||||||
|
nodeName = strings.TrimSpace(node.Name)
|
||||||
|
}
|
||||||
|
warnings = append(warnings, fmt.Sprintf("节点 %s 不在线,已跳过下发", nodeName))
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
return nil, err
|
||||||
|
}
|
||||||
|
trafficLimiterNames = append(trafficLimiterNames, ruleLimiterName)
|
||||||
|
}
|
||||||
|
runtimeLimiters.TrafficLimiter = strings.Join(trafficLimiterNames, ",")
|
||||||
|
|
||||||
if connLimiterConfig.Name != "" {
|
for _, connLimiterConfig := range connLimiterConfigs {
|
||||||
if err := h.ensureConnLimiterOnNode(fp.NodeID, connLimiterConfig); err != nil {
|
if err := h.ensureConnLimiterOnNode(fp.NodeID, connLimiterConfig); err != nil {
|
||||||
warnings = append(warnings, fmt.Sprintf("节点 %d 连接限制器下发失败: %v", fp.NodeID, err))
|
warnings = append(warnings, fmt.Sprintf("节点 %d 连接限制器下发失败: %v", fp.NodeID, err))
|
||||||
}
|
}
|
||||||
@@ -1876,27 +1881,35 @@ func (h *Handler) ensureConnLimiterOnNode(nodeID int64, cfg forwardLimiterConfig
|
|||||||
return nil
|
return nil
|
||||||
}
|
}
|
||||||
|
|
||||||
func buildConnLimiterConfig(forward *forwardRecord, userMaxConn int) forwardLimiterConfig {
|
func buildConnLimiterConfigs(forward *forwardRecord, userMaxConn int) []forwardLimiterConfig {
|
||||||
if forward == nil {
|
if forward == nil {
|
||||||
return forwardLimiterConfig{}
|
return nil
|
||||||
}
|
}
|
||||||
limits := make([]string, 0, 2)
|
|
||||||
if forward.MaxConn > 0 {
|
if forward.MaxConn > 0 {
|
||||||
limits = append(limits, fmt.Sprintf("$ %d", forward.MaxConn))
|
limits := []string{fmt.Sprintf("$ %d", forward.MaxConn)}
|
||||||
} else if userMaxConn > 0 {
|
if forward.IPMaxConn > 0 {
|
||||||
limits = append(limits, fmt.Sprintf("$ %d", userMaxConn))
|
limits = append(limits, fmt.Sprintf("$$ %d", forward.IPMaxConn))
|
||||||
|
}
|
||||||
|
return []forwardLimiterConfig{{Name: fmt.Sprintf("rule_conn_limit_%d", forward.ID), Limits: limits}}
|
||||||
|
}
|
||||||
|
configs := make([]forwardLimiterConfig, 0, 2)
|
||||||
|
if userMaxConn > 0 {
|
||||||
|
configs = append(configs, forwardLimiterConfig{Name: fmt.Sprintf("user_conn_limit_%d", forward.UserID), Limits: []string{fmt.Sprintf("$ %d", userMaxConn)}})
|
||||||
}
|
}
|
||||||
if forward.IPMaxConn > 0 {
|
if forward.IPMaxConn > 0 {
|
||||||
limits = append(limits, fmt.Sprintf("$$ %d", forward.IPMaxConn))
|
configs = append(configs, forwardLimiterConfig{Name: fmt.Sprintf("rule_conn_limit_%d", forward.ID), Limits: []string{fmt.Sprintf("$$ %d", forward.IPMaxConn)}})
|
||||||
}
|
}
|
||||||
if len(limits) == 0 {
|
return configs
|
||||||
return forwardLimiterConfig{}
|
}
|
||||||
|
|
||||||
|
func joinLimiterNames(configs []forwardLimiterConfig) string {
|
||||||
|
names := make([]string, 0, len(configs))
|
||||||
|
for _, cfg := range configs {
|
||||||
|
if cfg.Name != "" {
|
||||||
|
names = append(names, cfg.Name)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
name := fmt.Sprintf("user_conn_limit_%d", forward.UserID)
|
return strings.Join(names, ",")
|
||||||
if forward.MaxConn > 0 || forward.IPMaxConn > 0 {
|
|
||||||
name = fmt.Sprintf("rule_conn_limit_%d", forward.ID)
|
|
||||||
}
|
|
||||||
return forwardLimiterConfig{Name: name, Limits: limits}
|
|
||||||
}
|
}
|
||||||
|
|
||||||
func speedToLimitLine(key string, speed int) string {
|
func speedToLimitLine(key string, speed int) string {
|
||||||
|
|||||||
@@ -479,24 +479,30 @@ func TestBuildForwardServiceConfigs_IPv6BindIP(t *testing.T) {
|
|||||||
}
|
}
|
||||||
|
|
||||||
func TestBuildConnLimiterConfigCombinesTotalAndPerIP(t *testing.T) {
|
func TestBuildConnLimiterConfigCombinesTotalAndPerIP(t *testing.T) {
|
||||||
cfg := buildConnLimiterConfig(&forwardRecord{ID: 42, UserID: 9, MaxConn: 100, IPMaxConn: 5}, 37)
|
cfgs := buildConnLimiterConfigs(&forwardRecord{ID: 42, UserID: 9, MaxConn: 100, IPMaxConn: 5}, 37)
|
||||||
want := forwardLimiterConfig{Name: "rule_conn_limit_42", Limits: []string{"$ 100", "$$ 5"}}
|
want := []forwardLimiterConfig{{Name: "rule_conn_limit_42", Limits: []string{"$ 100", "$$ 5"}}}
|
||||||
if !reflect.DeepEqual(cfg, want) {
|
if !reflect.DeepEqual(cfgs, want) {
|
||||||
t.Fatalf("expected %+v, got %+v", want, cfg)
|
t.Fatalf("expected %+v, got %+v", want, cfgs)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestBuildConnLimiterConfigUsesUserTotalWithRulePerIP(t *testing.T) {
|
func TestBuildConnLimiterConfigUsesUserTotalWithRulePerIP(t *testing.T) {
|
||||||
cfg := buildConnLimiterConfig(&forwardRecord{ID: 42, UserID: 9, IPMaxConn: 5}, 37)
|
cfgs := buildConnLimiterConfigs(&forwardRecord{ID: 42, UserID: 9, IPMaxConn: 5}, 37)
|
||||||
want := forwardLimiterConfig{Name: "rule_conn_limit_42", Limits: []string{"$ 37", "$$ 5"}}
|
want := []forwardLimiterConfig{
|
||||||
if !reflect.DeepEqual(cfg, want) {
|
{Name: "user_conn_limit_9", Limits: []string{"$ 37"}},
|
||||||
t.Fatalf("expected %+v, got %+v", want, cfg)
|
{Name: "rule_conn_limit_42", Limits: []string{"$$ 5"}},
|
||||||
|
}
|
||||||
|
if !reflect.DeepEqual(cfgs, want) {
|
||||||
|
t.Fatalf("expected %+v, got %+v", want, cfgs)
|
||||||
|
}
|
||||||
|
if got := joinLimiterNames(cfgs); got != "user_conn_limit_9,rule_conn_limit_42" {
|
||||||
|
t.Fatalf("expected composite limiter names, got %q", got)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
func TestBuildTrafficLimiterPayloadCombinesTotalAndPerIP(t *testing.T) {
|
func TestBuildTrafficLimiterPayloadUsesOnlyPerIPRulesWhenTotalIsSeparate(t *testing.T) {
|
||||||
payload := buildTrafficLimiterPayload("rule_traffic_limit_42", intPtr(80), intPtr(40))
|
payload := buildTrafficLimiterPayload("rule_traffic_limit_42", nil, intPtr(40))
|
||||||
wantLimits := []string{"$ 10.0MB 10.0MB", "0.0.0.0/0 5.0MB 5.0MB", "::/0 5.0MB 5.0MB"}
|
wantLimits := []string{"0.0.0.0/0 5.0MB 5.0MB", "::/0 5.0MB 5.0MB"}
|
||||||
if payload["name"] != "rule_traffic_limit_42" {
|
if payload["name"] != "rule_traffic_limit_42" {
|
||||||
t.Fatalf("expected name rule_traffic_limit_42, got %v", payload["name"])
|
t.Fatalf("expected name rule_traffic_limit_42, got %v", payload["name"])
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ import (
|
|||||||
"net/http"
|
"net/http"
|
||||||
"net/http/httptest"
|
"net/http/httptest"
|
||||||
"net/url"
|
"net/url"
|
||||||
|
"reflect"
|
||||||
"strings"
|
"strings"
|
||||||
"sync"
|
"sync"
|
||||||
"testing"
|
"testing"
|
||||||
@@ -366,6 +367,143 @@ func TestUserMaxConnUpdateResyncsExistingForwards(t *testing.T) {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func TestUserMaxConnWithPerIPRuleSplitsRuntimeLimiters(t *testing.T) {
|
||||||
|
secret := "contract-jwt-secret"
|
||||||
|
router, r := setupContractRouter(t, secret)
|
||||||
|
server := httptest.NewServer(router)
|
||||||
|
defer server.Close()
|
||||||
|
|
||||||
|
userToken, err := auth.GenerateToken(3, "per_ip_user", 1, secret)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("generate user token: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
now := time.Now().UnixMilli()
|
||||||
|
if err := r.DB().Exec(`
|
||||||
|
INSERT INTO user(id, user, pwd, role_id, exp_time, flow, in_flow, out_flow, flow_reset_time, num, max_conn, created_time, updated_time, status)
|
||||||
|
VALUES(3, 'per_ip_user', 'pwd', 1, ?, 99999, 0, 0, 1, 10, 37, ?, ?, 1)
|
||||||
|
`, now+365*24*3600*1000, now, now).Error; err != nil {
|
||||||
|
t.Fatalf("insert user: %v", err)
|
||||||
|
}
|
||||||
|
if err := r.DB().Exec(`
|
||||||
|
INSERT INTO tunnel(id, name, traffic_ratio, type, protocol, flow, created_time, updated_time, status, in_ip, inx)
|
||||||
|
VALUES(11, 'user-per-ip-conn-tunnel', 1.0, 1, 'tls', 99999, ?, ?, 1, NULL, 0)
|
||||||
|
`, now, now).Error; err != nil {
|
||||||
|
t.Fatalf("insert tunnel: %v", err)
|
||||||
|
}
|
||||||
|
if err := r.DB().Exec(`
|
||||||
|
INSERT INTO node(id, name, secret, server_ip, server_ip_v4, server_ip_v6, port, interface_name, version, http, tls, socks, created_time, updated_time, status, tcp_listen_addr, udp_listen_addr, inx)
|
||||||
|
VALUES(21, 'user-per-ip-conn-node', 'user-per-ip-conn-secret', '10.23.0.1', '10.23.0.1', '', '32300-32310', '', 'v1', 1, 1, 1, ?, ?, 1, '[::]', '[::]', 0)
|
||||||
|
`, now, now).Error; err != nil {
|
||||||
|
t.Fatalf("insert node: %v", err)
|
||||||
|
}
|
||||||
|
if err := r.DB().Exec(`
|
||||||
|
INSERT INTO chain_tunnel(tunnel_id, chain_type, node_id, port, strategy, inx, protocol)
|
||||||
|
VALUES(11, 1, 21, 32301, 'round', 1, 'tls')
|
||||||
|
`).Error; err != nil {
|
||||||
|
t.Fatalf("insert chain_tunnel: %v", err)
|
||||||
|
}
|
||||||
|
if err := r.DB().Exec(`
|
||||||
|
INSERT INTO user_tunnel(id, user_id, tunnel_id, num, flow, in_flow, out_flow, flow_reset_time, exp_time, status)
|
||||||
|
VALUES(31, 3, 11, 10, 99999, 0, 0, 1, ?, 1)
|
||||||
|
`, now+365*24*3600*1000).Error; err != nil {
|
||||||
|
t.Fatalf("insert user_tunnel: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
var commandMu sync.Mutex
|
||||||
|
receivedCommands := make([]string, 0)
|
||||||
|
addCLimitersData := make([]json.RawMessage, 0)
|
||||||
|
var updateServiceData json.RawMessage
|
||||||
|
|
||||||
|
stopNode := startMockSessionForMaxConn(t, server.URL, "user-per-ip-conn-secret", func(cmdType string, data json.RawMessage) (bool, string) {
|
||||||
|
commandMu.Lock()
|
||||||
|
defer commandMu.Unlock()
|
||||||
|
receivedCommands = append(receivedCommands, cmdType)
|
||||||
|
if cmdType == "AddCLimiters" {
|
||||||
|
addCLimitersData = append(addCLimitersData, append([]byte(nil), data...))
|
||||||
|
}
|
||||||
|
if cmdType == "UpdateService" {
|
||||||
|
updateServiceData = append([]byte(nil), data...)
|
||||||
|
}
|
||||||
|
return false, ""
|
||||||
|
})
|
||||||
|
defer stopNode()
|
||||||
|
|
||||||
|
waitNodeStatus(t, r, 21, 1)
|
||||||
|
|
||||||
|
payload := map[string]interface{}{
|
||||||
|
"name": "user-per-ip-conn-forward",
|
||||||
|
"tunnelId": int64(11),
|
||||||
|
"remoteAddr": "1.1.1.1:443",
|
||||||
|
"strategy": "fifo",
|
||||||
|
"ipMaxConn": 7,
|
||||||
|
}
|
||||||
|
body, err := json.Marshal(payload)
|
||||||
|
if err != nil {
|
||||||
|
t.Fatalf("marshal payload: %v", err)
|
||||||
|
}
|
||||||
|
req := httptest.NewRequest(http.MethodPost, "/api/v1/forward/create", bytes.NewReader(body))
|
||||||
|
req.Header.Set("Authorization", userToken)
|
||||||
|
req.Header.Set("Content-Type", "application/json")
|
||||||
|
res := httptest.NewRecorder()
|
||||||
|
router.ServeHTTP(res, req)
|
||||||
|
|
||||||
|
var out response.R
|
||||||
|
if err := json.NewDecoder(res.Body).Decode(&out); err != nil {
|
||||||
|
t.Fatalf("decode response: %v", err)
|
||||||
|
}
|
||||||
|
if out.Code != 0 {
|
||||||
|
t.Fatalf("expected create success, got code=%d msg=%s", out.Code, out.Msg)
|
||||||
|
}
|
||||||
|
|
||||||
|
var forwardID int64
|
||||||
|
if err := r.DB().Raw("SELECT id FROM forward WHERE name = ?", "user-per-ip-conn-forward").Scan(&forwardID).Error; err != nil {
|
||||||
|
t.Fatalf("get forward ID: %v", err)
|
||||||
|
}
|
||||||
|
expectedRuleName := fmt.Sprintf("rule_conn_limit_%d", forwardID)
|
||||||
|
|
||||||
|
commandMu.Lock()
|
||||||
|
defer commandMu.Unlock()
|
||||||
|
if len(addCLimitersData) != 2 {
|
||||||
|
t.Fatalf("expected two AddCLimiters commands. Received: %v", receivedCommands)
|
||||||
|
}
|
||||||
|
if updateServiceData == nil {
|
||||||
|
t.Fatalf("expected UpdateService. Received: %v", receivedCommands)
|
||||||
|
}
|
||||||
|
|
||||||
|
gotLimits := make(map[string][]string)
|
||||||
|
for _, raw := range addCLimitersData {
|
||||||
|
var data map[string]interface{}
|
||||||
|
if err := json.Unmarshal(raw, &data); err != nil {
|
||||||
|
t.Fatalf("unmarshal AddCLimiters data: %v", err)
|
||||||
|
}
|
||||||
|
limits, ok := data["limits"].([]interface{})
|
||||||
|
if !ok {
|
||||||
|
t.Fatalf("expected limits array, got %T", data["limits"])
|
||||||
|
}
|
||||||
|
for _, limit := range limits {
|
||||||
|
gotLimits[fmt.Sprint(data["name"])] = append(gotLimits[fmt.Sprint(data["name"])], fmt.Sprint(limit))
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if !reflect.DeepEqual(gotLimits["user_conn_limit_3"], []string{"$ 37"}) {
|
||||||
|
t.Fatalf("expected user max limiter payload, got %v", gotLimits["user_conn_limit_3"])
|
||||||
|
}
|
||||||
|
if !reflect.DeepEqual(gotLimits[expectedRuleName], []string{"$$ 7"}) {
|
||||||
|
t.Fatalf("expected rule per-IP limiter payload, got %v", gotLimits[expectedRuleName])
|
||||||
|
}
|
||||||
|
|
||||||
|
var services []map[string]interface{}
|
||||||
|
if err := json.Unmarshal(updateServiceData, &services); err != nil {
|
||||||
|
t.Fatalf("unmarshal UpdateService data: %v", err)
|
||||||
|
}
|
||||||
|
expectedCLimiter := "user_conn_limit_3," + expectedRuleName
|
||||||
|
for _, service := range services {
|
||||||
|
if service["climiter"] != expectedCLimiter {
|
||||||
|
t.Fatalf("expected service climiter %s, got %v", expectedCLimiter, service["climiter"])
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
func startMockSessionForMaxConn(t *testing.T, baseURL string, nodeSecret string, onCommand func(cmdType string, data json.RawMessage) (bool, string)) func() {
|
func startMockSessionForMaxConn(t *testing.T, baseURL string, nodeSecret string, onCommand func(cmdType string, data json.RawMessage) (bool, string)) func() {
|
||||||
t.Helper()
|
t.Helper()
|
||||||
|
|
||||||
|
|||||||
@@ -73,7 +73,7 @@ func TestPerIPSpeedLimitRuntimePayload(t *testing.T) {
|
|||||||
|
|
||||||
var commandMu sync.Mutex
|
var commandMu sync.Mutex
|
||||||
receivedCommands := make([]string, 0)
|
receivedCommands := make([]string, 0)
|
||||||
var addLimitersData json.RawMessage
|
addLimitersData := make([]json.RawMessage, 0)
|
||||||
var updateServiceData json.RawMessage
|
var updateServiceData json.RawMessage
|
||||||
|
|
||||||
stopNode := startMockSessionForMaxConn(t, server.URL, "per-ip-speed-secret", func(cmdType string, data json.RawMessage) (bool, string) {
|
stopNode := startMockSessionForMaxConn(t, server.URL, "per-ip-speed-secret", func(cmdType string, data json.RawMessage) (bool, string) {
|
||||||
@@ -81,7 +81,7 @@ func TestPerIPSpeedLimitRuntimePayload(t *testing.T) {
|
|||||||
defer commandMu.Unlock()
|
defer commandMu.Unlock()
|
||||||
receivedCommands = append(receivedCommands, cmdType)
|
receivedCommands = append(receivedCommands, cmdType)
|
||||||
if cmdType == "AddLimiters" {
|
if cmdType == "AddLimiters" {
|
||||||
addLimitersData = append([]byte(nil), data...)
|
addLimitersData = append(addLimitersData, append([]byte(nil), data...))
|
||||||
}
|
}
|
||||||
if cmdType == "UpdateService" {
|
if cmdType == "UpdateService" {
|
||||||
updateServiceData = append([]byte(nil), data...)
|
updateServiceData = append([]byte(nil), data...)
|
||||||
@@ -123,34 +123,41 @@ func TestPerIPSpeedLimitRuntimePayload(t *testing.T) {
|
|||||||
t.Fatalf("get forward ID: %v", err)
|
t.Fatalf("get forward ID: %v", err)
|
||||||
}
|
}
|
||||||
expectedName := fmt.Sprintf("rule_traffic_limit_%d", forwardID)
|
expectedName := fmt.Sprintf("rule_traffic_limit_%d", forwardID)
|
||||||
expectedLimits := []string{"$ 10.0MB 10.0MB", "0.0.0.0/0 5.0MB 5.0MB", "::/0 5.0MB 5.0MB"}
|
expectedTotalName := fmt.Sprint(totalSpeedID)
|
||||||
|
expectedRuleLimits := []string{"0.0.0.0/0 5.0MB 5.0MB", "::/0 5.0MB 5.0MB"}
|
||||||
|
expectedTotalLimits := []string{"$ 10.0MB 10.0MB"}
|
||||||
|
|
||||||
commandMu.Lock()
|
commandMu.Lock()
|
||||||
defer commandMu.Unlock()
|
defer commandMu.Unlock()
|
||||||
if addLimitersData == nil {
|
if len(addLimitersData) != 2 {
|
||||||
t.Fatalf("expected AddLimiters to be sent. Received: %v", receivedCommands)
|
t.Fatalf("expected AddLimiters to be sent. Received: %v", receivedCommands)
|
||||||
}
|
}
|
||||||
if updateServiceData == nil {
|
if updateServiceData == nil {
|
||||||
t.Fatalf("expected UpdateService to be sent. Received: %v", receivedCommands)
|
t.Fatalf("expected UpdateService to be sent. Received: %v", receivedCommands)
|
||||||
}
|
}
|
||||||
|
|
||||||
var addData map[string]interface{}
|
gotLimiterLimits := make(map[string][]string)
|
||||||
if err := json.Unmarshal(addLimitersData, &addData); err != nil {
|
for _, raw := range addLimitersData {
|
||||||
t.Fatalf("unmarshal AddLimiters data: %v", err)
|
var addData map[string]interface{}
|
||||||
|
if err := json.Unmarshal(raw, &addData); err != nil {
|
||||||
|
t.Fatalf("unmarshal AddLimiters data: %v", err)
|
||||||
|
}
|
||||||
|
name := fmt.Sprint(addData["name"])
|
||||||
|
limits, ok := addData["limits"].([]interface{})
|
||||||
|
if !ok {
|
||||||
|
t.Fatalf("expected limits array, got %T", addData["limits"])
|
||||||
|
}
|
||||||
|
gotLimits := make([]string, 0, len(limits))
|
||||||
|
for _, limit := range limits {
|
||||||
|
gotLimits = append(gotLimits, fmt.Sprint(limit))
|
||||||
|
}
|
||||||
|
gotLimiterLimits[name] = gotLimits
|
||||||
}
|
}
|
||||||
if addData["name"] != expectedName {
|
if !reflect.DeepEqual(gotLimiterLimits[expectedTotalName], expectedTotalLimits) {
|
||||||
t.Fatalf("expected limiter name %s, got %v", expectedName, addData["name"])
|
t.Fatalf("expected total limits %v, got %v", expectedTotalLimits, gotLimiterLimits[expectedTotalName])
|
||||||
}
|
}
|
||||||
limits, ok := addData["limits"].([]interface{})
|
if !reflect.DeepEqual(gotLimiterLimits[expectedName], expectedRuleLimits) {
|
||||||
if !ok {
|
t.Fatalf("expected rule limits %v, got %v", expectedRuleLimits, gotLimiterLimits[expectedName])
|
||||||
t.Fatalf("expected limits array, got %T", addData["limits"])
|
|
||||||
}
|
|
||||||
gotLimits := make([]string, 0, len(limits))
|
|
||||||
for _, limit := range limits {
|
|
||||||
gotLimits = append(gotLimits, fmt.Sprint(limit))
|
|
||||||
}
|
|
||||||
if !reflect.DeepEqual(gotLimits, expectedLimits) {
|
|
||||||
t.Fatalf("expected limits %v, got %v", expectedLimits, gotLimits)
|
|
||||||
}
|
}
|
||||||
|
|
||||||
var services []map[string]interface{}
|
var services []map[string]interface{}
|
||||||
@@ -161,8 +168,9 @@ func TestPerIPSpeedLimitRuntimePayload(t *testing.T) {
|
|||||||
t.Fatalf("expected services in UpdateService")
|
t.Fatalf("expected services in UpdateService")
|
||||||
}
|
}
|
||||||
for _, service := range services {
|
for _, service := range services {
|
||||||
if service["limiter"] != expectedName {
|
expectedLimiter := expectedTotalName + "," + expectedName
|
||||||
t.Fatalf("expected service limiter %s, got %v", expectedName, service["limiter"])
|
if service["limiter"] != expectedLimiter {
|
||||||
|
t.Fatalf("expected service limiter %s, got %v", expectedLimiter, service["limiter"])
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,199 @@
|
|||||||
|
package service
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"fmt"
|
||||||
|
"sort"
|
||||||
|
"strconv"
|
||||||
|
"strings"
|
||||||
|
|
||||||
|
corelimiter "github.com/go-gost/core/limiter"
|
||||||
|
connlimiter "github.com/go-gost/core/limiter/conn"
|
||||||
|
trafficlimiter "github.com/go-gost/core/limiter/traffic"
|
||||||
|
xtraffic "github.com/go-gost/x/limiter/traffic"
|
||||||
|
"github.com/go-gost/x/registry"
|
||||||
|
)
|
||||||
|
|
||||||
|
func resolveTrafficLimiter(names string) trafficlimiter.TrafficLimiter {
|
||||||
|
parts := splitLimiterNames(names)
|
||||||
|
if len(parts) == 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if len(parts) == 1 {
|
||||||
|
return resolveSingleTrafficLimiter(parts[0])
|
||||||
|
}
|
||||||
|
limiters := make([]trafficlimiter.TrafficLimiter, 0, len(parts))
|
||||||
|
for _, part := range parts {
|
||||||
|
if lim := resolveSingleTrafficLimiter(part); lim != nil {
|
||||||
|
limiters = append(limiters, lim)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if len(limiters) == 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if len(limiters) == 1 {
|
||||||
|
return limiters[0]
|
||||||
|
}
|
||||||
|
return &compositeTrafficLimiter{limiters: limiters}
|
||||||
|
}
|
||||||
|
|
||||||
|
func resolveSingleTrafficLimiter(name string) trafficlimiter.TrafficLimiter {
|
||||||
|
lim := registry.TrafficLimiterRegistry().Get(name)
|
||||||
|
if lim != nil {
|
||||||
|
return lim
|
||||||
|
}
|
||||||
|
if val, err := strconv.Atoi(name); err == nil && val > 0 {
|
||||||
|
return xtraffic.NewTrafficLimiter(
|
||||||
|
xtraffic.LimitsOption(fmt.Sprintf("%s %dB %dB", xtraffic.ServiceLimitKey, val, val)),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
return xtraffic.NewTrafficLimiter(
|
||||||
|
xtraffic.LimitsOption(fmt.Sprintf("%s %s %s", xtraffic.ServiceLimitKey, name, name)),
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
func resolveConnLimiter(names string) connlimiter.ConnLimiter {
|
||||||
|
parts := splitLimiterNames(names)
|
||||||
|
if len(parts) == 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if len(parts) == 1 {
|
||||||
|
return registry.ConnLimiterRegistry().Get(parts[0])
|
||||||
|
}
|
||||||
|
limiters := make([]connlimiter.ConnLimiter, 0, len(parts))
|
||||||
|
for _, part := range parts {
|
||||||
|
if lim := registry.ConnLimiterRegistry().Get(part); lim != nil {
|
||||||
|
limiters = append(limiters, lim)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if len(limiters) == 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if len(limiters) == 1 {
|
||||||
|
return limiters[0]
|
||||||
|
}
|
||||||
|
return &compositeConnLimiter{limiters: limiters}
|
||||||
|
}
|
||||||
|
|
||||||
|
func splitLimiterNames(names string) []string {
|
||||||
|
parts := strings.Split(names, ",")
|
||||||
|
out := make([]string, 0, len(parts))
|
||||||
|
for _, part := range parts {
|
||||||
|
if part = strings.TrimSpace(part); part != "" {
|
||||||
|
out = append(out, part)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return out
|
||||||
|
}
|
||||||
|
|
||||||
|
type compositeTrafficLimiter struct {
|
||||||
|
limiters []trafficlimiter.TrafficLimiter
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l *compositeTrafficLimiter) In(ctx context.Context, key string, opts ...corelimiter.Option) trafficlimiter.Limiter {
|
||||||
|
limiters := make([]trafficlimiter.Limiter, 0, len(l.limiters))
|
||||||
|
for _, child := range l.limiters {
|
||||||
|
if lim := child.In(ctx, key, opts...); lim != nil {
|
||||||
|
limiters = append(limiters, lim)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return newCompositeTrafficChildLimiter(limiters)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l *compositeTrafficLimiter) Out(ctx context.Context, key string, opts ...corelimiter.Option) trafficlimiter.Limiter {
|
||||||
|
limiters := make([]trafficlimiter.Limiter, 0, len(l.limiters))
|
||||||
|
for _, child := range l.limiters {
|
||||||
|
if lim := child.Out(ctx, key, opts...); lim != nil {
|
||||||
|
limiters = append(limiters, lim)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return newCompositeTrafficChildLimiter(limiters)
|
||||||
|
}
|
||||||
|
|
||||||
|
type compositeTrafficChildLimiter struct {
|
||||||
|
limiters []trafficlimiter.Limiter
|
||||||
|
}
|
||||||
|
|
||||||
|
func newCompositeTrafficChildLimiter(limiters []trafficlimiter.Limiter) trafficlimiter.Limiter {
|
||||||
|
if len(limiters) == 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if len(limiters) == 1 {
|
||||||
|
return limiters[0]
|
||||||
|
}
|
||||||
|
sort.Slice(limiters, func(i, j int) bool {
|
||||||
|
return limiters[i].Limit() < limiters[j].Limit()
|
||||||
|
})
|
||||||
|
return &compositeTrafficChildLimiter{limiters: limiters}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l *compositeTrafficChildLimiter) Wait(ctx context.Context, n int) int {
|
||||||
|
for _, lim := range l.limiters {
|
||||||
|
if v := lim.Wait(ctx, n); v < n {
|
||||||
|
n = v
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return n
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l *compositeTrafficChildLimiter) Limit() int {
|
||||||
|
if len(l.limiters) == 0 {
|
||||||
|
return 0
|
||||||
|
}
|
||||||
|
return l.limiters[0].Limit()
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l *compositeTrafficChildLimiter) Set(n int) {}
|
||||||
|
|
||||||
|
type compositeConnLimiter struct {
|
||||||
|
limiters []connlimiter.ConnLimiter
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l *compositeConnLimiter) Limiter(key string) connlimiter.Limiter {
|
||||||
|
limiters := make([]connlimiter.Limiter, 0, len(l.limiters))
|
||||||
|
for _, child := range l.limiters {
|
||||||
|
if lim := child.Limiter(key); lim != nil {
|
||||||
|
limiters = append(limiters, lim)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return newCompositeConnChildLimiter(limiters)
|
||||||
|
}
|
||||||
|
|
||||||
|
type compositeConnChildLimiter struct {
|
||||||
|
limiters []connlimiter.Limiter
|
||||||
|
}
|
||||||
|
|
||||||
|
func newCompositeConnChildLimiter(limiters []connlimiter.Limiter) connlimiter.Limiter {
|
||||||
|
if len(limiters) == 0 {
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
if len(limiters) == 1 {
|
||||||
|
return limiters[0]
|
||||||
|
}
|
||||||
|
sort.Slice(limiters, func(i, j int) bool {
|
||||||
|
return limiters[i].Limit() < limiters[j].Limit()
|
||||||
|
})
|
||||||
|
return &compositeConnChildLimiter{limiters: limiters}
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l *compositeConnChildLimiter) Allow(n int) (allowed bool) {
|
||||||
|
var i int
|
||||||
|
for i = range l.limiters {
|
||||||
|
if allowed = l.limiters[i].Allow(n); !allowed {
|
||||||
|
break
|
||||||
|
}
|
||||||
|
}
|
||||||
|
if !allowed && i > 0 && n > 0 {
|
||||||
|
for _, lim := range l.limiters[:i] {
|
||||||
|
lim.Allow(-n)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return allowed
|
||||||
|
}
|
||||||
|
|
||||||
|
func (l *compositeConnChildLimiter) Limit() int {
|
||||||
|
if len(l.limiters) == 0 {
|
||||||
|
return 0
|
||||||
|
}
|
||||||
|
return l.limiters[0].Limit()
|
||||||
|
}
|
||||||
@@ -0,0 +1,79 @@
|
|||||||
|
package service
|
||||||
|
|
||||||
|
import (
|
||||||
|
"context"
|
||||||
|
"io"
|
||||||
|
"testing"
|
||||||
|
|
||||||
|
corelimiter "github.com/go-gost/core/limiter"
|
||||||
|
corelogger "github.com/go-gost/core/logger"
|
||||||
|
xconn "github.com/go-gost/x/limiter/conn"
|
||||||
|
xtraffic "github.com/go-gost/x/limiter/traffic"
|
||||||
|
xlogger "github.com/go-gost/x/logger"
|
||||||
|
"github.com/go-gost/x/registry"
|
||||||
|
)
|
||||||
|
|
||||||
|
func TestResolveTrafficLimiterComposesCommaSeparatedNames(t *testing.T) {
|
||||||
|
const totalName = "test_total_speed_composite"
|
||||||
|
const ruleName = "test_rule_speed_composite"
|
||||||
|
registry.TrafficLimiterRegistry().Unregister(totalName)
|
||||||
|
registry.TrafficLimiterRegistry().Unregister(ruleName)
|
||||||
|
defer registry.TrafficLimiterRegistry().Unregister(totalName)
|
||||||
|
defer registry.TrafficLimiterRegistry().Unregister(ruleName)
|
||||||
|
|
||||||
|
logger := xlogger.NewLogger(xlogger.OutputOption(io.Discard), xlogger.LevelOption(corelogger.ErrorLevel))
|
||||||
|
if err := registry.TrafficLimiterRegistry().Register(totalName, xtraffic.NewTrafficLimiter(xtraffic.LimitsOption("$ 10B 10B"), xtraffic.LoggerOption(logger))); err != nil {
|
||||||
|
t.Fatalf("register total limiter: %v", err)
|
||||||
|
}
|
||||||
|
if err := registry.TrafficLimiterRegistry().Register(ruleName, xtraffic.NewTrafficLimiter(xtraffic.LimitsOption("0.0.0.0/0 3B 3B"), xtraffic.LoggerOption(logger))); err != nil {
|
||||||
|
t.Fatalf("register rule limiter: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
lim := resolveTrafficLimiter(totalName + "," + ruleName)
|
||||||
|
if lim == nil {
|
||||||
|
t.Fatalf("expected composite traffic limiter")
|
||||||
|
}
|
||||||
|
serviceLimiter := lim.In(context.Background(), "192.0.2.1:1000", corelimiter.ScopeOption(corelimiter.ScopeService))
|
||||||
|
if serviceLimiter == nil || serviceLimiter.Limit() != 10 {
|
||||||
|
t.Fatalf("expected service-scope total limiter 10, got %#v", serviceLimiter)
|
||||||
|
}
|
||||||
|
connLimiter := lim.In(context.Background(), "192.0.2.1:1000", corelimiter.ScopeOption(corelimiter.ScopeConn))
|
||||||
|
if connLimiter == nil || connLimiter.Limit() != 3 {
|
||||||
|
t.Fatalf("expected conn-scope per-IP limiter 3, got %#v", connLimiter)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
func TestResolveConnLimiterComposesCommaSeparatedNames(t *testing.T) {
|
||||||
|
const totalName = "test_total_conn_composite"
|
||||||
|
const ruleName = "test_rule_conn_composite"
|
||||||
|
registry.ConnLimiterRegistry().Unregister(totalName)
|
||||||
|
registry.ConnLimiterRegistry().Unregister(ruleName)
|
||||||
|
defer registry.ConnLimiterRegistry().Unregister(totalName)
|
||||||
|
defer registry.ConnLimiterRegistry().Unregister(ruleName)
|
||||||
|
|
||||||
|
logger := xlogger.NewLogger(xlogger.OutputOption(io.Discard), xlogger.LevelOption(corelogger.ErrorLevel))
|
||||||
|
if err := registry.ConnLimiterRegistry().Register(totalName, xconn.NewConnLimiter(xconn.LimitsOption("$ 2"), xconn.LoggerOption(logger))); err != nil {
|
||||||
|
t.Fatalf("register total conn limiter: %v", err)
|
||||||
|
}
|
||||||
|
if err := registry.ConnLimiterRegistry().Register(ruleName, xconn.NewConnLimiter(xconn.LimitsOption("$$ 1"), xconn.LoggerOption(logger))); err != nil {
|
||||||
|
t.Fatalf("register rule conn limiter: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
lim := resolveConnLimiter(totalName + "," + ruleName)
|
||||||
|
if lim == nil {
|
||||||
|
t.Fatalf("expected composite conn limiter")
|
||||||
|
}
|
||||||
|
clientLimiter := lim.Limiter("192.0.2.1")
|
||||||
|
if clientLimiter == nil || clientLimiter.Limit() != 1 {
|
||||||
|
t.Fatalf("expected composite client limiter with strictest limit 1, got %#v", clientLimiter)
|
||||||
|
}
|
||||||
|
if !clientLimiter.Allow(1) {
|
||||||
|
t.Fatalf("expected first connection to be allowed")
|
||||||
|
}
|
||||||
|
if clientLimiter.Allow(1) {
|
||||||
|
t.Fatalf("expected per-IP rule limiter to reject second connection")
|
||||||
|
}
|
||||||
|
if !lim.Limiter("192.0.2.2").Allow(1) {
|
||||||
|
t.Fatalf("expected another client to share total limiter but have independent per-IP capacity")
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -3,7 +3,6 @@ package service
|
|||||||
import (
|
import (
|
||||||
"fmt"
|
"fmt"
|
||||||
"runtime"
|
"runtime"
|
||||||
"strconv"
|
|
||||||
"strings"
|
"strings"
|
||||||
"time"
|
"time"
|
||||||
|
|
||||||
@@ -31,7 +30,6 @@ import (
|
|||||||
logger_parser "github.com/go-gost/x/config/parsing/logger"
|
logger_parser "github.com/go-gost/x/config/parsing/logger"
|
||||||
selector_parser "github.com/go-gost/x/config/parsing/selector"
|
selector_parser "github.com/go-gost/x/config/parsing/selector"
|
||||||
tls_util "github.com/go-gost/x/internal/util/tls"
|
tls_util "github.com/go-gost/x/internal/util/tls"
|
||||||
xtraffic "github.com/go-gost/x/limiter/traffic"
|
|
||||||
cache_limiter "github.com/go-gost/x/limiter/traffic/cache"
|
cache_limiter "github.com/go-gost/x/limiter/traffic/cache"
|
||||||
"github.com/go-gost/x/metadata"
|
"github.com/go-gost/x/metadata"
|
||||||
mdutil "github.com/go-gost/x/metadata/util"
|
mdutil "github.com/go-gost/x/metadata/util"
|
||||||
@@ -185,20 +183,7 @@ func ParseService(cfg *config.ServiceConfig) (service.Service, error) {
|
|||||||
|
|
||||||
var trafficLimiter listener.Option
|
var trafficLimiter listener.Option
|
||||||
if cfg.Limiter != "" {
|
if cfg.Limiter != "" {
|
||||||
lim := registry.TrafficLimiterRegistry().Get(cfg.Limiter)
|
lim := resolveTrafficLimiter(cfg.Limiter)
|
||||||
if lim == nil {
|
|
||||||
// Try to parse as simple number (bandwidth in bytes/sec)
|
|
||||||
if val, err := strconv.Atoi(cfg.Limiter); err == nil && val > 0 {
|
|
||||||
lim = xtraffic.NewTrafficLimiter(
|
|
||||||
xtraffic.LimitsOption(fmt.Sprintf("%s %dB %dB", xtraffic.ServiceLimitKey, val, val)),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
if lim == nil {
|
|
||||||
lim = xtraffic.NewTrafficLimiter(
|
|
||||||
xtraffic.LimitsOption(fmt.Sprintf("%s %s %s", xtraffic.ServiceLimitKey, cfg.Limiter, cfg.Limiter)),
|
|
||||||
)
|
|
||||||
}
|
|
||||||
}
|
|
||||||
trafficLimiter = listener.TrafficLimiterOption(
|
trafficLimiter = listener.TrafficLimiterOption(
|
||||||
cache_limiter.NewCachedTrafficLimiter(
|
cache_limiter.NewCachedTrafficLimiter(
|
||||||
lim,
|
lim,
|
||||||
@@ -216,7 +201,7 @@ func ParseService(cfg *config.ServiceConfig) (service.Service, error) {
|
|||||||
listener.AuthOption(auth_parser.Info(cfg.Listener.Auth)),
|
listener.AuthOption(auth_parser.Info(cfg.Listener.Auth)),
|
||||||
listener.TLSConfigOption(tlsConfig),
|
listener.TLSConfigOption(tlsConfig),
|
||||||
listener.AdmissionOption(xadmission.AdmissionGroup(admissions...)),
|
listener.AdmissionOption(xadmission.AdmissionGroup(admissions...)),
|
||||||
listener.ConnLimiterOption(registry.ConnLimiterRegistry().Get(cfg.CLimiter)),
|
listener.ConnLimiterOption(resolveConnLimiter(cfg.CLimiter)),
|
||||||
listener.ServiceOption(cfg.Name),
|
listener.ServiceOption(cfg.Name),
|
||||||
listener.ProxyProtocolOption(ppv),
|
listener.ProxyProtocolOption(ppv),
|
||||||
listener.StatsOption(pStats),
|
listener.StatsOption(pStats),
|
||||||
|
|||||||
Reference in New Issue
Block a user