mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-10-10 19:36:36 +08:00
feat: add CLimiters support for websocket reporter
This commit is contained in:
@@ -102,3 +102,71 @@ type updateLimiterRequest struct {
|
|||||||
type deleteLimiterRequest struct {
|
type deleteLimiterRequest struct {
|
||||||
Limiter string `json:"limiter"`
|
Limiter string `json:"limiter"`
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func createConnLimiter(req createLimiterRequest) error {
|
||||||
|
name := strings.TrimSpace(req.Data.Name)
|
||||||
|
if name == "" {
|
||||||
|
return errors.New("limiter name is required")
|
||||||
|
}
|
||||||
|
req.Data.Name = name
|
||||||
|
|
||||||
|
if registry.ConnLimiterRegistry().IsRegistered(name) {
|
||||||
|
return errors.New("conn limiter " + name + " already exists")
|
||||||
|
}
|
||||||
|
|
||||||
|
v := parser.ParseConnLimiter(&req.Data)
|
||||||
|
|
||||||
|
if err := registry.ConnLimiterRegistry().Register(name, v); err != nil {
|
||||||
|
return errors.New("conn limiter " + name + " already exists")
|
||||||
|
}
|
||||||
|
|
||||||
|
if c := config.Global(); c != nil {
|
||||||
|
c.CLimiters = append(c.CLimiters, &req.Data)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func updateConnLimiter(req updateLimiterRequest) error {
|
||||||
|
name := strings.TrimSpace(req.Limiter)
|
||||||
|
req.Data.Name = name
|
||||||
|
if registry.ConnLimiterRegistry().IsRegistered(name) {
|
||||||
|
registry.ConnLimiterRegistry().Unregister(name)
|
||||||
|
}
|
||||||
|
|
||||||
|
v := parser.ParseConnLimiter(&req.Data)
|
||||||
|
|
||||||
|
if err := registry.ConnLimiterRegistry().Register(name, v); err != nil {
|
||||||
|
return errors.New("conn limiter " + name + " already exists")
|
||||||
|
}
|
||||||
|
|
||||||
|
if c := config.Global(); c != nil {
|
||||||
|
for i := range c.CLimiters {
|
||||||
|
if c.CLimiters[i].Name == name {
|
||||||
|
c.CLimiters[i] = &req.Data
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
}
|
||||||
|
c.CLimiters = append(c.CLimiters, &req.Data)
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|
||||||
|
func deleteConnLimiter(req deleteLimiterRequest) error {
|
||||||
|
name := strings.TrimSpace(req.Limiter)
|
||||||
|
|
||||||
|
if registry.ConnLimiterRegistry().IsRegistered(name) {
|
||||||
|
registry.ConnLimiterRegistry().Unregister(name)
|
||||||
|
}
|
||||||
|
|
||||||
|
if c := config.Global(); c != nil {
|
||||||
|
limiteres := c.CLimiters
|
||||||
|
c.CLimiters = nil
|
||||||
|
for _, s := range limiteres {
|
||||||
|
if s.Name == name {
|
||||||
|
continue
|
||||||
|
}
|
||||||
|
c.CLimiters = append(c.CLimiters, s)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return nil
|
||||||
|
}
|
||||||
|
|||||||
@@ -836,6 +836,18 @@ func (w *WebSocketReporter) routeCommand(cmd CommandMessage) {
|
|||||||
err = w.handleDeleteLimiter(cmd.Data)
|
err = w.handleDeleteLimiter(cmd.Data)
|
||||||
response.Type = "DeleteLimitersResponse"
|
response.Type = "DeleteLimitersResponse"
|
||||||
needSaveConfig = true
|
needSaveConfig = true
|
||||||
|
case "AddCLimiters":
|
||||||
|
err = w.handleAddCLimiter(cmd.Data)
|
||||||
|
response.Type = "AddCLimitersResponse"
|
||||||
|
needSaveConfig = true
|
||||||
|
case "UpdateCLimiters":
|
||||||
|
err = w.handleUpdateCLimiter(cmd.Data)
|
||||||
|
response.Type = "UpdateCLimitersResponse"
|
||||||
|
needSaveConfig = true
|
||||||
|
case "DeleteCLimiters":
|
||||||
|
err = w.handleDeleteCLimiter(cmd.Data)
|
||||||
|
response.Type = "DeleteCLimitersResponse"
|
||||||
|
needSaveConfig = true
|
||||||
|
|
||||||
// TCP Ping 诊断命令(只读,不需要保存配置)
|
// TCP Ping 诊断命令(只读,不需要保存配置)
|
||||||
case "TcpPing":
|
case "TcpPing":
|
||||||
@@ -1130,6 +1142,67 @@ func (w *WebSocketReporter) handleDeleteLimiter(data interface{}) error {
|
|||||||
return deleteLimiter(deleteReq)
|
return deleteLimiter(deleteReq)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
func (w *WebSocketReporter) handleAddCLimiter(data interface{}) error {
|
||||||
|
jsonData, err := json.Marshal(data)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("序列化数据失败: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
var limiterConfig config.LimiterConfig
|
||||||
|
if err := json.Unmarshal(jsonData, &limiterConfig); err != nil {
|
||||||
|
return fmt.Errorf("解析限流器配置失败: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
req := createLimiterRequest{Data: limiterConfig}
|
||||||
|
return createConnLimiter(req)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (w *WebSocketReporter) handleUpdateCLimiter(data interface{}) error {
|
||||||
|
jsonData, err := json.Marshal(data)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("序列化数据失败: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
var updateReq struct {
|
||||||
|
Limiter string `json:"limiter"`
|
||||||
|
Data config.LimiterConfig `json:"data"`
|
||||||
|
}
|
||||||
|
|
||||||
|
if err := json.Unmarshal(jsonData, &updateReq); err != nil {
|
||||||
|
var limiterConfig config.LimiterConfig
|
||||||
|
if err := json.Unmarshal(jsonData, &limiterConfig); err != nil {
|
||||||
|
return fmt.Errorf("解析更新请求失败: %v", err)
|
||||||
|
}
|
||||||
|
updateReq.Limiter = limiterConfig.Name
|
||||||
|
updateReq.Data = limiterConfig
|
||||||
|
}
|
||||||
|
|
||||||
|
req := updateLimiterRequest{
|
||||||
|
Limiter: updateReq.Limiter,
|
||||||
|
Data: updateReq.Data,
|
||||||
|
}
|
||||||
|
return updateConnLimiter(req)
|
||||||
|
}
|
||||||
|
|
||||||
|
func (w *WebSocketReporter) handleDeleteCLimiter(data interface{}) error {
|
||||||
|
jsonData, err := json.Marshal(data)
|
||||||
|
if err != nil {
|
||||||
|
return fmt.Errorf("序列化数据失败: %v", err)
|
||||||
|
}
|
||||||
|
|
||||||
|
var deleteReq deleteLimiterRequest
|
||||||
|
|
||||||
|
if err := json.Unmarshal(jsonData, &deleteReq); err != nil {
|
||||||
|
var limiterName string
|
||||||
|
if err := json.Unmarshal(jsonData, &limiterName); err != nil {
|
||||||
|
return fmt.Errorf("解析删除请求失败: %v", err)
|
||||||
|
}
|
||||||
|
deleteReq.Limiter = limiterName
|
||||||
|
}
|
||||||
|
|
||||||
|
return deleteConnLimiter(deleteReq)
|
||||||
|
}
|
||||||
|
|
||||||
// handleSetProtocol 处理设置屏蔽协议的命令
|
// handleSetProtocol 处理设置屏蔽协议的命令
|
||||||
func (w *WebSocketReporter) handleSetProtocol(data interface{}) error {
|
func (w *WebSocketReporter) handleSetProtocol(data interface{}) error {
|
||||||
jsonData, err := json.Marshal(data)
|
jsonData, err := json.Marshal(data)
|
||||||
|
|||||||
Reference in New Issue
Block a user