feat: 实现每用户流量统计与定时落库
流量统计从全站按日聚合升级为 per-user 维度,并引入 60 秒定时 落库机制,服务崩溃最多丢失 60 秒数据。 后端: - model/vpn.go: 新增 UserTrafficStat 模型(user_id + date 联合唯一索引) - db/db.go: AutoMigrate 注册新模型 - vpn/tunnel.go: tunnelConn 增加 flushedRx/flushedTx 快照字段 + flushDelta() 增量计算;recordTraffic 改为 per-user 双写 (user_traffic_stats + traffic_stats);连接断开 defer 改增量落库 - vpn/service.go: 新增 flushDone 字段 + trafficFlusher(60s 定时 落库)+ flushAllTraffic;Stop() 先停 flusher 再最终 flush 全部 在线连接增量;TotalLiveTraffic 改返回未落库增量避免与 DB 重复; 新增 UserLiveTraffic(userID);ClientInfo 增加 RxBytes/TxBytes - handler/traffic.go: 新建 5 个 handler(admin 3 个 + user 2 个) - router.go: 注册 5 条新路由 新增 API: - GET /api/admin/traffic/today 所有用户今日流量排行 - GET /api/admin/traffic/history?days=N 全站近 N 天流量历史 - GET /api/admin/traffic/users/:id?days=N 指定用户近 N 天流量 - GET /api/me/traffic/today 自己的今日流量 - GET /api/me/traffic?days=N 自己的近 N 天流量 前端: - 安装 chart.js + vue-chartjs - 新建 TrafficChart.vue 可复用柱状图组件(上行/下行双柱) - AdminView.vue: 在线客户端表格增加 RX/TX 列 + 全站 7 天流量 图表 + 用户今日流量排行表 - ProfileView.vue: 新增流量统计卡片(今日上行/下行/合计)+ 7 天 流量柱状图 - zh.ts/en.ts: 新增 traffic 相关国际化
This commit is contained in:
+61
-2
@@ -27,6 +27,7 @@ type VpnService struct {
|
||||
switchx *PacketSwitch
|
||||
tun *TUNInterface
|
||||
tunDone chan struct{}
|
||||
flushDone chan struct{}
|
||||
running bool
|
||||
startedAt time.Time
|
||||
clients map[*tunnelConn]struct{}
|
||||
@@ -50,16 +51,59 @@ func (s *VpnService) StartedAt() time.Time {
|
||||
return s.startedAt
|
||||
}
|
||||
|
||||
const trafficFlushPeriod = 60 * time.Second
|
||||
|
||||
func (s *VpnService) TotalLiveTraffic() (rx, tx int64) {
|
||||
s.mu.RLock()
|
||||
for c := range s.clients {
|
||||
rx += c.rxBytes.Load()
|
||||
tx += c.txBytes.Load()
|
||||
rx += c.rxBytes.Load() - c.flushedRx.Load()
|
||||
tx += c.txBytes.Load() - c.flushedTx.Load()
|
||||
}
|
||||
s.mu.RUnlock()
|
||||
return
|
||||
}
|
||||
|
||||
func (s *VpnService) UserLiveTraffic(userID uint) (rx, tx int64) {
|
||||
s.mu.RLock()
|
||||
for c := range s.clients {
|
||||
if c.user.ID == userID {
|
||||
rx += c.rxBytes.Load() - c.flushedRx.Load()
|
||||
tx += c.txBytes.Load() - c.flushedTx.Load()
|
||||
}
|
||||
}
|
||||
s.mu.RUnlock()
|
||||
return
|
||||
}
|
||||
|
||||
func (s *VpnService) trafficFlusher() {
|
||||
ticker := time.NewTicker(trafficFlushPeriod)
|
||||
defer ticker.Stop()
|
||||
for {
|
||||
select {
|
||||
case <-ticker.C:
|
||||
s.flushAllTraffic()
|
||||
case <-s.flushDone:
|
||||
return
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (s *VpnService) flushAllTraffic() {
|
||||
s.mu.RLock()
|
||||
conns := make([]*tunnelConn, 0, len(s.clients))
|
||||
for c := range s.clients {
|
||||
conns = append(conns, c)
|
||||
}
|
||||
s.mu.RUnlock()
|
||||
|
||||
for _, c := range conns {
|
||||
rx, tx := c.flushDelta()
|
||||
if rx > 0 || tx > 0 {
|
||||
recordTraffic(c.user.ID, rx, tx)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (s *VpnService) Settings() model.VpnSetting {
|
||||
s.mu.RLock()
|
||||
defer s.mu.RUnlock()
|
||||
@@ -150,11 +194,13 @@ func (s *VpnService) ApplySettings(settings model.VpnSetting, reservations4, res
|
||||
s.switchx = NewPacketSwitch(settings.AllowClientToClient)
|
||||
s.tun = tun
|
||||
s.tunDone = make(chan struct{})
|
||||
s.flushDone = make(chan struct{})
|
||||
s.running = true
|
||||
s.startedAt = time.Now()
|
||||
s.mu.Unlock()
|
||||
|
||||
go s.serveTUN()
|
||||
go s.trafficFlusher()
|
||||
|
||||
subnet4 := ipNet.String()
|
||||
var subnet6Str string
|
||||
@@ -205,10 +251,21 @@ func (s *VpnService) Stop() error {
|
||||
s.running = false
|
||||
tun := s.tun
|
||||
done := s.tunDone
|
||||
flushDone := s.flushDone
|
||||
s.flushDone = nil
|
||||
clients := s.clients
|
||||
s.clients = make(map[*tunnelConn]struct{})
|
||||
s.mu.Unlock()
|
||||
|
||||
if flushDone != nil {
|
||||
close(flushDone)
|
||||
}
|
||||
for c := range clients {
|
||||
rx, tx := c.flushDelta()
|
||||
if rx > 0 || tx > 0 {
|
||||
recordTraffic(c.user.ID, rx, tx)
|
||||
}
|
||||
}
|
||||
for c := range clients {
|
||||
c.close()
|
||||
}
|
||||
@@ -379,6 +436,8 @@ type ClientInfo struct {
|
||||
IP string `json:"ip"`
|
||||
IP6 string `json:"ip6,omitempty"`
|
||||
ConnectedAt string `json:"connected_at"`
|
||||
RxBytes int64 `json:"rx_bytes"`
|
||||
TxBytes int64 `json:"tx_bytes"`
|
||||
}
|
||||
|
||||
func (s *VpnService) KickUser(userID uint) int {
|
||||
|
||||
+27
-2
@@ -40,11 +40,21 @@ type tunnelConn struct {
|
||||
ready atomic.Bool
|
||||
rxBytes atomic.Int64
|
||||
txBytes atomic.Int64
|
||||
flushedRx atomic.Int64
|
||||
flushedTx atomic.Int64
|
||||
}
|
||||
|
||||
func (c *tunnelConn) AssignedIP() net.IP { return c.assignedIP }
|
||||
func (c *tunnelConn) AssignedIP6() net.IP { return c.assignedIP6 }
|
||||
|
||||
func (c *tunnelConn) flushDelta() (rx, tx int64) {
|
||||
curRx := c.rxBytes.Load()
|
||||
curTx := c.txBytes.Load()
|
||||
rx = curRx - c.flushedRx.Swap(curRx)
|
||||
tx = curTx - c.flushedTx.Swap(curTx)
|
||||
return
|
||||
}
|
||||
|
||||
func (c *tunnelConn) WritePacket(data []byte) error {
|
||||
if !c.ready.Load() || len(data) == 0 {
|
||||
return nil
|
||||
@@ -82,6 +92,8 @@ func (c *tunnelConn) info() ClientInfo {
|
||||
Username: c.user.Username,
|
||||
IP: c.assignedIP.String(),
|
||||
ConnectedAt: c.connectedAt.Format("2006-01-02 15:04:05"),
|
||||
RxBytes: c.rxBytes.Load(),
|
||||
TxBytes: c.txBytes.Load(),
|
||||
}
|
||||
if c.assignedIP6 != nil {
|
||||
ci.IP6 = c.assignedIP6.String()
|
||||
@@ -136,7 +148,8 @@ func runTunnel(conn *websocket.Conn, user *model.User) {
|
||||
|
||||
VPN.registerClient(tc)
|
||||
defer func() {
|
||||
recordTraffic(tc.rxBytes.Load(), tc.txBytes.Load())
|
||||
rx, tx := tc.flushDelta()
|
||||
recordTraffic(tc.user.ID, rx, tx)
|
||||
VPN.unregisterClient(tc)
|
||||
}()
|
||||
|
||||
@@ -236,11 +249,23 @@ func runTunnel(conn *websocket.Conn, user *model.User) {
|
||||
}
|
||||
}
|
||||
|
||||
func recordTraffic(rx, tx int64) {
|
||||
func recordTraffic(userID uint, rx, tx int64) {
|
||||
if rx == 0 && tx == 0 {
|
||||
return
|
||||
}
|
||||
today := time.Now().Format("2006-01-02")
|
||||
|
||||
userStat := model.UserTrafficStat{UserID: userID, Date: today, RxBytes: rx, TxBytes: tx}
|
||||
if err := db.DB.Clauses(clause.OnConflict{
|
||||
Columns: []clause.Column{{Name: "user_id"}, {Name: "date"}},
|
||||
DoUpdates: clause.Assignments(map[string]interface{}{
|
||||
"rx_bytes": gorm.Expr("rx_bytes + ?", rx),
|
||||
"tx_bytes": gorm.Expr("tx_bytes + ?", tx),
|
||||
}),
|
||||
}).Create(&userStat).Error; err != nil {
|
||||
log.Printf("记录用户流量失败: %v", err)
|
||||
}
|
||||
|
||||
stat := model.TrafficStat{Date: today, RxBytes: rx, TxBytes: tx}
|
||||
if err := db.DB.Clauses(clause.OnConflict{
|
||||
Columns: []clause.Column{{Name: "date"}},
|
||||
|
||||
Reference in New Issue
Block a user