refactor: batch flow upload processing (#474)

* test: cover flow upload batch semantics

* refactor: batch flow upload persistence

* refactor: batch flow upload processing

* test: harden flow upload batch regression coverage
This commit is contained in:
sagit
2026-04-26 20:45:48 +08:00
committed by GitHub
parent a625884d61
commit 87a1a34ad5
9 changed files with 943 additions and 72 deletions
@@ -0,0 +1,183 @@
package handler
import (
"log"
"sort"
"strings"
"time"
"go-backend/internal/store/model"
"go-backend/internal/store/repo"
)
type flowPolicyTarget struct {
UserID int64
UserTunnelID int64
}
type flowUploadBatch struct {
flowDeltas []repo.FlowUploadCounterDelta
quotaUsage map[int64]int64
policyTargets []flowPolicyTarget
forwardTraffic map[int64]tunnelTrafficDelta
orphanServices map[string]struct{}
peerShareForwardItems map[string]flowItem
peerShareRuntimeItems map[int64]flowItem
}
func (h *Handler) buildFlowUploadBatch(items []flowItem, metas map[int64]repo.FlowUploadForwardMeta) flowUploadBatch {
batch := flowUploadBatch{
quotaUsage: make(map[int64]int64),
forwardTraffic: make(map[int64]tunnelTrafficDelta),
orphanServices: make(map[string]struct{}),
peerShareForwardItems: make(map[string]flowItem),
peerShareRuntimeItems: make(map[int64]flowItem),
}
policySeen := map[flowPolicyTarget]struct{}{}
flowSeen := map[int64]int{}
for _, item := range items {
serviceName := strings.TrimSpace(item.N)
if serviceName == "" || serviceName == "web_api" {
continue
}
if runtimeID, ok := parsePeerShareRuntimeServiceID(serviceName); ok {
merged := batch.peerShareRuntimeItems[runtimeID]
merged.N = serviceName
merged.U += item.U
merged.D += item.D
batch.peerShareRuntimeItems[runtimeID] = merged
continue
}
forwardID, userID, userTunnelID, ok := parseFlowServiceIDs(serviceName)
if !ok {
continue
}
normalized := normalizeForwardRuntimeServiceName(serviceName)
merged := batch.peerShareForwardItems[normalized]
merged.N = normalized
merged.U += item.U
merged.D += item.D
batch.peerShareForwardItems[normalized] = merged
meta, exists := metas[forwardID]
if !exists {
batch.orphanServices[serviceName] = struct{}{}
continue
}
raw := batch.forwardTraffic[forwardID]
raw.bytesIn += item.D
raw.bytesOut += item.U
batch.forwardTraffic[forwardID] = raw
scaledIn := int64(float64(item.D)*meta.TrafficRatio) * meta.TunnelFlow
scaledOut := int64(float64(item.U)*meta.TrafficRatio) * meta.TunnelFlow
if idx, ok := flowSeen[forwardID]; ok {
batch.flowDeltas[idx].InFlow += scaledIn
batch.flowDeltas[idx].OutFlow += scaledOut
} else {
flowSeen[forwardID] = len(batch.flowDeltas)
batch.flowDeltas = append(batch.flowDeltas, repo.FlowUploadCounterDelta{
ForwardID: forwardID,
UserID: userID,
UserTunnelID: userTunnelID,
InFlow: scaledIn,
OutFlow: scaledOut,
})
}
batch.quotaUsage[userID] += scaledIn + scaledOut
target := flowPolicyTarget{UserID: userID, UserTunnelID: userTunnelID}
if _, seen := policySeen[target]; !seen {
policySeen[target] = struct{}{}
batch.policyTargets = append(batch.policyTargets, target)
}
}
sort.Slice(batch.policyTargets, func(i, j int) bool {
if batch.policyTargets[i].UserID == batch.policyTargets[j].UserID {
return batch.policyTargets[i].UserTunnelID < batch.policyTargets[j].UserTunnelID
}
return batch.policyTargets[i].UserID < batch.policyTargets[j].UserID
})
return batch
}
func (h *Handler) applyFlowUploadBatch(nodeID int64, batch flowUploadBatch, now time.Time) {
if h == nil || h.repo == nil {
return
}
h.applyFlowDeltasWithFallback(nodeID, batch.flowDeltas)
for userID, quota := range h.applyQuotaUsageWithFallback(nodeID, batch.quotaUsage, now) {
h.enforceUserQuotaIfNeeded(userID, quota)
}
for _, target := range batch.policyTargets {
if target.UserID <= 0 || target.UserTunnelID <= 0 {
continue
}
h.enforceFlowPolicies(target.UserID, target.UserTunnelID)
}
for serviceName := range batch.orphanServices {
h.sendDeleteOrphanedForwardService(nodeID, serviceName)
}
for serviceName, item := range batch.peerShareForwardItems {
forwardID, _, _, ok := parseFlowServiceIDs(serviceName)
if ok {
h.processPeerShareFlowFromForward(forwardID, nodeID, serviceName, item)
}
}
for runtimeID, item := range batch.peerShareRuntimeItems {
h.processPeerShareFlow(runtimeID, item)
}
}
func (h *Handler) applyFlowDeltasWithFallback(nodeID int64, deltas []repo.FlowUploadCounterDelta) {
if h == nil || h.repo == nil || len(deltas) == 0 {
return
}
if err := h.repo.ApplyFlowUploadDeltasBatch(deltas); err == nil {
return
} else {
log.Printf("flow upload write failed op=flow.batch_apply node_id=%d err=%v", nodeID, err)
}
for _, delta := range deltas {
if err := h.repo.AddFlow(delta.ForwardID, delta.UserID, delta.UserTunnelID, delta.InFlow, delta.OutFlow); err != nil {
log.Printf("flow upload write failed op=flow.single_apply node_id=%d forward_id=%d user_id=%d user_tunnel_id=%d err=%v", nodeID, delta.ForwardID, delta.UserID, delta.UserTunnelID, err)
}
}
}
func (h *Handler) applyQuotaUsageWithFallback(nodeID int64, usages map[int64]int64, now time.Time) map[int64]*model.UserQuotaView {
if h == nil || h.repo == nil || len(usages) == 0 {
return map[int64]*model.UserQuotaView{}
}
quotaViews, err := h.repo.AddUserQuotaUsageBatch(usages, now)
if err == nil {
return quotaViews
}
log.Printf("flow upload write failed op=quota.batch_apply node_id=%d err=%v", nodeID, err)
userIDs := make([]int64, 0, len(usages))
for userID := range usages {
if userID > 0 {
userIDs = append(userIDs, userID)
}
}
sort.Slice(userIDs, func(i, j int) bool { return userIDs[i] < userIDs[j] })
quotaViews = make(map[int64]*model.UserQuotaView, len(userIDs))
for _, userID := range userIDs {
quota, singleErr := h.repo.AddUserQuotaUsage(userID, usages[userID], now)
if singleErr != nil {
log.Printf("flow upload write failed op=quota.single_apply node_id=%d user_id=%d err=%v", nodeID, userID, singleErr)
continue
}
if quota != nil {
quotaViews[userID] = quota
}
}
return quotaViews
}
@@ -0,0 +1,254 @@
package handler
import (
"path/filepath"
"testing"
"time"
"go-backend/internal/store/model"
"go-backend/internal/store/repo"
)
func TestBuildFlowUploadBatchAggregatesForwardQuotaPeerShareAndCleanupTargets(t *testing.T) {
h := &Handler{}
metas := map[int64]repo.FlowUploadForwardMeta{
20: {
ForwardID: 20,
TunnelID: 1,
TrafficRatio: 2,
TunnelFlow: 3,
},
}
batch := h.buildFlowUploadBatch([]flowItem{
{N: "20_2_10", U: 70, D: 50},
{N: "20_2_10_tcp", U: 40, D: 30},
{N: "99_2_10", U: 12, D: 8},
{N: "fed_svc_17", U: 9, D: 1},
}, metas)
if len(batch.flowDeltas) != 1 {
t.Fatalf("expected 1 flow delta, got %d", len(batch.flowDeltas))
}
delta := batch.flowDeltas[0]
if delta.ForwardID != 20 || delta.UserID != 2 || delta.UserTunnelID != 10 {
t.Fatalf("unexpected flow delta identity: %#v", delta)
}
if delta.InFlow != 480 || delta.OutFlow != 660 {
t.Fatalf("expected scaled flow in=480 out=660, got in=%d out=%d", delta.InFlow, delta.OutFlow)
}
if batch.quotaUsage[2] != 1140 {
t.Fatalf("expected quota usage 1140, got %d", batch.quotaUsage[2])
}
if len(batch.policyTargets) != 1 {
t.Fatalf("expected 1 policy target, got %d", len(batch.policyTargets))
}
if batch.policyTargets[0].UserID != 2 || batch.policyTargets[0].UserTunnelID != 10 {
t.Fatalf("unexpected policy target: %#v", batch.policyTargets[0])
}
traffic := batch.forwardTraffic[20]
if traffic.bytesIn != 80 || traffic.bytesOut != 110 {
t.Fatalf("expected raw traffic in=80 out=110, got in=%d out=%d", traffic.bytesIn, traffic.bytesOut)
}
if _, ok := batch.orphanServices["99_2_10"]; !ok {
t.Fatalf("expected orphan service cleanup target for 99_2_10")
}
if item, ok := batch.peerShareForwardItems["99_2_10"]; !ok || item.U != 12 || item.D != 8 {
t.Fatalf("expected orphan forward to remain eligible for peer-share accounting, got %#v ok=%v", item, ok)
}
if item, ok := batch.peerShareForwardItems["20_2_10"]; !ok || item.U != 110 || item.D != 80 {
t.Fatalf("expected merged peer-share forward item, got %#v ok=%v", item, ok)
}
if item, ok := batch.peerShareRuntimeItems[17]; !ok || item.U != 9 || item.D != 1 {
t.Fatalf("expected merged peer-share runtime item, got %#v ok=%v", item, ok)
}
}
func TestApplyFlowUploadBatchContinuesPolicyAndPeerShareSideEffectsWhenQuotaBatchFails(t *testing.T) {
r, err := repo.Open(filepath.Join(t.TempDir(), "flow-upload-batch-quota-fail.db"))
if err != nil {
t.Fatalf("open repo: %v", err)
}
defer r.Close()
now := time.Now()
nowMs := now.UnixMilli()
if err := r.DB().Create(&model.User{ID: 2, User: "flow-user", Pwd: "pwd", RoleID: 1, ExpTime: 2727251700000, Flow: 99999, Num: 99999, CreatedTime: nowMs, Status: 1}).Error; err != nil {
t.Fatalf("seed user: %v", err)
}
if err := r.DB().Create(&model.Tunnel{ID: 1, Name: "tunnel-1", TrafficRatio: 1, Type: 1, Protocol: "tls", Flow: 1, CreatedTime: nowMs, UpdatedTime: nowMs, Status: 1}).Error; err != nil {
t.Fatalf("seed tunnel: %v", err)
}
if err := r.DB().Create(&model.UserTunnel{ID: 10, UserID: 2, TunnelID: 1, Num: 99999, Flow: 0, ExpTime: 2727251700000, Status: 1}).Error; err != nil {
t.Fatalf("seed user tunnel: %v", err)
}
if err := r.DB().Create(&model.Forward{ID: 20, UserID: 2, UserName: "flow-user", Name: "forward-20", TunnelID: 1, RemoteAddr: "1.1.1.1:80", Strategy: "fifo", CreatedTime: nowMs, UpdatedTime: nowMs, Status: 1}).Error; err != nil {
t.Fatalf("seed forward: %v", err)
}
if err := r.CreatePeerShare(&repo.PeerShare{Name: "share", NodeID: 1, Token: "token", MaxBandwidth: 0, CurrentFlow: 0, PortRangeStart: 31000, PortRangeEnd: 31010, IsActive: 1, CreatedTime: nowMs, UpdatedTime: nowMs}); err != nil {
t.Fatalf("create peer share: %v", err)
}
share, err := r.GetPeerShareByToken("token")
if err != nil || share == nil {
t.Fatalf("load peer share: %v", err)
}
if err := r.DB().Exec(`
INSERT INTO peer_share_runtime(share_id, node_id, reservation_id, resource_key, binding_id, role, chain_name, service_name, protocol, strategy, port, target, applied, status, created_time, updated_time)
VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`, share.ID, 1, "svc-r1", "svc-rk1", "", "forward", "", "20_2_10", "tcp", "fifo", 31001, "", 1, 1, nowMs, nowMs).Error; err != nil {
t.Fatalf("insert peer share runtime: %v", err)
}
if err := r.DB().Exec(`
CREATE TRIGGER fail_user_quota_insert
BEFORE INSERT ON user_quota
BEGIN
SELECT RAISE(FAIL, 'quota insert blocked for test');
END;
`).Error; err != nil {
t.Fatalf("create quota failure trigger: %v", err)
}
h := &Handler{repo: r}
h.applyFlowUploadBatch(1, flowUploadBatch{
flowDeltas: []repo.FlowUploadCounterDelta{{ForwardID: 20, UserID: 2, UserTunnelID: 10, InFlow: 80, OutFlow: 120}},
quotaUsage: map[int64]int64{2: 200},
policyTargets: []flowPolicyTarget{{UserID: 2, UserTunnelID: 10}},
peerShareForwardItems: map[string]flowItem{"20_2_10": {N: "20_2_10", U: 120, D: 80}},
}, now)
if got := mustQueryInt(t, r, `SELECT status FROM forward WHERE id = 20`); got != 0 {
t.Fatalf("expected flow-policy enforcement to pause forward after quota failure, got status=%d", got)
}
updatedShare, err := r.GetPeerShare(share.ID)
if err != nil || updatedShare == nil {
t.Fatalf("reload peer share: %v", err)
}
if updatedShare.CurrentFlow != 200 {
t.Fatalf("expected peer-share flow accounting to continue after quota failure, got %d", updatedShare.CurrentFlow)
}
}
func TestApplyFlowUploadBatchContinuesPeerShareSideEffectsWhenFlowBatchFails(t *testing.T) {
r, err := repo.Open(filepath.Join(t.TempDir(), "flow-upload-batch-flow-fail.db"))
if err != nil {
t.Fatalf("open repo: %v", err)
}
defer r.Close()
now := time.Now()
nowMs := now.UnixMilli()
if err := r.DB().Create(&model.User{ID: 2, User: "flow-user", Pwd: "pwd", RoleID: 1, ExpTime: 2727251700000, Flow: 99999, Num: 99999, CreatedTime: nowMs, Status: 1}).Error; err != nil {
t.Fatalf("seed user: %v", err)
}
if err := r.DB().Create(&model.Tunnel{ID: 1, Name: "tunnel-1", TrafficRatio: 1, Type: 1, Protocol: "tls", Flow: 1, CreatedTime: nowMs, UpdatedTime: nowMs, Status: 1}).Error; err != nil {
t.Fatalf("seed tunnel: %v", err)
}
if err := r.DB().Create(&model.UserTunnel{ID: 10, UserID: 2, TunnelID: 1, Num: 99999, Flow: 0, ExpTime: 2727251700000, Status: 1}).Error; err != nil {
t.Fatalf("seed user tunnel: %v", err)
}
if err := r.DB().Create(&model.Forward{ID: 20, UserID: 2, UserName: "flow-user", Name: "forward-20", TunnelID: 1, RemoteAddr: "1.1.1.1:80", Strategy: "fifo", CreatedTime: nowMs, UpdatedTime: nowMs, Status: 1}).Error; err != nil {
t.Fatalf("seed forward: %v", err)
}
if err := r.DB().Create(&model.Forward{ID: 21, UserID: 2, UserName: "flow-user", Name: "forward-21", TunnelID: 1, RemoteAddr: "1.1.1.1:81", Strategy: "fifo", CreatedTime: nowMs, UpdatedTime: nowMs, Status: 1}).Error; err != nil {
t.Fatalf("seed second forward: %v", err)
}
if err := r.CreatePeerShare(&repo.PeerShare{Name: "share", NodeID: 1, Token: "token", MaxBandwidth: 0, CurrentFlow: 0, PortRangeStart: 31000, PortRangeEnd: 31010, IsActive: 1, CreatedTime: nowMs, UpdatedTime: nowMs}); err != nil {
t.Fatalf("create peer share: %v", err)
}
share, err := r.GetPeerShareByToken("token")
if err != nil || share == nil {
t.Fatalf("load peer share: %v", err)
}
if err := r.DB().Exec(`
INSERT INTO peer_share_runtime(share_id, node_id, reservation_id, resource_key, binding_id, role, chain_name, service_name, protocol, strategy, port, target, applied, status, created_time, updated_time)
VALUES(?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
`, share.ID, 1, "svc-r1", "svc-rk1", "", "forward", "", "20_2_10", "tcp", "fifo", 31001, "", 1, 1, nowMs, nowMs).Error; err != nil {
t.Fatalf("insert peer share runtime: %v", err)
}
if err := r.DB().Exec(`
CREATE TRIGGER fail_forward_flow_update
BEFORE UPDATE ON forward
WHEN NEW.id = 21 AND (NEW.in_flow != OLD.in_flow OR NEW.out_flow != OLD.out_flow)
BEGIN
SELECT RAISE(FAIL, 'forward flow update blocked for test');
END;
`).Error; err != nil {
t.Fatalf("create flow failure trigger: %v", err)
}
h := &Handler{repo: r}
h.applyFlowUploadBatch(1, flowUploadBatch{
flowDeltas: []repo.FlowUploadCounterDelta{
{ForwardID: 20, UserID: 2, UserTunnelID: 10, InFlow: 80, OutFlow: 120},
{ForwardID: 21, UserID: 2, UserTunnelID: 10, InFlow: 30, OutFlow: 40},
},
quotaUsage: map[int64]int64{2: 200},
policyTargets: []flowPolicyTarget{{UserID: 2, UserTunnelID: 10}},
peerShareForwardItems: map[string]flowItem{"20_2_10": {N: "20_2_10", U: 120, D: 80}},
}, now)
if got := mustQueryInt(t, r, `SELECT status FROM forward WHERE id = 20`); got != 0 {
t.Fatalf("expected flow-policy enforcement to pause forward after flow batch failure, got status=%d", got)
}
updatedShare, err := r.GetPeerShare(share.ID)
if err != nil || updatedShare == nil {
t.Fatalf("reload peer share: %v", err)
}
if updatedShare.CurrentFlow != 200 {
t.Fatalf("expected peer-share flow accounting to continue after flow batch failure, got %d", updatedShare.CurrentFlow)
}
if got := mustQueryInt(t, r, `SELECT in_flow FROM forward WHERE id = 20`); got != 80 {
t.Fatalf("expected flow fallback to persist forward 20 in_flow=80, got %d", got)
}
if got := mustQueryInt(t, r, `SELECT in_flow FROM forward WHERE id = 21`); got != 0 {
t.Fatalf("expected failed forward 21 delta to remain unapplied, got %d", got)
}
if got := mustQueryInt(t, r, `SELECT in_flow FROM user WHERE id = 2`); got != 80 {
t.Fatalf("expected flow fallback to preserve successful user totals, got %d", got)
}
if got := mustQueryInt(t, r, `SELECT in_flow FROM user_tunnel WHERE id = 10`); got != 80 {
t.Fatalf("expected flow fallback to preserve successful user_tunnel totals, got %d", got)
}
}
func TestApplyFlowUploadBatchFallsBackToPerUserQuotaUpdates(t *testing.T) {
r, err := repo.Open(filepath.Join(t.TempDir(), "flow-upload-batch-quota-fallback.db"))
if err != nil {
t.Fatalf("open repo: %v", err)
}
defer r.Close()
now := time.Now()
nowMs := now.UnixMilli()
dayKey := int64(now.Year()*10000 + int(now.Month())*100 + now.Day())
monthKey := int64(now.Year()*100 + int(now.Month()))
if err := r.DB().Exec(`INSERT INTO user(id, user, pwd, role_id, exp_time, flow, in_flow, out_flow, flow_reset_time, num, created_time, updated_time, status) VALUES(2, 'u2', 'pwd', 1, 2727251700000, 99999, 0, 0, 1, 99999, ?, ?, 1)`, nowMs, nowMs).Error; err != nil {
t.Fatalf("insert user 2: %v", err)
}
if err := r.DB().Exec(`INSERT INTO user(id, user, pwd, role_id, exp_time, flow, in_flow, out_flow, flow_reset_time, num, created_time, updated_time, status) VALUES(3, 'u3', 'pwd', 1, 2727251700000, 99999, 0, 0, 1, 99999, ?, ?, 1)`, nowMs, nowMs).Error; err != nil {
t.Fatalf("insert user 3: %v", err)
}
if err := r.DB().Exec(`INSERT INTO user_quota(user_id, daily_limit_gb, monthly_limit_gb, daily_used_bytes, monthly_used_bytes, day_key, month_key, disabled_by_quota, disabled_at, paused_forward_ids, created_time, updated_time) VALUES(2, 0, 0, 0, 0, ?, ?, 0, 0, '', ?, ?), (3, 0, 0, 0, 0, ?, ?, 0, 0, '', ?, ?)`, dayKey, monthKey, nowMs, nowMs, dayKey, monthKey, nowMs, nowMs).Error; err != nil {
t.Fatalf("insert user quotas: %v", err)
}
if err := r.DB().Exec(`
CREATE TRIGGER fail_user_3_quota_update
BEFORE UPDATE ON user_quota
WHEN NEW.user_id = 3 AND (NEW.daily_used_bytes != OLD.daily_used_bytes OR NEW.monthly_used_bytes != OLD.monthly_used_bytes)
BEGIN
SELECT RAISE(FAIL, 'quota update blocked for user 3');
END;
`).Error; err != nil {
t.Fatalf("create quota fallback trigger: %v", err)
}
h := &Handler{repo: r}
h.applyFlowUploadBatch(1, flowUploadBatch{quotaUsage: map[int64]int64{2: 200, 3: 300}}, now)
if got := mustQueryInt(t, r, `SELECT daily_used_bytes FROM user_quota WHERE user_id = 2`); got != 200 {
t.Fatalf("expected quota fallback to persist user 2 usage, got %d", got)
}
if got := mustQueryInt(t, r, `SELECT daily_used_bytes FROM user_quota WHERE user_id = 3`); got != 0 {
t.Fatalf("expected failed user 3 quota delta to remain unapplied, got %d", got)
}
}
+10 -4
View File
@@ -7,6 +7,7 @@ import (
"encoding/json"
"fmt"
"io"
"log"
"net/http"
"net/url"
"sort"
@@ -797,11 +798,16 @@ func (h *Handler) flowUpload(w http.ResponseWriter, r *http.Request) {
if err == nil && strings.TrimSpace(raw) != "" {
var items []flowItem
if json.Unmarshal([]byte(raw), &items) == nil {
nowMs := time.Now().UnixMilli()
h.recordTunnelMetricsFromFlowItems(node.ID, items, nowMs)
for _, item := range items {
h.processFlowItem(node.ID, item)
now := time.Now()
forwardIDs := collectFlowUploadForwardIDs(items)
metas, metaErr := h.repo.GetFlowUploadForwardMetas(forwardIDs)
if metaErr != nil {
log.Printf("flow upload metadata lookup failed node_id=%d err=%v", node.ID, metaErr)
metas = map[int64]repo.FlowUploadForwardMeta{}
}
batch := h.buildFlowUploadBatch(items, metas)
h.recordTunnelMetricsFromForwardBatch(node.ID, batch.forwardTraffic, metas, now.UnixMilli())
h.applyFlowUploadBatch(node.ID, batch, now)
}
}
@@ -6,6 +6,7 @@ import (
"time"
"go-backend/internal/store/model"
"go-backend/internal/store/repo"
)
type tunnelTrafficDelta struct {
@@ -21,75 +22,42 @@ func unixMilliBucketMinute(nowMs int64) int64 {
return nowMs - (nowMs % minuteMs)
}
func (h *Handler) recordTunnelMetricsFromFlowItems(nodeID int64, items []flowItem, nowMs int64) {
if h == nil || h.repo == nil {
return
}
if nodeID <= 0 || len(items) == 0 {
return
func collectFlowUploadForwardIDs(items []flowItem) []int64 {
ids := make([]int64, 0, len(items))
seen := make(map[int64]struct{}, len(items))
for _, item := range items {
forwardID, _, _, ok := parseFlowServiceIDs(strings.TrimSpace(item.N))
if !ok || forwardID <= 0 {
continue
}
if _, exists := seen[forwardID]; exists {
continue
}
seen[forwardID] = struct{}{}
ids = append(ids, forwardID)
}
return ids
}
func (h *Handler) recordTunnelMetricsFromForwardBatch(nodeID int64, forwardDeltas map[int64]tunnelTrafficDelta, metas map[int64]repo.FlowUploadForwardMeta, nowMs int64) {
if h == nil || h.repo == nil || nodeID <= 0 || len(forwardDeltas) == 0 {
return
}
bucketTs := unixMilliBucketMinute(nowMs)
if bucketTs <= 0 {
return
}
forwardDeltas := make(map[int64]tunnelTrafficDelta)
var skippedParse, skippedZero int
for _, item := range items {
name := strings.TrimSpace(item.N)
if name == "" || name == "web_api" {
continue
}
forwardID, _, _, ok := parseFlowServiceIDs(name)
if !ok {
skippedParse++
continue
}
if item.D == 0 && item.U == 0 {
skippedZero++
continue
}
d := forwardDeltas[forwardID]
d.bytesIn += item.D
d.bytesOut += item.U
forwardDeltas[forwardID] = d
}
if len(forwardDeltas) == 0 {
if len(items) > 0 {
log.Printf("monitoring debug op=tunnel_metric.no_forward_deltas node_id=%d items=%d skipped_parse=%d skipped_zero=%d", nodeID, len(items), skippedParse, skippedZero)
}
return
}
forwardIDs := make([]int64, 0, len(forwardDeltas))
for id := range forwardDeltas {
forwardIDs = append(forwardIDs, id)
}
forwardTunnelMap, err := h.repo.MapForwardIDsToTunnelIDs(forwardIDs)
if err != nil {
log.Printf("monitoring write skipped op=tunnel_metric.map_forward_to_tunnel node_id=%d err=%v", nodeID, err)
return
}
if len(forwardTunnelMap) == 0 {
log.Printf("monitoring debug op=tunnel_metric.no_tunnel_map node_id=%d forward_ids=%v", nodeID, forwardIDs)
return
}
tunnelAgg := make(map[int64]tunnelTrafficDelta)
for forwardID, delta := range forwardDeltas {
tunnelID := forwardTunnelMap[forwardID]
if tunnelID <= 0 {
meta, ok := metas[forwardID]
if !ok || meta.TunnelID <= 0 {
continue
}
a := tunnelAgg[tunnelID]
a.bytesIn += delta.bytesIn
a.bytesOut += delta.bytesOut
tunnelAgg[tunnelID] = a
}
if len(tunnelAgg) == 0 {
return
current := tunnelAgg[meta.TunnelID]
current.bytesIn += delta.bytesIn
current.bytesOut += delta.bytesOut
tunnelAgg[meta.TunnelID] = current
}
metrics := make([]*model.TunnelMetric, 0, len(tunnelAgg))
@@ -98,14 +66,11 @@ func (h *Handler) recordTunnelMetricsFromFlowItems(nodeID int64, items []flowIte
continue
}
metrics = append(metrics, &model.TunnelMetric{
TunnelID: tunnelID,
NodeID: nodeID,
Timestamp: bucketTs,
BytesIn: delta.bytesIn,
BytesOut: delta.bytesOut,
Connections: 0,
Errors: 0,
AvgLatencyMs: 0,
TunnelID: tunnelID,
NodeID: nodeID,
Timestamp: bucketTs,
BytesIn: delta.bytesIn,
BytesOut: delta.bytesOut,
})
}
if len(metrics) == 0 {
@@ -114,7 +79,7 @@ func (h *Handler) recordTunnelMetricsFromFlowItems(nodeID int64, items []flowIte
if err := h.repo.UpsertTunnelMetricBuckets(metrics); err != nil {
log.Printf("monitoring write failed op=tunnel_metric.upsert_buckets node_id=%d bucket_ts=%d count=%d err=%v", nodeID, bucketTs, len(metrics), err)
} else {
log.Printf("monitoring ok op=tunnel_metric.upsert_buckets node_id=%d bucket_ts=%d count=%d", nodeID, bucketTs, len(metrics))
return
}
log.Printf("monitoring ok op=tunnel_metric.upsert_buckets node_id=%d bucket_ts=%d count=%d", nodeID, bucketTs, len(metrics))
}