mirror of
https://github.com/Sagit-chu/flvx.git
synced 2026-10-05 01:26:37 +08:00
chore: update monitor rendering and release rc2
This commit is contained in:
@@ -36,7 +36,7 @@ func TestRecordNodeMetric(t *testing.T) {
|
||||
svc.RecordNodeMetric(1, info)
|
||||
svc.flushNodeMetrics()
|
||||
|
||||
metrics, err := r.GetNodeMetrics(1, 0, time.Now().UnixMilli()+1000)
|
||||
metrics, err := r.GetNodeMetrics(1, time.Now().UnixMilli()-60000, time.Now().UnixMilli()+1000)
|
||||
if err != nil {
|
||||
t.Fatalf("get metrics: %v", err)
|
||||
}
|
||||
@@ -86,7 +86,7 @@ func TestRecordNodeMetricAutoFlush(t *testing.T) {
|
||||
|
||||
time.Sleep(100 * time.Millisecond)
|
||||
|
||||
metrics, err := r.GetNodeMetrics(1, 0, time.Now().UnixMilli()+1000)
|
||||
metrics, err := r.GetNodeMetrics(1, time.Now().UnixMilli()-60000, time.Now().UnixMilli()+1000)
|
||||
if err != nil {
|
||||
t.Fatalf("get metrics: %v", err)
|
||||
}
|
||||
@@ -123,7 +123,7 @@ func TestIngestionServiceStart(t *testing.T) {
|
||||
|
||||
<-ctx.Done()
|
||||
|
||||
metrics, err := r.GetNodeMetrics(1, 0, time.Now().UnixMilli()+1000)
|
||||
metrics, err := r.GetNodeMetrics(1, time.Now().UnixMilli()-60000, time.Now().UnixMilli()+1000)
|
||||
if err != nil {
|
||||
t.Fatalf("get metrics: %v", err)
|
||||
}
|
||||
@@ -198,7 +198,7 @@ func TestGetMetricsWithTimeRange(t *testing.T) {
|
||||
|
||||
svc.flushNodeMetrics()
|
||||
|
||||
metrics, err := svc.GetMetrics(1, 0, now+1000)
|
||||
metrics, err := svc.GetMetrics(1, now-60000, now+1000)
|
||||
if err != nil {
|
||||
t.Fatalf("get metrics: %v", err)
|
||||
}
|
||||
@@ -224,7 +224,7 @@ func TestPruneMetrics(t *testing.T) {
|
||||
|
||||
svc.pruneMetrics()
|
||||
|
||||
metrics, err := r.GetNodeMetrics(1, 0, time.Now().UnixMilli()+1000)
|
||||
metrics, err := r.GetNodeMetrics(1, time.Now().UnixMilli()-60000, time.Now().UnixMilli()+1000)
|
||||
if err != nil {
|
||||
t.Fatalf("get metrics: %v", err)
|
||||
}
|
||||
@@ -255,7 +255,7 @@ func TestMultipleNodes(t *testing.T) {
|
||||
svc.flushNodeMetrics()
|
||||
|
||||
for nodeID := int64(1); nodeID <= 3; nodeID++ {
|
||||
metrics, err := r.GetNodeMetrics(nodeID, 0, time.Now().UnixMilli()+1000)
|
||||
metrics, err := r.GetNodeMetrics(nodeID, time.Now().UnixMilli()-60000, time.Now().UnixMilli()+1000)
|
||||
if err != nil {
|
||||
t.Fatalf("get metrics for node %d: %v", nodeID, err)
|
||||
}
|
||||
@@ -279,7 +279,7 @@ func TestZeroValues(t *testing.T) {
|
||||
svc.RecordNodeMetric(1, info)
|
||||
svc.flushNodeMetrics()
|
||||
|
||||
metrics, err := r.GetNodeMetrics(1, 0, time.Now().UnixMilli()+1000)
|
||||
metrics, err := r.GetNodeMetrics(1, time.Now().UnixMilli()-60000, time.Now().UnixMilli()+1000)
|
||||
if err != nil {
|
||||
t.Fatalf("get metrics: %v", err)
|
||||
}
|
||||
|
||||
@@ -3320,15 +3320,61 @@ func (r *Repository) GetNodeMetrics(nodeID int64, startMs, endMs int64) ([]model
|
||||
if r == nil || r.db == nil {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
rangeMs := endMs - startMs
|
||||
const maxRawRangeMs = int64(60 * 60 * 1000) // 1 hour — return raw data for short ranges
|
||||
const targetPoints = 500 // target number of chart points for downsampled data
|
||||
|
||||
// For short ranges, return raw data (full resolution).
|
||||
if rangeMs <= maxRawRangeMs {
|
||||
var metrics []model.NodeMetric
|
||||
err := r.db.Where("node_id = ? AND timestamp >= ? AND timestamp <= ?", nodeID, startMs, endMs).
|
||||
Order("timestamp ASC").
|
||||
Limit(5000).
|
||||
Find(&metrics).Error
|
||||
return metrics, err
|
||||
}
|
||||
|
||||
// For longer ranges, downsample via SQL aggregation to keep the response small and fast.
|
||||
bucketMs := rangeMs / targetPoints
|
||||
if bucketMs < 1000 {
|
||||
bucketMs = 1000 // minimum 1-second buckets
|
||||
}
|
||||
|
||||
bucketExpr := fmt.Sprintf("(timestamp / %d * %d)", bucketMs, bucketMs)
|
||||
groupExpr := fmt.Sprintf("timestamp / %d", bucketMs)
|
||||
|
||||
var metrics []model.NodeMetric
|
||||
err := r.db.Where("node_id = ? AND timestamp >= ? AND timestamp <= ?", nodeID, startMs, endMs).
|
||||
Order("timestamp DESC").
|
||||
Limit(5000).
|
||||
Find(&metrics).Error
|
||||
if len(metrics) > 1 {
|
||||
for i, j := 0, len(metrics)-1; i < j; i, j = i+1, j-1 {
|
||||
metrics[i], metrics[j] = metrics[j], metrics[i]
|
||||
}
|
||||
err := r.db.Model(&model.NodeMetric{}).
|
||||
Select(
|
||||
fmt.Sprintf(
|
||||
"? AS node_id, "+
|
||||
"CAST(%s AS INTEGER) AS timestamp, "+
|
||||
"AVG(cpu_usage) AS cpu_usage, "+
|
||||
"AVG(mem_usage) AS mem_usage, "+
|
||||
"AVG(disk_usage) AS disk_usage, "+
|
||||
"CAST(AVG(net_in_bytes) AS INTEGER) AS net_in_bytes, "+
|
||||
"CAST(AVG(net_out_bytes) AS INTEGER) AS net_out_bytes, "+
|
||||
"CAST(AVG(net_in_speed) AS INTEGER) AS net_in_speed, "+
|
||||
"CAST(AVG(net_out_speed) AS INTEGER) AS net_out_speed, "+
|
||||
"AVG(load1) AS load1, "+
|
||||
"AVG(load5) AS load5, "+
|
||||
"AVG(load15) AS load15, "+
|
||||
"CAST(AVG(tcp_conns) AS INTEGER) AS tcp_conns, "+
|
||||
"CAST(AVG(udp_conns) AS INTEGER) AS udp_conns, "+
|
||||
"CAST(MAX(uptime) AS INTEGER) AS uptime",
|
||||
bucketExpr,
|
||||
),
|
||||
nodeID,
|
||||
).
|
||||
Where("node_id = ? AND timestamp >= ? AND timestamp <= ?", nodeID, startMs, endMs).
|
||||
Group(groupExpr).
|
||||
Order("timestamp ASC").
|
||||
Limit(targetPoints + 100). // safety margin
|
||||
Scan(&metrics).Error
|
||||
|
||||
if metrics == nil {
|
||||
metrics = make([]model.NodeMetric, 0)
|
||||
}
|
||||
return metrics, err
|
||||
}
|
||||
|
||||
@@ -1175,7 +1175,7 @@ func TestMetricBatchInsert(t *testing.T) {
|
||||
t.Fatalf("batch insert: %v", err)
|
||||
}
|
||||
|
||||
retrieved, err := repo.GetNodeMetrics(1, 0, now+1000)
|
||||
retrieved, err := repo.GetNodeMetrics(1, now-10000, now+1000)
|
||||
if err != nil {
|
||||
t.Fatalf("get metrics: %v", err)
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user