feat(obs): 去掉宿主机网卡趋势,磁盘读写改按速率展示

Agent 不再采集网卡累计字节,看板与节点网络图仅保留访问日志已提供/接收。
磁盘 IO 按小时换算为 B/s 曲线,摘要为近 24 小时平均速率,并同步 Swagger。
This commit is contained in:
ryan
2026-07-18 13:17:13 +08:00
parent 55f8c9a527
commit 26057514a1
20 changed files with 87 additions and 278 deletions
@@ -57,7 +57,7 @@ func BuildSnapshot(cfg *config.Config, stateStore *state.Store) *protocol.NodeMe
metric.StorageTotalBytes = storageTotal
metric.StorageUsedBytes = storageUsed
metric.NetworkRxBytes, metric.NetworkTxBytes = edgeobs.ReadLinuxNetworkTotals()
// Host NIC totals are not collected (product no longer surfaces host NIC trends).
metric.DiskReadBytes, metric.DiskWriteBytes = edgeobs.ReadLinuxDiskTotals()
if stateStore == nil {
@@ -164,8 +164,7 @@ func persistNodeMetricSnapshot(ctx context.Context, nodeID string, snapshot *Nod
StorageTotalBytes: snapshot.StorageTotalBytes,
DiskReadBytes: snapshot.DiskReadBytes,
DiskWriteBytes: snapshot.DiskWriteBytes,
NetworkRxBytes: snapshot.NetworkRxBytes,
NetworkTxBytes: snapshot.NetworkTxBytes,
// NetworkRx/Tx no longer collected from agents; CH columns remain 0.
}
return model.InsertOpenFlareMetricSnapshot(ctx, record)
}
+1 -4
View File
@@ -354,12 +354,9 @@ func compressCapacityTrendPoints(points []observability.CapacityTrendPoint) [][]
func compressNetworkTrendPoints(points []observability.NetworkTrendPoint) [][]any {
rows := make([][]any, 0, len(points))
for _, point := range points {
// Compact layout (stable positions):
// [0] bucket, [1] host_rx, [2] host_tx, [3] bytes_received, [4] bytes_provided, [5] reported_nodes
// Compact layout: [0] bucket, [1] bytes_received, [2] bytes_provided, [3] reported_nodes
rows = append(rows, []any{
point.BucketStartedAt,
point.NetworkRxBytes,
point.NetworkTxBytes,
point.BytesReceived,
point.BytesProvided,
point.ReportedNodes,
@@ -149,7 +149,7 @@ func TestGetOverviewStructure(t *testing.T) {
require.Len(t, row, 4)
}
for _, row := range overview.Trends.Network24h {
require.Len(t, row, 6)
require.Len(t, row, 4)
}
for _, row := range overview.Trends.DiskIO24h {
require.Len(t, row, 4)
@@ -51,8 +51,6 @@ type NodeMetricSnapshotView struct {
StorageTotalBytes int64 `json:"storage_total_bytes"`
DiskReadBytes int64 `json:"disk_read_bytes"`
DiskWriteBytes int64 `json:"disk_write_bytes"`
NetworkRxBytes int64 `json:"network_rx_bytes"`
NetworkTxBytes int64 `json:"network_tx_bytes"`
OpenrestyConnections int64 `json:"openresty_connections"`
}
@@ -95,13 +93,10 @@ type CapacityTrendPoint struct {
ReportedNodes int `json:"reported_nodes"`
}
// NetworkTrendPoint is a network trend bucket.
// Host network_* is L3 (宿主机网卡).
// bytes_received/provided are L1 business bytes from access logs.
// NetworkTrendPoint is a business-byte trend bucket from access logs (L1).
// Host NIC trends are intentionally not exposed.
type NetworkTrendPoint struct {
BucketStartedAt time.Time `json:"bucket_started_at"`
NetworkRxBytes int64 `json:"network_rx_bytes"`
NetworkTxBytes int64 `json:"network_tx_bytes"`
BytesReceived int64 `json:"bytes_received"` // sum(request_length)
BytesProvided int64 `json:"bytes_provided"` // sum(bytes_sent)
ReportedNodes int `json:"reported_nodes"`
@@ -135,11 +130,6 @@ type diskCounterState struct {
seen bool
}
type networkCounterState struct {
rx int64
tx int64
seen bool
}
func buildTrafficWindowSummaryFromAccessLogs(
ctx context.Context,
@@ -194,8 +184,6 @@ func BuildMetricSnapshotViews(
StorageTotalBytes: snapshot.StorageTotalBytes,
DiskReadBytes: snapshot.DiskReadBytes,
DiskWriteBytes: snapshot.DiskWriteBytes,
NetworkRxBytes: snapshot.NetworkRxBytes,
NetworkTxBytes: snapshot.NetworkTxBytes,
}
if matched := matchEdgeHealth(snapshot.CapturedAt, edgeHealth); matched != nil {
view.OpenrestyConnections = matched.Connections
@@ -284,8 +272,8 @@ func buildHealthSummary(
}
// BuildNodeTrends builds 24h trend series.
// Business traffic (requests/errors/UV and provided/received bytes) comes from access logs.
// Host capacity/disk/network come from metric snapshots (hourly when available).
// Business traffic (requests/errors and provided/received bytes) comes from access logs.
// Host capacity/disk come from metric snapshots (hourly when available). Host NIC is not tracked.
func BuildNodeTrends(
ctx context.Context,
now time.Time,
@@ -296,8 +284,7 @@ func BuildNodeTrends(
trafficTrend := BuildTrafficTrendPointsFromAccessLogs(ctx, now, nodeID, trendSince)
capacityTrend := BuildCapacityTrendPoints(now, snapshots)
networkTrend := BuildNetworkTrendPoints(now, snapshots)
// Overlay L1 business bytes onto network points.
networkTrend := emptyNetworkTrendPoints(now)
applyAccessLogBytesToNetworkTrend(ctx, now, nodeID, trendSince, networkTrend)
diskIOTrend := BuildDiskIOTrendPoints(now, snapshots)
@@ -305,8 +292,6 @@ func BuildNodeTrends(
if metricErr == nil && len(metricHourly) > 0 {
capacityTrend = BuildCapacityTrendPointsFromHourly(now, metricHourly)
diskIOTrend = BuildDiskIOTrendPointsFromHourly(now, metricHourly)
networkTrend = BuildNetworkTrendPointsFromHourly(now, metricHourly)
applyAccessLogBytesToNetworkTrend(ctx, now, nodeID, trendSince, networkTrend)
}
return NodeTrends{
@@ -522,85 +507,12 @@ func BuildCapacityTrendPointsFromHourly(now time.Time, hourly []*model.OpenFlare
return points
}
// BuildNetworkTrendPoints builds 24h host-network trend buckets.
// Host network counters must be process-lifetime cumulative values; this function
// converts consecutive samples into deltas.
func BuildNetworkTrendPoints(
now time.Time,
snapshots []*model.OpenFlareMetricSnapshot,
) []NetworkTrendPoint {
start := trendWindowStart(now)
points := make([]NetworkTrendPoint, observabilityTrendBuckets)
accumulators := make([]snapshotTrendAccumulator, observabilityTrendBuckets)
for index := range points {
points[index].BucketStartedAt = start.Add(time.Duration(index) * time.Hour)
accumulators[index].nodes = make(map[string]struct{})
}
sort.Slice(snapshots, func(i int, j int) bool {
if snapshots[i].CapturedAt.Equal(snapshots[j].CapturedAt) {
return snapshots[i].NodeID < snapshots[j].NodeID
}
return snapshots[i].CapturedAt.Before(snapshots[j].CapturedAt)
})
previousHostByNode := make(map[string]networkCounterState, len(snapshots))
for _, snapshot := range snapshots {
if snapshot == nil {
continue
}
nodeKey := snapshot.NodeID
if nodeKey == "" {
nodeKey = unknownTrendNodeKey
}
previous := previousHostByNode[nodeKey]
previousHostByNode[nodeKey] = networkCounterState{
rx: snapshot.NetworkRxBytes,
tx: snapshot.NetworkTxBytes,
seen: true,
}
if !previous.seen {
continue
}
index, ok := trendBucketIndex(snapshot.CapturedAt, start)
if !ok {
continue
}
points[index].NetworkRxBytes += nonNegativeDelta(snapshot.NetworkRxBytes, previous.rx)
points[index].NetworkTxBytes += nonNegativeDelta(snapshot.NetworkTxBytes, previous.tx)
if snapshot.NodeID != "" {
accumulators[index].nodes[snapshot.NodeID] = struct{}{}
}
}
for index := range points {
points[index].ReportedNodes = len(accumulators[index].nodes)
}
return points
}
// BuildNetworkTrendPointsFromHourly builds 24h host-network trend buckets from metric hourly aggregates.
// Business bytes (已提供/接收) are applied separately via applyAccessLogBytesToNetworkTrend.
func BuildNetworkTrendPointsFromHourly(
now time.Time,
metricHourly []*model.OpenFlareMetricHourly,
) []NetworkTrendPoint {
func emptyNetworkTrendPoints(now time.Time) []NetworkTrendPoint {
start := trendWindowStart(now)
points := make([]NetworkTrendPoint, observabilityTrendBuckets)
for index := range points {
points[index].BucketStartedAt = start.Add(time.Duration(index) * time.Hour)
}
for _, row := range metricHourly {
if row == nil {
continue
}
index, ok := trendBucketIndex(row.Hour, start)
if !ok {
continue
}
points[index].NetworkRxBytes += row.NetworkRxBytes
points[index].NetworkTxBytes += row.NetworkTxBytes
if row.ReportedNodes > points[index].ReportedNodes {
points[index].ReportedNodes = row.ReportedNodes
}
}
return points
}
@@ -113,27 +113,16 @@ func TestBuildCapacityTrendPointsFromHourlyFillsBuckets(t *testing.T) {
}
}
func TestBuildNetworkTrendPointsUsesCounterDeltas(t *testing.T) {
func TestEmptyNetworkTrendPointsHas24Buckets(t *testing.T) {
t.Parallel()
now := time.Date(2026, 7, 10, 9, 30, 0, 0, time.UTC)
base := now.Truncate(time.Hour)
snapshots := []*model.OpenFlareMetricSnapshot{
{NodeID: "n1", CapturedAt: base.Add(10 * time.Minute), NetworkRxBytes: 1000, NetworkTxBytes: 2000},
{NodeID: "n1", CapturedAt: base.Add(20 * time.Minute), NetworkRxBytes: 1500, NetworkTxBytes: 2600},
points := emptyNetworkTrendPoints(now)
if len(points) != observabilityTrendBuckets {
t.Fatalf("len(points) = %d, want %d", len(points), observabilityTrendBuckets)
}
points := BuildNetworkTrendPoints(now, snapshots)
current := points[len(points)-1]
if current.NetworkRxBytes != 500 {
t.Fatalf("network_rx_bytes = %d, want 500", current.NetworkRxBytes)
}
if current.NetworkTxBytes != 600 {
t.Fatalf("network_tx_bytes = %d, want 600", current.NetworkTxBytes)
}
if current.BytesReceived != 0 || current.BytesProvided != 0 {
t.Fatalf("business bytes should be 0 without access logs overlay, got received=%d provided=%d",
current.BytesReceived, current.BytesProvided)
if !points[0].BucketStartedAt.Before(points[len(points)-1].BucketStartedAt) {
t.Fatalf("bucket order invalid: first=%v last=%v", points[0].BucketStartedAt, points[len(points)-1].BucketStartedAt)
}
}
+2 -2
View File
@@ -92,8 +92,8 @@ func TestHeartbeatPayloadBindingAndFrpsObservationInsert(t *testing.T) {
Snapshot: &agent.NodeMetricSnapshot{
CapturedAtUnix: now.Unix(),
CPUUsagePercent: 12.5,
NetworkRxBytes: 1024,
NetworkTxBytes: 2048,
DiskReadBytes: 100,
DiskWriteBytes: 200,
},
HealthEvents: []agent.NodeHealthEvent{},
})
@@ -48,9 +48,8 @@ func BuildSnapshot(cfg *config.Config, stateStore *state.Store) *service.AgentNo
metric.MemoryTotalBytes, metric.MemoryUsedBytes = edgeobs.ReadMemInfo()
metric.StorageTotalBytes, metric.StorageUsedBytes = edgeobs.StatFilesystem(cfg.DataDir)
metric.NetworkRxBytes, metric.NetworkTxBytes = edgeobs.ReadLinuxNetworkTotals()
// Host NIC totals are not collected.
metric.DiskReadBytes, metric.DiskWriteBytes = edgeobs.ReadLinuxDiskTotals()
if stateStore == nil {
return metric
}