mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-10-08 02:36:37 +08:00
2.0版本测试
This commit is contained in:
+1
-1
@@ -119,7 +119,7 @@ func main() {
|
||||
log := xlogger.NewLogger()
|
||||
logger.SetDefault(log)
|
||||
|
||||
wsReporter := socket.StartWebSocketReporterWithConfig(config.Addr, config.Secret, config.Http, config.Tls, config.Socks, "1.2.3")
|
||||
wsReporter := socket.StartWebSocketReporterWithConfig(config.Addr, config.Secret, config.Http, config.Tls, config.Socks, "2.0.0")
|
||||
defer wsReporter.Stop()
|
||||
service.SetHTTPReportURL(config.Addr, config.Secret)
|
||||
|
||||
|
||||
@@ -0,0 +1,206 @@
|
||||
package service
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// GlobalTrafficManager 全局流量管理器(所有服务共享)
|
||||
type GlobalTrafficManager struct {
|
||||
mu sync.RWMutex
|
||||
serviceTraffic map[string]*ServiceTraffic // key: 服务名, value: 流量数据
|
||||
ctx context.Context
|
||||
cancel context.CancelFunc
|
||||
reportTicker *time.Ticker
|
||||
}
|
||||
|
||||
// ServiceTraffic 单个服务的流量累积
|
||||
type ServiceTraffic struct {
|
||||
mu sync.Mutex
|
||||
ServiceName string
|
||||
UpBytes int64 // 上行流量(累积)
|
||||
DownBytes int64 // 下行流量(累积)
|
||||
}
|
||||
|
||||
var (
|
||||
globalManager *GlobalTrafficManager
|
||||
globalManagerOnce sync.Once
|
||||
)
|
||||
|
||||
// GetGlobalTrafficManager 获取全局流量管理器单例
|
||||
func GetGlobalTrafficManager() *GlobalTrafficManager {
|
||||
globalManagerOnce.Do(func() {
|
||||
ctx, cancel := context.WithCancel(context.Background())
|
||||
globalManager = &GlobalTrafficManager{
|
||||
serviceTraffic: make(map[string]*ServiceTraffic),
|
||||
ctx: ctx,
|
||||
cancel: cancel,
|
||||
reportTicker: time.NewTicker(5 * time.Second),
|
||||
}
|
||||
// 启动定时上报协程
|
||||
go globalManager.startReporting()
|
||||
})
|
||||
return globalManager
|
||||
}
|
||||
|
||||
// AddTraffic 添加流量到指定服务(由各服务调用)
|
||||
func (m *GlobalTrafficManager) AddTraffic(serviceName string, upBytes, downBytes int64) {
|
||||
if upBytes == 0 && downBytes == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
|
||||
// 获取或创建服务流量记录
|
||||
traffic, exists := m.serviceTraffic[serviceName]
|
||||
if !exists {
|
||||
traffic = &ServiceTraffic{
|
||||
ServiceName: serviceName,
|
||||
}
|
||||
m.serviceTraffic[serviceName] = traffic
|
||||
}
|
||||
|
||||
// 累加流量
|
||||
traffic.mu.Lock()
|
||||
traffic.UpBytes += upBytes
|
||||
traffic.DownBytes += downBytes
|
||||
traffic.mu.Unlock()
|
||||
}
|
||||
|
||||
// startReporting 启动定时上报协程(每5秒执行一次)
|
||||
func (m *GlobalTrafficManager) startReporting() {
|
||||
|
||||
for {
|
||||
select {
|
||||
case <-m.reportTicker.C:
|
||||
m.collectAndReport()
|
||||
|
||||
case <-m.ctx.Done():
|
||||
fmt.Printf("⏹️ 全局流量上报器已停止\n")
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// collectAndReport 收集所有服务流量并合并上报
|
||||
func (m *GlobalTrafficManager) collectAndReport() {
|
||||
m.mu.Lock()
|
||||
|
||||
// 如果没有流量,直接返回
|
||||
if len(m.serviceTraffic) == 0 {
|
||||
m.mu.Unlock()
|
||||
return
|
||||
}
|
||||
|
||||
// 复制当前所有流量数据(避免长时间持锁)
|
||||
trafficSnapshot := make(map[string]*ServiceTraffic)
|
||||
reportData := make(map[string]struct {
|
||||
up int64
|
||||
down int64
|
||||
})
|
||||
|
||||
for name, traffic := range m.serviceTraffic {
|
||||
traffic.mu.Lock()
|
||||
if traffic.UpBytes > 0 || traffic.DownBytes > 0 {
|
||||
trafficSnapshot[name] = traffic
|
||||
reportData[name] = struct {
|
||||
up int64
|
||||
down int64
|
||||
}{
|
||||
up: traffic.UpBytes,
|
||||
down: traffic.DownBytes,
|
||||
}
|
||||
}
|
||||
traffic.mu.Unlock()
|
||||
}
|
||||
m.mu.Unlock()
|
||||
|
||||
// 如果没有需要上报的流量,返回
|
||||
if len(reportData) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
// 构建上报数据数组(保持每个服务独立)
|
||||
reportItems := make([]TrafficReportItem, 0, len(reportData))
|
||||
var totalUp, totalDown int64
|
||||
|
||||
for serviceName, data := range reportData {
|
||||
reportItems = append(reportItems, TrafficReportItem{
|
||||
N: serviceName, // 保持服务名不变
|
||||
U: data.up,
|
||||
D: data.down,
|
||||
})
|
||||
totalUp += data.up
|
||||
totalDown += data.down
|
||||
}
|
||||
|
||||
// 批量发送上报请求(一次HTTP请求包含所有服务)
|
||||
success, err := sendBatchTrafficReport(m.ctx, reportItems)
|
||||
if err != nil {
|
||||
fmt.Printf("❌ 全局流量上报失败: %v (总流量: ↑%d ↓%d, %d个服务)\n", err, totalUp, totalDown, len(reportItems))
|
||||
return
|
||||
}
|
||||
|
||||
if !success {
|
||||
fmt.Printf("⚠️ 全局流量上报未成功 (总流量: ↑%d ↓%d, %d个服务)\n", totalUp, totalDown, len(reportItems))
|
||||
return
|
||||
}
|
||||
|
||||
// 上报成功,清空已上报的流量
|
||||
m.clearReportedTraffic(reportData)
|
||||
}
|
||||
|
||||
// clearReportedTraffic 清空已成功上报的流量
|
||||
func (m *GlobalTrafficManager) clearReportedTraffic(reportedData map[string]struct {
|
||||
up int64
|
||||
down int64
|
||||
}) {
|
||||
m.mu.Lock()
|
||||
defer m.mu.Unlock()
|
||||
|
||||
for serviceName, reported := range reportedData {
|
||||
if traffic, exists := m.serviceTraffic[serviceName]; exists {
|
||||
traffic.mu.Lock()
|
||||
// 减去已上报的流量
|
||||
traffic.UpBytes -= reported.up
|
||||
traffic.DownBytes -= reported.down
|
||||
|
||||
// 如果流量归零,从map中删除该服务记录(避免内存泄漏)
|
||||
if traffic.UpBytes <= 0 && traffic.DownBytes <= 0 {
|
||||
traffic.mu.Unlock()
|
||||
delete(m.serviceTraffic, serviceName)
|
||||
} else {
|
||||
traffic.mu.Unlock()
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// Stop 停止全局流量管理器
|
||||
func (m *GlobalTrafficManager) Stop() {
|
||||
if m.reportTicker != nil {
|
||||
m.reportTicker.Stop()
|
||||
}
|
||||
if m.cancel != nil {
|
||||
m.cancel()
|
||||
}
|
||||
fmt.Printf("🛑 全局流量管理器已停止\n")
|
||||
}
|
||||
|
||||
// GetServiceTraffic 获取指定服务的当前流量(用于调试)
|
||||
func (m *GlobalTrafficManager) GetServiceTraffic(serviceName string) (upBytes, downBytes int64) {
|
||||
m.mu.RLock()
|
||||
defer m.mu.RUnlock()
|
||||
|
||||
if traffic, exists := m.serviceTraffic[serviceName]; exists {
|
||||
traffic.mu.Lock()
|
||||
upBytes = traffic.UpBytes
|
||||
downBytes = traffic.DownBytes
|
||||
traffic.mu.Unlock()
|
||||
}
|
||||
return
|
||||
}
|
||||
|
||||
@@ -403,19 +403,15 @@ func (s *defaultService) observeStats(ctx context.Context) {
|
||||
TotalErrs: st.Get(stats.KindTotalErrs),
|
||||
},
|
||||
}
|
||||
|
||||
// 将流量累积到全局管理器,而不是立即上报
|
||||
if outputBytes > 0 || inputBytes > 0 {
|
||||
reportItems := TrafficReportItem{
|
||||
N: s.name,
|
||||
U: int64(outputBytes),
|
||||
D: int64(inputBytes),
|
||||
}
|
||||
success, err := sendTrafficReport(ctx, reportItems)
|
||||
if err != nil {
|
||||
fmt.Printf("发送流量报告失败: %v", err)
|
||||
} else if success {
|
||||
if xstats, ok := st.(*xstats.Stats); ok {
|
||||
xstats.ResetTraffic(st.Get(stats.KindInputBytes)-inputBytes, st.Get(stats.KindOutputBytes)-outputBytes)
|
||||
}
|
||||
globalManager := GetGlobalTrafficManager()
|
||||
globalManager.AddTraffic(s.name, int64(outputBytes), int64(inputBytes))
|
||||
|
||||
// 立即重置流量计数(因为已经记录到全局管理器中)
|
||||
if xstats, ok := st.(*xstats.Stats); ok {
|
||||
xstats.ResetTraffic(st.Get(stats.KindInputBytes)-inputBytes, st.Get(stats.KindOutputBytes)-outputBytes)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -41,8 +41,8 @@ func SetHTTPReportURL(addr string, secret string) {
|
||||
}
|
||||
}
|
||||
|
||||
// sendTrafficReport 发送流量报告到HTTP接口
|
||||
func sendTrafficReport(ctx context.Context, reportItems TrafficReportItem) (bool, error) {
|
||||
// sendBatchTrafficReport 批量发送多个服务的流量报告到HTTP接口
|
||||
func sendBatchTrafficReport(ctx context.Context, reportItems []TrafficReportItem) (bool, error) {
|
||||
jsonData, err := json.Marshal(reportItems)
|
||||
if err != nil {
|
||||
return false, fmt.Errorf("序列化报告数据失败: %v", err)
|
||||
@@ -112,6 +112,7 @@ func sendTrafficReport(ctx context.Context, reportItems TrafficReportItem) (bool
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
// sendConfigReport 发送配置报告到HTTP接口
|
||||
func sendConfigReport(ctx context.Context) (bool, error) {
|
||||
if configReportURL == "" {
|
||||
|
||||
Reference in New Issue
Block a user