From 3373e5ade9bbff27e0902c07ee27fb67fc9198c9 Mon Sep 17 00:00:00 2001 From: sagitchu Date: Tue, 28 Apr 2026 00:18:32 +0800 Subject: [PATCH] fix: preserve shared limiters with per-IP rules --- .../internal/http/handler/control_plane.go | 81 ++++--- .../http/handler/control_plane_test.go | 28 ++- .../contract/max_conn_limit_contract_test.go | 138 ++++++++++++ .../per_ip_speed_limit_contract_test.go | 50 +++-- .../parsing/service/composite_limiter.go | 199 ++++++++++++++++++ .../parsing/service/composite_limiter_test.go | 79 +++++++ go-gost/x/config/parsing/service/parse.go | 19 +- 7 files changed, 511 insertions(+), 83 deletions(-) create mode 100644 go-gost/x/config/parsing/service/composite_limiter.go create mode 100644 go-gost/x/config/parsing/service/composite_limiter_test.go diff --git a/go-backend/internal/http/handler/control_plane.go b/go-backend/internal/http/handler/control_plane.go index 3bb7b25..75af395 100644 --- a/go-backend/internal/http/handler/control_plane.go +++ b/go-backend/internal/http/handler/control_plane.go @@ -292,27 +292,13 @@ func (h *Handler) syncForwardServicesWithWarnings(forward *forwardRecord, method if user != nil && user.MaxConn > 0 { userMaxConn = user.MaxConn } - connLimiterConfig := buildConnLimiterConfig(forward, userMaxConn) + connLimiterConfigs := buildConnLimiterConfigs(forward, userMaxConn) for _, fp := range ports { - runtimeLimiters := forwardRuntimeLimiters{ConnLimiter: connLimiterConfig.Name} - if ipSpeed != nil { - runtimeLimiters.TrafficLimiter = fmt.Sprintf("rule_traffic_limit_%d", forward.ID) - if err := h.ensureTrafficLimiterOnNode(fp.NodeID, runtimeLimiters.TrafficLimiter, speed, 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 - } - } else if limiterID != nil && speed != nil { - runtimeLimiters.TrafficLimiter = strconv.FormatInt(*limiterID, 10) + runtimeLimiters := forwardRuntimeLimiters{ConnLimiter: joinLimiterNames(connLimiterConfigs)} + trafficLimiterNames := make([]string, 0, 2) + if limiterID != nil && speed != nil { + totalLimiterName := strconv.FormatInt(*limiterID, 10) 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 isNodeOfflineOrTimeoutError(err) { @@ -326,9 +312,28 @@ func (h *Handler) syncForwardServicesWithWarnings(forward *forwardRecord, method } 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 { warnings = append(warnings, fmt.Sprintf("节点 %d 连接限制器下发失败: %v", fp.NodeID, err)) } @@ -1876,27 +1881,35 @@ func (h *Handler) ensureConnLimiterOnNode(nodeID int64, cfg forwardLimiterConfig return nil } -func buildConnLimiterConfig(forward *forwardRecord, userMaxConn int) forwardLimiterConfig { +func buildConnLimiterConfigs(forward *forwardRecord, userMaxConn int) []forwardLimiterConfig { if forward == nil { - return forwardLimiterConfig{} + return nil } - limits := make([]string, 0, 2) if forward.MaxConn > 0 { - limits = append(limits, fmt.Sprintf("$ %d", forward.MaxConn)) - } else if userMaxConn > 0 { - limits = append(limits, fmt.Sprintf("$ %d", userMaxConn)) + limits := []string{fmt.Sprintf("$ %d", forward.MaxConn)} + if forward.IPMaxConn > 0 { + 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 { - 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 forwardLimiterConfig{} + return configs +} + +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) - if forward.MaxConn > 0 || forward.IPMaxConn > 0 { - name = fmt.Sprintf("rule_conn_limit_%d", forward.ID) - } - return forwardLimiterConfig{Name: name, Limits: limits} + return strings.Join(names, ",") } func speedToLimitLine(key string, speed int) string { diff --git a/go-backend/internal/http/handler/control_plane_test.go b/go-backend/internal/http/handler/control_plane_test.go index c3fb79a..86c6433 100644 --- a/go-backend/internal/http/handler/control_plane_test.go +++ b/go-backend/internal/http/handler/control_plane_test.go @@ -479,24 +479,30 @@ func TestBuildForwardServiceConfigs_IPv6BindIP(t *testing.T) { } func TestBuildConnLimiterConfigCombinesTotalAndPerIP(t *testing.T) { - cfg := buildConnLimiterConfig(&forwardRecord{ID: 42, UserID: 9, MaxConn: 100, IPMaxConn: 5}, 37) - want := forwardLimiterConfig{Name: "rule_conn_limit_42", Limits: []string{"$ 100", "$$ 5"}} - if !reflect.DeepEqual(cfg, want) { - t.Fatalf("expected %+v, got %+v", want, cfg) + cfgs := buildConnLimiterConfigs(&forwardRecord{ID: 42, UserID: 9, MaxConn: 100, IPMaxConn: 5}, 37) + want := []forwardLimiterConfig{{Name: "rule_conn_limit_42", Limits: []string{"$ 100", "$$ 5"}}} + if !reflect.DeepEqual(cfgs, want) { + t.Fatalf("expected %+v, got %+v", want, cfgs) } } func TestBuildConnLimiterConfigUsesUserTotalWithRulePerIP(t *testing.T) { - cfg := buildConnLimiterConfig(&forwardRecord{ID: 42, UserID: 9, IPMaxConn: 5}, 37) - want := forwardLimiterConfig{Name: "rule_conn_limit_42", Limits: []string{"$ 37", "$$ 5"}} - if !reflect.DeepEqual(cfg, want) { - t.Fatalf("expected %+v, got %+v", want, cfg) + cfgs := buildConnLimiterConfigs(&forwardRecord{ID: 42, UserID: 9, IPMaxConn: 5}, 37) + want := []forwardLimiterConfig{ + {Name: "user_conn_limit_9", Limits: []string{"$ 37"}}, + {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) { - payload := buildTrafficLimiterPayload("rule_traffic_limit_42", intPtr(80), intPtr(40)) - wantLimits := []string{"$ 10.0MB 10.0MB", "0.0.0.0/0 5.0MB 5.0MB", "::/0 5.0MB 5.0MB"} +func TestBuildTrafficLimiterPayloadUsesOnlyPerIPRulesWhenTotalIsSeparate(t *testing.T) { + payload := buildTrafficLimiterPayload("rule_traffic_limit_42", nil, intPtr(40)) + wantLimits := []string{"0.0.0.0/0 5.0MB 5.0MB", "::/0 5.0MB 5.0MB"} if payload["name"] != "rule_traffic_limit_42" { t.Fatalf("expected name rule_traffic_limit_42, got %v", payload["name"]) } diff --git a/go-backend/tests/contract/max_conn_limit_contract_test.go b/go-backend/tests/contract/max_conn_limit_contract_test.go index b356be1..5e1176a 100644 --- a/go-backend/tests/contract/max_conn_limit_contract_test.go +++ b/go-backend/tests/contract/max_conn_limit_contract_test.go @@ -7,6 +7,7 @@ import ( "net/http" "net/http/httptest" "net/url" + "reflect" "strings" "sync" "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() { t.Helper() diff --git a/go-backend/tests/contract/per_ip_speed_limit_contract_test.go b/go-backend/tests/contract/per_ip_speed_limit_contract_test.go index faaa1ec..04f2314 100644 --- a/go-backend/tests/contract/per_ip_speed_limit_contract_test.go +++ b/go-backend/tests/contract/per_ip_speed_limit_contract_test.go @@ -73,7 +73,7 @@ func TestPerIPSpeedLimitRuntimePayload(t *testing.T) { var commandMu sync.Mutex receivedCommands := make([]string, 0) - var addLimitersData json.RawMessage + addLimitersData := make([]json.RawMessage, 0) var updateServiceData json.RawMessage 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() receivedCommands = append(receivedCommands, cmdType) if cmdType == "AddLimiters" { - addLimitersData = append([]byte(nil), data...) + addLimitersData = append(addLimitersData, append([]byte(nil), data...)) } if cmdType == "UpdateService" { updateServiceData = append([]byte(nil), data...) @@ -123,34 +123,41 @@ func TestPerIPSpeedLimitRuntimePayload(t *testing.T) { t.Fatalf("get forward ID: %v", err) } 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() defer commandMu.Unlock() - if addLimitersData == nil { + if len(addLimitersData) != 2 { t.Fatalf("expected AddLimiters to be sent. Received: %v", receivedCommands) } if updateServiceData == nil { t.Fatalf("expected UpdateService to be sent. Received: %v", receivedCommands) } - var addData map[string]interface{} - if err := json.Unmarshal(addLimitersData, &addData); err != nil { - t.Fatalf("unmarshal AddLimiters data: %v", err) + gotLimiterLimits := make(map[string][]string) + for _, raw := range addLimitersData { + 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 { - t.Fatalf("expected limiter name %s, got %v", expectedName, addData["name"]) + if !reflect.DeepEqual(gotLimiterLimits[expectedTotalName], expectedTotalLimits) { + t.Fatalf("expected total limits %v, got %v", expectedTotalLimits, gotLimiterLimits[expectedTotalName]) } - 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)) - } - if !reflect.DeepEqual(gotLimits, expectedLimits) { - t.Fatalf("expected limits %v, got %v", expectedLimits, gotLimits) + if !reflect.DeepEqual(gotLimiterLimits[expectedName], expectedRuleLimits) { + t.Fatalf("expected rule limits %v, got %v", expectedRuleLimits, gotLimiterLimits[expectedName]) } var services []map[string]interface{} @@ -161,8 +168,9 @@ func TestPerIPSpeedLimitRuntimePayload(t *testing.T) { t.Fatalf("expected services in UpdateService") } for _, service := range services { - if service["limiter"] != expectedName { - t.Fatalf("expected service limiter %s, got %v", expectedName, service["limiter"]) + expectedLimiter := expectedTotalName + "," + expectedName + if service["limiter"] != expectedLimiter { + t.Fatalf("expected service limiter %s, got %v", expectedLimiter, service["limiter"]) } } } diff --git a/go-gost/x/config/parsing/service/composite_limiter.go b/go-gost/x/config/parsing/service/composite_limiter.go new file mode 100644 index 0000000..e1304b1 --- /dev/null +++ b/go-gost/x/config/parsing/service/composite_limiter.go @@ -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() +} diff --git a/go-gost/x/config/parsing/service/composite_limiter_test.go b/go-gost/x/config/parsing/service/composite_limiter_test.go new file mode 100644 index 0000000..2f7a2d4 --- /dev/null +++ b/go-gost/x/config/parsing/service/composite_limiter_test.go @@ -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") + } +} diff --git a/go-gost/x/config/parsing/service/parse.go b/go-gost/x/config/parsing/service/parse.go index 636419c..584398a 100644 --- a/go-gost/x/config/parsing/service/parse.go +++ b/go-gost/x/config/parsing/service/parse.go @@ -3,7 +3,6 @@ package service import ( "fmt" "runtime" - "strconv" "strings" "time" @@ -31,7 +30,6 @@ import ( logger_parser "github.com/go-gost/x/config/parsing/logger" selector_parser "github.com/go-gost/x/config/parsing/selector" 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" "github.com/go-gost/x/metadata" mdutil "github.com/go-gost/x/metadata/util" @@ -185,20 +183,7 @@ func ParseService(cfg *config.ServiceConfig) (service.Service, error) { var trafficLimiter listener.Option if cfg.Limiter != "" { - lim := registry.TrafficLimiterRegistry().Get(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)), - ) - } - } + lim := resolveTrafficLimiter(cfg.Limiter) trafficLimiter = listener.TrafficLimiterOption( cache_limiter.NewCachedTrafficLimiter( lim, @@ -216,7 +201,7 @@ func ParseService(cfg *config.ServiceConfig) (service.Service, error) { listener.AuthOption(auth_parser.Info(cfg.Listener.Auth)), listener.TLSConfigOption(tlsConfig), listener.AdmissionOption(xadmission.AdmissionGroup(admissions...)), - listener.ConnLimiterOption(registry.ConnLimiterRegistry().Get(cfg.CLimiter)), + listener.ConnLimiterOption(resolveConnLimiter(cfg.CLimiter)), listener.ServiceOption(cfg.Name), listener.ProxyProtocolOption(ppv), listener.StatsOption(pStats),