Files
zogo/rate/rate.go
T
4566704 45864a61dd feat(rate): 自 go-hua 迁入秒级带宽限速器
- 按秒配额补充令牌, Add 阻塞限速, 配合 conn 使用; 附测试、例程与 README
2026-09-20 12:25:36 +08:00

110 lines
2.6 KiB
Go
Raw Permalink Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
// Package rate 提供按秒配额的带宽限速,常配合 conn 包使用。
package rate
import (
"sync/atomic"
"time"
)
// Rate 带宽限速器:每秒重置配额,配额用尽时 Add 阻塞等待
type Rate struct {
Limit int32 `json:"limit"` // 限制速度 字节 0为不限制
Surplus int32 `json:"surplus"` // 剩余流量 字节
Now int32 `json:"now"` // 接收(下载/下行) 当前流量 字节
Max int32 `json:"max"` // 接收(下载/下行) 最大流量 字节
IsStop chan bool `json:"-"`
}
// NewRate 创建限速器,bandwidth 为带宽上限(单位 Mbps)
func NewRate(bandwidth int) *Rate {
// 带宽 应该是Mbps 要换算成Mbyte
limit := bandwidth * 1024 * 1024 / 8
r := new(Rate)
r.Limit = int32(limit)
r.Now = 0
r.Surplus = int32(limit)
r.IsStop = make(chan bool) // 初始化停止通道,否则 Stop() 会 panic
return r
}
// SetLimit 动态调整带宽上限,bandwidth 为 Mbps(0 表示不限速)
func (r *Rate) SetLimit(bandwidth int) {
// 带宽 应该是Mbps 要换算成Mbyte
limit := 0
if bandwidth > 0 {
limit = bandwidth * 1024 * 1024 / 8
} else {
limit = 0
}
atomic.StoreInt32(&r.Limit, int32(limit))
}
// Start 启动每秒配额重置协程
func (r *Rate) Start() {
go r.proc()
}
// Stop 停止配额重置协程
func (r *Rate) Stop() {
r.IsStop <- true
}
// proc 配额重置循环
func (r *Rate) proc() {
ticker := time.NewTicker(time.Second * 1)
for {
select {
case <-ticker.C:
r.reset()
case <-r.IsStop:
ticker.Stop()
return
}
}
}
// reset 统计上一秒用量并重置配额
func (r *Rate) reset() {
n := r.Limit - atomic.LoadInt32(&r.Surplus)
atomic.StoreInt32(&r.Now, n)
atomic.StoreInt32(&r.Surplus, r.Limit)
if n > 0 {
now := atomic.LoadInt32(&r.Now)
if now > 0 && now > atomic.LoadInt32(&r.Max) {
atomic.StoreInt32(&r.Max, now)
}
}
//fmt.Printf("now:%d limit:%d Surplus:%d \n", n, r.Limit, r.Limit)
}
// GetNow 获取当前秒已用量(字节)
func (r *Rate) GetNow() int {
n := atomic.LoadInt32(&r.Now)
return int(n)
}
// ResetMax 取出并清零当前秒用量
func (r *Rate) ResetMax() int {
n := atomic.SwapInt32(&r.Now, 0)
return int(n)
}
// Add 计入本秒用量;配额用尽时阻塞等待下一秒配额释放
func (r *Rate) Add(size int) {
if atomic.LoadInt32(&r.Surplus) > 0 || atomic.LoadInt32(&r.Limit) == 0 {
atomic.AddInt32(&r.Surplus, -int32(size))
return
}
for {
//fmt.Println("等待")
time.Sleep(time.Millisecond * 10)
if atomic.LoadInt32(&r.Surplus) > 0 {
atomic.AddInt32(&r.Surplus, -int32(size))
return
}
}
}