Files
pangolin/server/internal/agentd/agent.go
T
wangjia a6b35268ab feat(agent): #7 真实流量采集 — clash_api 节点总量 → usage_daily
sing-box 该构建未编入 v2ray_api、clash /connections 也不暴露用户,故无法做
准确 per-user 计费。改取节点累计总量(clash downloadTotal/uploadTotal)做
窗口增量,按当前 dp_uuid 均摊上报(单用户=准确,多用户近似,已注释标注)。

- render.go:节点 sing-box 配置加 experimental.clash_api(loopback+secret)。
- ClashUsageSource:轮询 /connections 总量,delta + 重启回绕保护;clash 的
  up/down 是代理视角,对调为用户视角(下载→bytes_down)。
- SingBox.DpUUIDs() 供归属;Agent.UseClashUsage() 在 cmd/agent 接上。

已节点 live 验证:usage_daily 实时入库,下载增量正确落 bytes_down。
局限:多 dp_uuid 时均摊(准确计费需重编 sing-box 带 with_v2ray_api)。

Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
2026-06-19 06:01:24 +08:00

208 lines
5.6 KiB
Go

package agentd
import (
"context"
"errors"
"math/rand"
"sync/atomic"
"time"
agentv1 "github.com/wangjia/pangolin/server/internal/pb/agentv1"
)
// errResync is returned within a session to force a reconnect + full Register.
var errResync = errors.New("agentd: full resync required")
// Agent is the top-level node daemon. Construct with New and drive with Run.
type Agent struct {
cfg Config
sb *SingBox
nodeUUID atomic.Value // string; written once by Run, read by all loops + tests
dial Dialer
enrollDial EnrollDialer
load LoadSource
usage UsageSource
clock func() time.Time
lastCommandID atomic.Int64
shutdown context.CancelFunc
restarter Restarter
}
// Option customises an Agent (primarily for tests / dependency injection).
type Option func(*Agent)
func WithRestarter(r Restarter) Option { return func(a *Agent) { a.restarter = r } }
func WithDialer(d Dialer) Option { return func(a *Agent) { a.dial = d } }
func WithEnrollDialer(d EnrollDialer) Option { return func(a *Agent) { a.enrollDial = d } }
func WithLoadSource(l LoadSource) Option { return func(a *Agent) { a.load = l } }
func WithUsageSource(u UsageSource) Option { return func(a *Agent) { a.usage = u } }
func WithClock(f func() time.Time) Option { return func(a *Agent) { a.clock = f } }
// New builds an Agent. Production callers pass nothing extra and get mTLS dialers;
// tests inject bufconn dialers, fake restarters, etc.
func New(cfg Config, opts ...Option) *Agent {
cfg = cfg.withDefaults()
a := &Agent{cfg: cfg, clock: time.Now}
for _, o := range opts {
o(a)
}
if a.dial == nil {
if cfg.Insecure {
a.dial = insecureDialer(cfg.ControlPlaneAddr)
} else {
a.dial = NewMTLSDialer(cfg)
}
}
if a.enrollDial == nil {
if cfg.Insecure {
a.enrollDial = EnrollDialer(insecureDialer(cfg.ControlPlaneAddr))
} else {
a.enrollDial = NewEnrollDialer(cfg)
}
}
a.sb = NewSingBox(cfg, a.restarter)
a.sb.clock = a.clock
if a.load == nil {
a.load = defaultLoadSource{sb: a.sb}
}
if a.usage == nil {
a.usage = nopUsageSource{}
}
return a
}
// SingBox exposes the underlying manager (tests / introspection).
func (a *Agent) SingBox() *SingBox { return a.sb }
// UseClashUsage wires the real usage source: reads node-total traffic from the
// local clash_api and reports per-window deltas (attributed across dp_uuids).
// Call after New (needs a.sb). Production main enables this; tests leave it off.
func (a *Agent) UseClashUsage() { a.usage = newClashUsageSource(a.sb) }
// NodeUUID returns the enrolled node UUID (empty until Run enrolls).
func (a *Agent) NodeUUID() string {
v, _ := a.nodeUUID.Load().(string)
return v
}
// Run enrolls (if needed), restores state, then maintains the control-plane
// connection with exponential-backoff reconnect until ctx is cancelled.
func (a *Agent) Run(ctx context.Context) error {
ctx, cancel := context.WithCancel(ctx)
defer cancel()
a.shutdown = cancel
uuid, err := EnsureEnrolled(ctx, a.cfg, a.enrollDial)
if err != nil {
return err
}
a.nodeUUID.Store(uuid)
if err := a.sb.LoadState(); err != nil {
return err
}
go a.sb.Run(ctx)
go a.runTTL(ctx)
a.sb.markDirty() // render whatever we restored before the first connect
backoff := a.cfg.BackoffMin
for {
if ctx.Err() != nil {
return nil
}
registered, err := a.runSession(ctx)
if ctx.Err() != nil {
return nil
}
if registered {
backoff = a.cfg.BackoffMin // healthy session → reset backoff
}
if err != nil && !errors.Is(err, errResync) {
logf("session ended: %v; reconnecting in %s", err, backoff)
}
if !sleepCtx(ctx, jitter(backoff)) {
return nil
}
backoff = nextBackoff(backoff, a.cfg.BackoffMax)
}
}
// runSession establishes one connection: Register (full snapshot), then runs
// heartbeat/subscribe/usage until the first of them fails. Returns whether
// Register succeeded (used to reset backoff).
func (a *Agent) runSession(ctx context.Context) (registered bool, err error) {
conn, err := a.dial(ctx)
if err != nil {
return false, err
}
defer conn.Close()
client := agentv1.NewAgentServiceClient(conn)
snap, err := client.Register(ctx, &agentv1.RegisterRequest{
NodeUUID: a.NodeUUID(),
AgentVersion: a.cfg.AgentVersion,
LocalConfigVersion: a.sb.ConfigVersion(),
})
if err != nil {
return false, err
}
// Full snapshot is authoritative: overwrite local credential table + inbounds.
a.sb.ApplyConfig(snap, true)
if snap.LastCommandID > a.lastCommandID.Load() {
a.lastCommandID.Store(snap.LastCommandID)
}
sctx, cancel := context.WithCancel(ctx)
defer cancel()
errc := make(chan error, 3)
go func() { errc <- a.runHeartbeat(sctx, client) }()
go func() { errc <- a.runSubscribe(sctx, client) }()
go func() { errc <- a.runUsage(sctx, client) }()
first := <-errc
cancel()
<-errc
<-errc
return true, first
}
// ─── backoff helpers ─────────────────────────────────────────────────────────
func nextBackoff(cur, max time.Duration) time.Duration {
next := cur * 2
if next > max {
next = max
}
return next
}
// jitter applies ±20% randomisation to avoid thundering-herd reconnects.
func jitter(d time.Duration) time.Duration {
if d <= 0 {
return 0
}
delta := float64(d) * 0.2
return d + time.Duration((rand.Float64()*2-1)*delta)
}
// sleepCtx sleeps for d unless ctx is cancelled first. Returns false if cancelled.
func sleepCtx(ctx context.Context, d time.Duration) bool {
if d <= 0 {
return ctx.Err() == nil
}
t := time.NewTimer(d)
defer t.Stop()
select {
case <-ctx.Done():
return false
case <-t.C:
return true
}
}