From 160e63558fbcfe3ed7483cf8f5dee51e9a2eb8ee Mon Sep 17 00:00:00 2001
From: ryan
- 独立访客 {formatCompactNumber(traffic.unique_visitors)} · 错误{' '}
+ 窗口UV(估) {formatCompactNumber(traffic.unique_visitors)} · 错误{' '}
{formatCompactNumber(traffic.error_count)} · 估算 QPS{' '}
{traffic.estimated_qps.toFixed(2)}
{trafficSummary - ? `近 60 秒 · UV ${formatMetricCount(trafficSummary.unique_visitor_count)}` + ? `近 60 秒 · 窗口UV ${formatMetricCount(trafficSummary.unique_visitor_count)}` : '暂无窗口流量摘要'}
diff --git a/internal/apps/admin/status/clickhouse.go b/internal/apps/admin/status/clickhouse.go index a2f8ea05..f84c2d5a 100644 --- a/internal/apps/admin/status/clickhouse.go +++ b/internal/apps/admin/status/clickhouse.go @@ -6,16 +6,19 @@ package status import ( "net/http" + "github.com/Rain-kl/Wavelet/internal/apps/openflare/chwriter" + "github.com/Rain-kl/Wavelet/internal/apps/risk_control" "github.com/Rain-kl/Wavelet/internal/common/response" "github.com/Rain-kl/Wavelet/internal/config" "github.com/Rain-kl/Wavelet/internal/db" + "github.com/Rain-kl/Wavelet/internal/db/batchwriter" analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics" "github.com/gin-gonic/gin" ) // GetClickHouseStatus returns ClickHouse operational metrics for administrators. // @Summary 获取 ClickHouse 运行指标 -// @Description 返回 ClickHouse parts、mutation、async_insert 队列等运维指标,需要管理员权限 +// @Description 返回 ClickHouse parts、mutation、async_insert 队列及进程内 batch writer 指标,需要管理员权限 // @Tags admin // @Produce json // @Security SessionCookie @@ -36,5 +39,15 @@ func GetClickHouseStatus(c *gin.Context) { response.AbortInternal(c, "获取 ClickHouse 运行指标失败") return } + stats.BatchWriters = collectBatchWriterStats() c.JSON(http.StatusOK, response.OK(stats)) -} \ No newline at end of file +} + +func collectBatchWriterStats() []batchwriter.Stats { + out := chwriter.WriterStats() + if out == nil { + out = make([]batchwriter.Stats, 0, 1) + } + out = append(out, risk_control.LogWriterStats()) + return out +} diff --git a/internal/apps/openflare/chwriter/dedup.go b/internal/apps/openflare/chwriter/dedup.go index 3ffcecd2..8490fc9f 100644 --- a/internal/apps/openflare/chwriter/dedup.go +++ b/internal/apps/openflare/chwriter/dedup.go @@ -25,7 +25,7 @@ func newDedupSet() *dedupSet { // markIfNew records key when it has not been seen within dedupTTL. func (s *dedupSet) markIfNew(key string) bool { - if key == "" { + if s == nil || key == "" { return false } @@ -33,19 +33,33 @@ func (s *dedupSet) markIfNew(key string) bool { s.mu.Lock() defer s.mu.Unlock() - // Periodically clean up all expired keys (e.g., every 30 seconds) - if now.Sub(s.lastCleanup) >= 30*time.Second { - for existing, expiresAt := range s.keys { - if now.After(expiresAt) { - delete(s.keys, existing) - } - } - s.lastCleanup = now - } + s.cleanupExpiredLocked(now) if expiresAt, exists := s.keys[key]; exists && now.Before(expiresAt) { return false } s.keys[key] = now.Add(dedupTTL) return true -} \ No newline at end of file +} + +// unmark removes a key so a later enqueue or flush retry may accept it again. +func (s *dedupSet) unmark(key string) { + if s == nil || key == "" { + return + } + s.mu.Lock() + defer s.mu.Unlock() + delete(s.keys, key) +} + +func (s *dedupSet) cleanupExpiredLocked(now time.Time) { + if now.Sub(s.lastCleanup) < 30*time.Second { + return + } + for existing, expiresAt := range s.keys { + if now.After(expiresAt) { + delete(s.keys, existing) + } + } + s.lastCleanup = now +} diff --git a/internal/apps/openflare/chwriter/dedup_test.go b/internal/apps/openflare/chwriter/dedup_test.go index e36542e0..a89187de 100644 --- a/internal/apps/openflare/chwriter/dedup_test.go +++ b/internal/apps/openflare/chwriter/dedup_test.go @@ -3,7 +3,16 @@ package chwriter -import "testing" +import ( + "context" + "errors" + "sync" + "testing" + "time" + + "github.com/Rain-kl/Wavelet/internal/db/batchwriter" + analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" +) func TestDedupSetMarkIfNew(t *testing.T) { t.Parallel() @@ -21,4 +30,157 @@ func TestDedupSetMarkIfNew(t *testing.T) { if set.markIfNew("") { t.Fatal("markIfNew() = true, want false on empty key") } -} \ No newline at end of file +} + +func TestDedupSetUnmarkAllowsRetry(t *testing.T) { + t.Parallel() + + set := newDedupSet() + if !set.markIfNew("k") { + t.Fatal("markIfNew() = false, want true") + } + set.unmark("k") + if !set.markIfNew("k") { + t.Fatal("markIfNew() after unmark = false, want true") + } +} + +func TestQueueWithDedupDoesNotMarkWhenEnqueueFails(t *testing.T) { + t.Parallel() + + cfg := batchwriter.DefaultConfig() + cfg.QueueSize = 1 + cfg.MaxBatchSize = 10 + cfg.FlushInterval = time.Hour + + // Block the worker so the queue stays full after one enqueue. + block := make(chan struct{}) + writer, err := batchwriter.New[int](cfg, func(context.Context, []int) error { + <-block + return nil + }) + if err != nil { + t.Fatalf("New() error = %v", err) + } + writer.Start(context.Background()) + t.Cleanup(func() { + close(block) + stopCtx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + _ = writer.Stop(stopCtx) + }) + + // Fill the channel buffer (and the worker's current receive slot may empty one). + // Keep enqueueing until full so subsequent queueWithDedup fails. + for i := 0; i < cfg.QueueSize+2; i++ { + _ = writer.TryEnqueue(i) + if writer.IsFull() { + break + } + } + if !writer.IsFull() { + t.Fatal("writer not full after filling; cannot test enqueue failure path") + } + + dedup := newDedupSet() + queueWithDedup(writer, dedup, "dedup-key", 99) + // Key must not remain marked after failed enqueue. + if !dedup.markIfNew("dedup-key") { + t.Fatal("dedup key still marked after failed enqueue; want unmark") + } +} + +func TestQueueWithDedupMarksOnlyOnSuccess(t *testing.T) { + t.Parallel() + + cfg := batchwriter.DefaultConfig() + cfg.MaxBatchSize = 100 + cfg.FlushInterval = time.Hour + + writer, err := batchwriter.New[int](cfg, func(context.Context, []int) error { return nil }) + if err != nil { + t.Fatalf("New() error = %v", err) + } + writer.Start(context.Background()) + t.Cleanup(func() { + stopCtx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + _ = writer.Stop(stopCtx) + }) + + dedup := newDedupSet() + queueWithDedup(writer, dedup, "ok-key", 1) + if dedup.markIfNew("ok-key") { + t.Fatal("markIfNew() = true after successful enqueue, want false (key marked)") + } +} + +func TestFlushErrorHandlerUnmarksKeys(t *testing.T) { + t.Parallel() + + dedup := newDedupSet() + flushErr := errors.New("ch down") + + var ( + mu sync.Mutex + errCount int + ) + + cfg := batchwriter.Config{ + Name: "test_obs", + QueueSize: 10, + MaxBatchSize: 1, + FlushInterval: time.Hour, + } + keyFn := func(s analyticsmodel.NodeMetricSnapshot) string { + return metricSnapshotKey(s) + } + writer, err := batchwriter.New( + cfg, + func(context.Context, []analyticsmodel.NodeMetricSnapshot) error { return flushErr }, + batchwriter.WithFlushErrorHandler[analyticsmodel.NodeMetricSnapshot](func(_ context.Context, items []analyticsmodel.NodeMetricSnapshot, err error) { + mu.Lock() + errCount++ + mu.Unlock() + for _, item := range items { + dedup.unmark(keyFn(item)) + } + }), + ) + if err != nil { + t.Fatalf("New() error = %v", err) + } + writer.Start(context.Background()) + t.Cleanup(func() { + stopCtx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + _ = writer.Stop(stopCtx) + }) + + item := analyticsmodel.NodeMetricSnapshot{ + NodeID: "n1", + CapturedAt: time.Unix(1, 0).UTC(), + } + key := keyFn(item) + if !dedup.markIfNew(key) { + t.Fatal("markIfNew failed") + } + if !writer.TryEnqueue(item) { + t.Fatal("TryEnqueue failed") + } + + deadline := time.Now().Add(time.Second) + for { + mu.Lock() + ready := errCount >= 1 + mu.Unlock() + if ready || time.Now().After(deadline) { + break + } + time.Sleep(5 * time.Millisecond) + } + + if !dedup.markIfNew(key) { + t.Fatal("key still marked after flush error unmark; want available for retry") + } +} diff --git a/internal/apps/openflare/chwriter/writer.go b/internal/apps/openflare/chwriter/writer.go index 7e5fd60b..d9f5634c 100644 --- a/internal/apps/openflare/chwriter/writer.go +++ b/internal/apps/openflare/chwriter/writer.go @@ -14,6 +14,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/config" "github.com/Rain-kl/Wavelet/internal/db/batchwriter" "github.com/Rain-kl/Wavelet/internal/lifecycle" + "github.com/Rain-kl/Wavelet/internal/model" analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics" "github.com/Rain-kl/Wavelet/pkg/logger" @@ -33,6 +34,10 @@ const ( nodeAccessLogMinBatchSize = 50 nodeAccessLogFlushEvery = 2 * time.Second nodeAccessLogMaxFlushWait = 5 * time.Second + + // flushAttempts is total tries (1 initial + short retries) before giving up a batch. + flushAttempts = 2 + flushRetryBackoff = 50 * time.Millisecond ) var ( @@ -65,11 +70,36 @@ func Init(ctx context.Context) { frpsDedup = newDedupSet() frpcDedup = newDedupSet() - metricSnapshotWriter = mustNewObservabilityWriter("metric_snapshots", analyticsrepo.BatchInsertNodeMetricSnapshots) - requestReportWriter = mustNewObservabilityWriter("request_reports", analyticsrepo.BatchInsertNodeRequestReports) - openrestyWriter = mustNewObservabilityWriter("openresty_obs", analyticsrepo.BatchInsertNodeObsOpenresty) - frpsWriter = mustNewObservabilityWriter("frps_obs", analyticsrepo.BatchInsertNodeObsFrps) - frpcWriter = mustNewObservabilityWriter("frpc_obs", analyticsrepo.BatchInsertNodeObsFrpc) + metricSnapshotWriter = mustNewObservabilityWriter( + "metric_snapshots", + withFlushRetries(analyticsrepo.BatchInsertNodeMetricSnapshots), + metricSnapshotDedup, + metricSnapshotKey, + ) + requestReportWriter = mustNewObservabilityWriter( + "request_reports", + withFlushRetries(analyticsrepo.BatchInsertNodeRequestReports), + requestReportDedup, + requestReportKey, + ) + openrestyWriter = mustNewObservabilityWriter( + "openresty_obs", + withFlushRetries(analyticsrepo.BatchInsertNodeObsOpenresty), + openrestyDedup, + openrestyKey, + ) + frpsWriter = mustNewObservabilityWriter( + "frps_obs", + withFlushRetries(analyticsrepo.BatchInsertNodeObsFrps), + frpsDedup, + frpsKey, + ) + frpcWriter = mustNewObservabilityWriter( + "frpc_obs", + withFlushRetries(analyticsrepo.BatchInsertNodeObsFrpc), + frpcDedup, + frpcKey, + ) nodeAccessLogWriter = mustNewNodeAccessLogWriter() metricSnapshotWriter.Start(ctx) @@ -79,6 +109,7 @@ func Init(ctx context.Context) { frpcWriter.Start(ctx) nodeAccessLogWriter.Start(ctx) + wireModelInsertHooks() lifecycle.OnShutdown("openflare_chwriter", Stop) }) } @@ -108,69 +139,49 @@ func Stop(ctx context.Context) error { return firstErr } +// WriterStats returns queue depth and failure counters for all OpenFlare writers. +func WriterStats() []batchwriter.Stats { + writers := []statsProvider{ + metricSnapshotWriter, + requestReportWriter, + openrestyWriter, + frpsWriter, + frpcWriter, + nodeAccessLogWriter, + } + out := make([]batchwriter.Stats, 0, len(writers)) + for _, w := range writers { + if w == nil { + continue + } + out = append(out, w.Stats()) + } + return out +} + // QueueMetricSnapshot enqueues a metric snapshot for asynchronous flush. func QueueMetricSnapshot(snapshot analyticsmodel.NodeMetricSnapshot) { - if metricSnapshotWriter == nil { - return - } - key := fmt.Sprintf("%s|%d", snapshot.NodeID, snapshot.CapturedAt.UTC().UnixNano()) - if !metricSnapshotDedup.markIfNew(key) { - return - } - metricSnapshotWriter.TryEnqueue(snapshot) + queueWithDedup(metricSnapshotWriter, metricSnapshotDedup, metricSnapshotKey(snapshot), snapshot) } // QueueRequestReport enqueues a request report for asynchronous flush. func QueueRequestReport(report analyticsmodel.NodeRequestReport) { - if requestReportWriter == nil { - return - } - key := fmt.Sprintf( - "%s|%d|%d", - report.NodeID, - report.WindowStartedAt.UTC().UnixNano(), - report.WindowEndedAt.UTC().UnixNano(), - ) - if !requestReportDedup.markIfNew(key) { - return - } - requestReportWriter.TryEnqueue(report) + queueWithDedup(requestReportWriter, requestReportDedup, requestReportKey(report), report) } // QueueOpenrestyObservation enqueues an OpenResty observation for asynchronous flush. func QueueOpenrestyObservation(observation analyticsmodel.NodeObsOpenresty) { - if openrestyWriter == nil { - return - } - key := fmt.Sprintf("%s|%d", observation.NodeID, observation.CapturedAt.UTC().UnixNano()) - if !openrestyDedup.markIfNew(key) { - return - } - openrestyWriter.TryEnqueue(observation) + queueWithDedup(openrestyWriter, openrestyDedup, openrestyKey(observation), observation) } // QueueFrpsObservation enqueues an FRPS observation for asynchronous flush. func QueueFrpsObservation(observation analyticsmodel.NodeObsFrps) { - if frpsWriter == nil { - return - } - key := fmt.Sprintf("%s|%d", observation.NodeID, observation.CapturedAt.UTC().UnixNano()) - if !frpsDedup.markIfNew(key) { - return - } - frpsWriter.TryEnqueue(observation) + queueWithDedup(frpsWriter, frpsDedup, frpsKey(observation), observation) } // QueueFrpcObservation enqueues an FRPC observation for asynchronous flush. func QueueFrpcObservation(observation analyticsmodel.NodeObsFrpc) { - if frpcWriter == nil { - return - } - key := fmt.Sprintf("%s|%d", observation.NodeID, observation.CapturedAt.UTC().UnixNano()) - if !frpcDedup.markIfNew(key) { - return - } - frpcWriter.TryEnqueue(observation) + queueWithDedup(frpcWriter, frpcDedup, frpcKey(observation), observation) } // QueueNodeAccessLogs enqueues node access logs for asynchronous flush. @@ -183,7 +194,26 @@ func QueueNodeAccessLogs(logs []analyticsmodel.NodeAccessLog) { } } -func mustNewObservabilityWriter[T any](name string, flush batchwriter.FlushFunc[T]) *batchwriter.Writer[T] { +func queueWithDedup[T any](writer *batchwriter.Writer[T], dedup *dedupSet, key string, item T) { + if writer == nil { + return + } + // Mark first so concurrent duplicates still collapse; release on enqueue failure + // so a full queue does not permanently suppress the item. + if !dedup.markIfNew(key) { + return + } + if !writer.TryEnqueue(item) { + dedup.unmark(key) + } +} + +func mustNewObservabilityWriter[T any]( + name string, + flush batchwriter.FlushFunc[T], + dedup *dedupSet, + keyFn func(T) string, +) *batchwriter.Writer[T] { cfg := batchwriter.Config{ Name: name, QueueSize: observabilityQueueSize, @@ -196,8 +226,14 @@ func mustNewObservabilityWriter[T any](name string, flush batchwriter.FlushFunc[ cfg, flush, withObservabilityDropHandler[T](name), - batchwriter.WithFlushErrorHandler[T](func(ctx context.Context, batchSize int, err error) { - logger.ErrorF(ctx, "[OpenFlare] flush %s failed (batch=%d): %v", name, batchSize, err) + batchwriter.WithFlushErrorHandler[T](func(ctx context.Context, items []T, err error) { + logger.ErrorF(ctx, "[OpenFlare] flush %s failed (batch=%d): %v", name, len(items), err) + if dedup == nil || keyFn == nil { + return + } + for _, item := range items { + dedup.unmark(keyFn(item)) + } }), ) if err != nil { @@ -215,12 +251,14 @@ func mustNewNodeAccessLogWriter() *batchwriter.Writer[analyticsmodel.NodeAccessL FlushInterval: nodeAccessLogFlushEvery, MaxFlushWait: nodeAccessLogMaxFlushWait, } - writer, err := batchwriter.New[analyticsmodel.NodeAccessLog](cfg, analyticsrepo.BatchInsertNodeAccessLogs, + writer, err := batchwriter.New[analyticsmodel.NodeAccessLog]( + cfg, + withFlushRetries(analyticsrepo.BatchInsertNodeAccessLogs), batchwriter.WithDropHandler[analyticsmodel.NodeAccessLog](func(item analyticsmodel.NodeAccessLog) { logger.WarnF(context.Background(), "[OpenFlare] node access log queue full, dropping log for node %s path %s", item.NodeID, item.Path) }), - batchwriter.WithFlushErrorHandler[analyticsmodel.NodeAccessLog](func(ctx context.Context, batchSize int, err error) { - logger.ErrorF(ctx, "[OpenFlare] flush node access logs failed (batch=%d): %v", batchSize, err) + batchwriter.WithFlushErrorHandler[analyticsmodel.NodeAccessLog](func(ctx context.Context, items []analyticsmodel.NodeAccessLog, err error) { + logger.ErrorF(ctx, "[OpenFlare] flush node access logs failed (batch=%d): %v", len(items), err) }), ) if err != nil { @@ -235,10 +273,74 @@ func withObservabilityDropHandler[T any](name string) batchwriter.Option[T] { }) } +// withFlushRetries wraps a flush function with a short retry to ride out brief CH blips. +func withFlushRetries[T any](flush batchwriter.FlushFunc[T]) batchwriter.FlushFunc[T] { + return func(ctx context.Context, items []T) error { + var err error + for attempt := 1; attempt <= flushAttempts; attempt++ { + err = flush(ctx, items) + if err == nil { + return nil + } + if attempt == flushAttempts { + break + } + select { + case <-ctx.Done(): + return ctx.Err() + case <-time.After(flushRetryBackoff * time.Duration(attempt)): + } + } + return err + } +} + +func wireModelInsertHooks() { + model.SetObservabilityInsertHooks(model.ObservabilityInsertHooks{ + QueueMetricSnapshot: QueueMetricSnapshot, + QueueRequestReport: QueueRequestReport, + QueueOpenrestyObservation: QueueOpenrestyObservation, + QueueFrpsObservation: QueueFrpsObservation, + QueueFrpcObservation: QueueFrpcObservation, + }) + model.SetAccessLogInsertHooks(model.AccessLogInsertHooks{ + QueueNodeAccessLogs: QueueNodeAccessLogs, + }) +} + +func metricSnapshotKey(snapshot analyticsmodel.NodeMetricSnapshot) string { + return fmt.Sprintf("%s|%d", snapshot.NodeID, snapshot.CapturedAt.UTC().UnixNano()) +} + +func requestReportKey(report analyticsmodel.NodeRequestReport) string { + return fmt.Sprintf( + "%s|%d|%d", + report.NodeID, + report.WindowStartedAt.UTC().UnixNano(), + report.WindowEndedAt.UTC().UnixNano(), + ) +} + +func openrestyKey(observation analyticsmodel.NodeObsOpenresty) string { + return fmt.Sprintf("%s|%d", observation.NodeID, observation.CapturedAt.UTC().UnixNano()) +} + +func frpsKey(observation analyticsmodel.NodeObsFrps) string { + return fmt.Sprintf("%s|%d", observation.NodeID, observation.CapturedAt.UTC().UnixNano()) +} + +func frpcKey(observation analyticsmodel.NodeObsFrpc) string { + return fmt.Sprintf("%s|%d", observation.NodeID, observation.CapturedAt.UTC().UnixNano()) +} + type batchStopper interface { Stop(ctx context.Context) error } +type statsProvider interface { + Stats() batchwriter.Stats +} + func running() bool { return metricSnapshotWriter != nil && metricSnapshotWriter.Running() -} \ No newline at end of file +} diff --git a/internal/apps/openflare/dashboard/logics.go b/internal/apps/openflare/dashboard/logics.go index 2c1a97d7..a0fc8461 100644 --- a/internal/apps/openflare/dashboard/logics.go +++ b/internal/apps/openflare/dashboard/logics.go @@ -117,6 +117,16 @@ func buildOverviewView(ctx context.Context) (*OverviewView, error) { if err != nil { return nil, err } + // Latest-per-node health: dedicated LIMIT 1 BY queries (not a global raw LIMIT). + latestSnapshotRows, err := model.ListOpenFlareLatestMetricSnapshotsSince(ctx, "", since) + if err != nil { + return nil, err + } + latestTrafficRows, err := model.ListOpenFlareLatestRequestReportsSince(ctx, "", since) + if err != nil { + return nil, err + } + // Bounded raw windows remain for distributions and trend fallbacks; trends prefer hourly rollups. snapshots, err := model.ListOpenFlareMetricSnapshotsSince(ctx, "", since, dashboardOverviewSnapshotLimit) if err != nil { return nil, err @@ -146,8 +156,8 @@ func buildOverviewView(ctx context.Context) (*OverviewView, error) { var cpuNodeCount int var memoryNodeCount int - latestSnapshots := observability.LatestMetricSnapshotsByNode(snapshots) - latestTrafficReports := observability.LatestTrafficReportsByNode(reports) + latestSnapshots := observability.LatestMetricSnapshotsByNode(latestSnapshotRows) + latestTrafficReports := observability.LatestTrafficReportsByNode(latestTrafficRows) activeEventsByNode := observability.ActiveHealthEventsByNode(activeEvents) for _, node := range nodes { diff --git a/internal/apps/openflare/dashboard/logics_test.go b/internal/apps/openflare/dashboard/logics_test.go index be27379f..57c7c4b2 100644 --- a/internal/apps/openflare/dashboard/logics_test.go +++ b/internal/apps/openflare/dashboard/logics_test.go @@ -58,6 +58,32 @@ func TestGetOverviewStructure(t *testing.T) { OpenrestyStatus: "unknown", }).Error) + // Seed older + newer snapshots per node; health must use latest-per-node, not a global raw limit. + require.NoError(t, model.InsertOpenFlareMetricSnapshot(ctx, &model.OpenFlareMetricSnapshot{ + NodeID: "node-dashboard-1", + CapturedAt: now.Add(-2 * time.Hour), + CPUUsagePercent: 10, + MemoryUsedBytes: 1, + MemoryTotalBytes: 10, + })) + require.NoError(t, model.InsertOpenFlareMetricSnapshot(ctx, &model.OpenFlareMetricSnapshot{ + NodeID: "node-dashboard-1", + CapturedAt: now.Add(-time.Minute), + CPUUsagePercent: 55, + MemoryUsedBytes: 5, + MemoryTotalBytes: 10, + StorageUsedBytes: 2, + StorageTotalBytes: 10, + })) + require.NoError(t, model.InsertOpenFlareRequestReport(ctx, &model.OpenFlareRequestReport{ + NodeID: "node-dashboard-1", + WindowStartedAt: now.Add(-2 * time.Minute), + WindowEndedAt: now.Add(-time.Minute), + RequestCount: 12, + ErrorCount: 1, + UniqueVisitorCount: 4, + })) + overview, err := GetOverview(ctx) require.NoError(t, err) require.NotNil(t, overview) @@ -69,14 +95,14 @@ func TestGetOverviewStructure(t *testing.T) { assert.Equal(t, 0, overview.Summary.OfflineNodes) assert.Equal(t, 0, overview.Summary.UnhealthyNodes) - assert.Equal(t, int64(0), overview.Traffic.RequestCount) - assert.Equal(t, int64(0), overview.Traffic.UniqueVisitors) - assert.Equal(t, int64(0), overview.Traffic.ErrorCount) - assert.Equal(t, float64(0), overview.Traffic.EstimatedQPS) - assert.Equal(t, 0, overview.Traffic.ReportedNodes) + assert.Equal(t, int64(12), overview.Traffic.RequestCount) + assert.Equal(t, int64(4), overview.Traffic.UniqueVisitors) + assert.Equal(t, int64(1), overview.Traffic.ErrorCount) + assert.InDelta(t, 0.2, overview.Traffic.EstimatedQPS, 0.0001) + assert.Equal(t, 1, overview.Traffic.ReportedNodes) - assert.Equal(t, float64(0), overview.Capacity.AverageCPUUsagePercent) - assert.Equal(t, float64(0), overview.Capacity.AverageMemoryUsagePercent) + assert.Equal(t, 55.0, overview.Capacity.AverageCPUUsagePercent) + assert.Equal(t, 50.0, overview.Capacity.AverageMemoryUsagePercent) assert.Equal(t, 0, overview.Capacity.HighCPUNodes) assert.Equal(t, 0, overview.Capacity.HighMemoryNodes) assert.Equal(t, 0, overview.Capacity.HighStorageNodes) @@ -120,10 +146,20 @@ func TestGetOverviewStructure(t *testing.T) { assert.Equal(t, "Edge 1", onlineNode[2]) assert.Equal(t, "online", onlineNode[6]) assert.Equal(t, "healthy", onlineNode[7]) + // Latest-per-node health fields (indexes match compressDashboardNodes). + assert.Equal(t, 55.0, onlineNode[11]) // cpu_usage_percent from latest snapshot + assert.Equal(t, 50.0, onlineNode[12]) // memory_usage_percent + assert.Equal(t, int64(12), onlineNode[14]) + assert.Equal(t, int64(1), onlineNode[15]) + assert.Equal(t, int64(4), onlineNode[16]) pendingNode := nodeByID["node-dashboard-2"] require.NotNil(t, pendingNode) assert.Equal(t, "Edge 2", pendingNode[2]) assert.Equal(t, "pending", pendingNode[6]) assert.Equal(t, "unknown", pendingNode[7]) + + assert.Equal(t, 55.0, overview.Capacity.AverageCPUUsagePercent) + assert.Equal(t, 1, overview.Traffic.ReportedNodes) + assert.Equal(t, int64(4), overview.Traffic.UniqueVisitors) } diff --git a/internal/apps/openflare/option/logics.go b/internal/apps/openflare/option/logics.go index 0edd6127..3693c29b 100644 --- a/internal/apps/openflare/option/logics.go +++ b/internal/apps/openflare/option/logics.go @@ -59,6 +59,9 @@ type databaseCleanupResult struct { Target string `json:"target"` TargetLabel string `json:"target_label"` DeletedCount int64 `json:"deleted_count"` + EligibleCount int64 `json:"eligible_count,omitempty"` + CleanupMode string `json:"cleanup_mode,omitempty"` + TableTTLDays int `json:"table_ttl_days,omitempty"` DeleteAll bool `json:"delete_all"` RetentionDays *int `json:"retention_days,omitempty"` } @@ -198,6 +201,9 @@ func cleanupDatabaseObservability(ctx context.Context, input databaseCleanupInpu Target: result.Target, TargetLabel: result.TargetLabel, DeletedCount: result.DeletedCount, + EligibleCount: result.EligibleCount, + CleanupMode: result.CleanupMode, + TableTTLDays: result.TableTTLDays, DeleteAll: result.DeleteAll, RetentionDays: result.RetentionDays, }, nil diff --git a/internal/apps/openflare/option/logics_test.go b/internal/apps/openflare/option/logics_test.go index c89dd45f..a6b224ca 100644 --- a/internal/apps/openflare/option/logics_test.go +++ b/internal/apps/openflare/option/logics_test.go @@ -146,21 +146,26 @@ func TestCleanupDatabaseObservabilityDeletesRows(t *testing.T) { }, })) - retention := 7 - result, err := cleanupDatabaseObservability(ctx, databaseCleanupInput{ + // Retention shorter than table TTL (90d for access logs) must be rejected. + shortRetention := 7 + _, err := cleanupDatabaseObservability(ctx, databaseCleanupInput{ Target: "node_access_logs", - RetentionDays: &retention, + RetentionDays: &shortRetention, + }) + require.Error(t, err) + + // Full truncate still hard-deletes all rows. + result, err := cleanupDatabaseObservability(ctx, databaseCleanupInput{ + Target: "node_access_logs", }) require.NoError(t, err) assert.Equal(t, "node_access_logs", result.Target) assert.Equal(t, "访问日志", result.TargetLabel) - assert.Equal(t, int64(1), result.DeletedCount) - assert.False(t, result.DeleteAll) - require.NotNil(t, result.RetentionDays) - assert.Equal(t, 7, *result.RetentionDays) + assert.Equal(t, int64(2), result.DeletedCount) + assert.True(t, result.DeleteAll) + assert.Equal(t, "truncate", result.CleanupMode) rows, err := model.ListOpenFlareAccessLogs(ctx, model.OpenFlareAccessLogQuery{Page: 0, PageSize: 10}) require.NoError(t, err) - require.Len(t, rows, 1) - assert.Equal(t, "/recent", rows[0].Path) + assert.Empty(t, rows) } diff --git a/internal/apps/openflare/tasks/database_cleanup.go b/internal/apps/openflare/tasks/database_cleanup.go index 6bbae7c0..55901dd0 100644 --- a/internal/apps/openflare/tasks/database_cleanup.go +++ b/internal/apps/openflare/tasks/database_cleanup.go @@ -34,9 +34,19 @@ var databaseCleanupTargets = map[string]string{ DatabaseCleanupTargetAccessLogs: "访问日志", DatabaseCleanupTargetMetricSnapshots: "性能快照", DatabaseCleanupTargetRequestReports: "请求聚合", - DatabaseCleanupTargetObsOpenresty: "OpenResty 观测", - DatabaseCleanupTargetObsFrps: "FRPS 观测", - DatabaseCleanupTargetObsFrpc: "FRPC 观测", + DatabaseCleanupTargetObsOpenresty: "OpenResty 观测", + DatabaseCleanupTargetObsFrps: "FRPS 观测", + DatabaseCleanupTargetObsFrpc: "FRPC 观测", +} + +// databaseCleanupTableTTLDays maps API targets to ClickHouse DDL TTL days. +var databaseCleanupTableTTLDays = map[string]int{ + DatabaseCleanupTargetAccessLogs: analyticsrepo.TableTTLDaysNodeAccessLogs, + DatabaseCleanupTargetMetricSnapshots: analyticsrepo.TableTTLDaysNodeMetricSnapshots, + DatabaseCleanupTargetRequestReports: analyticsrepo.TableTTLDaysNodeRequestReports, + DatabaseCleanupTargetObsOpenresty: analyticsrepo.TableTTLDaysNodeObs, + DatabaseCleanupTargetObsFrps: analyticsrepo.TableTTLDaysNodeObs, + DatabaseCleanupTargetObsFrpc: analyticsrepo.TableTTLDaysNodeObs, } // DatabaseCleanupInput describes a manual observability cleanup request. @@ -46,14 +56,21 @@ type DatabaseCleanupInput struct { } // DatabaseCleanupResult summarizes a manual observability cleanup run. +// +// Semantics: +// - delete_all / cleanup_mode=truncate: DeletedCount is hard-deleted rows (TRUNCATE). +// - retention path / cleanup_mode=ttl_materialize: DeletedCount is always 0; +// EligibleCount estimates rows past the table DDL TTL (not an arbitrary younger cutoff). type DatabaseCleanupResult struct { - Target string `json:"target"` - TargetLabel string `json:"target_label"` - DeletedCount int64 `json:"deleted_count"` - CleanupMode string `json:"cleanup_mode,omitempty"` - DeleteAll bool `json:"delete_all"` - RetentionDays *int `json:"retention_days,omitempty"` - Cutoff *time.Time `json:"cutoff,omitempty"` + Target string `json:"target"` + TargetLabel string `json:"target_label"` + DeletedCount int64 `json:"deleted_count"` + EligibleCount int64 `json:"eligible_count,omitempty"` + CleanupMode string `json:"cleanup_mode,omitempty"` + TableTTLDays int `json:"table_ttl_days,omitempty"` + DeleteAll bool `json:"delete_all"` + RetentionDays *int `json:"retention_days,omitempty"` + Cutoff *time.Time `json:"cutoff,omitempty"` } // DatabaseAutoCleanupSummary summarizes a scheduled auto-cleanup run. @@ -63,7 +80,17 @@ type DatabaseAutoCleanupSummary struct { Results []DatabaseCleanupResult `json:"results"` } +// TableTTLDaysForCleanupTarget returns the DDL TTL days for a cleanup target. +func TableTTLDaysForCleanupTarget(target string) (int, bool) { + days, ok := databaseCleanupTableTTLDays[strings.TrimSpace(target)] + return days, ok +} + // CleanupDatabaseObservability deletes observability rows for the given target. +// +// When RetentionDays is nil, rows are hard-deleted via TRUNCATE. +// When RetentionDays is set, ClickHouse only force-materializes the table TTL policy: +// retention_days shorter than the table TTL is rejected (do not fake success). func CleanupDatabaseObservability(ctx context.Context, input DatabaseCleanupInput) (*DatabaseCleanupResult, error) { target := strings.TrimSpace(input.Target) targetLabel, ok := databaseCleanupTargets[target] @@ -74,10 +101,12 @@ func CleanupDatabaseObservability(ctx context.Context, input DatabaseCleanupInpu return nil, errors.New("retention_days 必须为大于 0 的整数") } + tableTTLDays := databaseCleanupTableTTLDays[target] result := &DatabaseCleanupResult{ - Target: target, - TargetLabel: targetLabel, - DeleteAll: input.RetentionDays == nil, + Target: target, + TargetLabel: targetLabel, + DeleteAll: input.RetentionDays == nil, + TableTTLDays: tableTTLDays, } if input.RetentionDays == nil { @@ -86,24 +115,37 @@ func CleanupDatabaseObservability(ctx context.Context, input DatabaseCleanupInpu return nil, err } result.DeletedCount = deleted + result.EligibleCount = deleted result.CleanupMode = mode return result, nil } retentionDays := *input.RetentionDays - cutoff := time.Now().UTC().Add(-time.Duration(retentionDays) * 24 * time.Hour) - deleted, mode, err := deleteObservabilityRowsBefore(ctx, target, cutoff) + if retentionDays < tableTTLDays { + return nil, fmt.Errorf( + "retention_days 不能小于表 TTL(%d 天);ClickHouse 仅支持按表 TTL 物化过期,更短保留请使用清空全部或调整 DDL", + tableTTLDays, + ) + } + + // MATERIALIZE TTL only enforces DDL policy; cutoff reported is the table TTL boundary. + tableCutoff := time.Now().UTC().Add(-time.Duration(tableTTLDays) * 24 * time.Hour) + eligible, mode, err := materializeObservabilityTableTTL(ctx, target) if err != nil { return nil, err } - result.DeletedCount = deleted + result.DeletedCount = 0 + result.EligibleCount = eligible result.CleanupMode = mode result.RetentionDays = &retentionDays - result.Cutoff = &cutoff + result.Cutoff = &tableCutoff return result, nil } // RunDatabaseAutoCleanupOnce runs retention-based cleanup for all observability targets. +// +// Configured retention shorter than a target's table TTL is clamped up to the table TTL +// so the scheduled job can force-materialize each table policy without failing. func RunDatabaseAutoCleanupOnce(ctx context.Context, now time.Time) (*DatabaseAutoCleanupSummary, error) { enabled, err := repository.GetBoolByKey(ctx, model.ConfigKeyDatabaseAutoCleanupEnabled) if err != nil { @@ -128,9 +170,13 @@ func RunDatabaseAutoCleanupOnce(ctx context.Context, now time.Time) (*DatabaseAu DatabaseCleanupTargetObsFrps, DatabaseCleanupTargetObsFrpc, } { + effectiveDays := retentionDays + if ttl, ok := databaseCleanupTableTTLDays[target]; ok && effectiveDays < ttl { + effectiveDays = ttl + } result, err := CleanupDatabaseObservability(ctx, DatabaseCleanupInput{ Target: target, - RetentionDays: &retentionDays, + RetentionDays: &effectiveDays, }) if err != nil { return nil, err @@ -172,29 +218,37 @@ func deleteAllObservabilityRows(ctx context.Context, target string) (int64, stri return deleted, analyticsrepo.CleanupModeTruncate, nil } -func deleteObservabilityRowsBefore(ctx context.Context, target string, cutoff time.Time) (int64, string, error) { +// materializeObservabilityTableTTL triggers table-TTL materialize (or memory-store delete-before +// with the table TTL cutoff for tests) and returns the eligible/estimate row count. +func materializeObservabilityTableTTL(ctx context.Context, target string) (int64, string, error) { + ttlDays, ok := databaseCleanupTableTTLDays[target] + if !ok { + return 0, "", errors.New("unsupported cleanup target") + } + cutoff := time.Now().UTC().Add(-time.Duration(ttlDays) * 24 * time.Hour) + var ( - deleted int64 - err error + eligible int64 + err error ) switch target { case DatabaseCleanupTargetAccessLogs: - deleted, err = model.DeleteOpenFlareAccessLogsBefore(ctx, cutoff) + eligible, err = model.DeleteOpenFlareAccessLogsBefore(ctx, cutoff) case DatabaseCleanupTargetMetricSnapshots: - deleted, err = model.DeleteOpenFlareMetricSnapshotsBefore(ctx, cutoff) + eligible, err = model.DeleteOpenFlareMetricSnapshotsBefore(ctx, cutoff) case DatabaseCleanupTargetRequestReports: - deleted, err = model.DeleteOpenFlareRequestReportsBefore(ctx, cutoff) + eligible, err = model.DeleteOpenFlareRequestReportsBefore(ctx, cutoff) case DatabaseCleanupTargetObsOpenresty: - deleted, err = model.DeleteOpenFlareNodeObservationOpenrestyBefore(ctx, cutoff) + eligible, err = model.DeleteOpenFlareNodeObservationOpenrestyBefore(ctx, cutoff) case DatabaseCleanupTargetObsFrps: - deleted, err = model.DeleteOpenFlareNodeObservationFrpsBefore(ctx, cutoff) + eligible, err = model.DeleteOpenFlareNodeObservationFrpsBefore(ctx, cutoff) case DatabaseCleanupTargetObsFrpc: - deleted, err = model.DeleteOpenFlareNodeObservationFrpcBefore(ctx, cutoff) + eligible, err = model.DeleteOpenFlareNodeObservationFrpcBefore(ctx, cutoff) default: return 0, "", errors.New("unsupported cleanup target") } if err != nil { return 0, "", err } - return deleted, analyticsrepo.CleanupModeTTLMaterialize, nil + return eligible, analyticsrepo.CleanupModeTTLMaterialize, nil } diff --git a/internal/apps/openflare/tasks/database_cleanup_test.go b/internal/apps/openflare/tasks/database_cleanup_test.go index cfe65daa..c397a8c2 100644 --- a/internal/apps/openflare/tasks/database_cleanup_test.go +++ b/internal/apps/openflare/tasks/database_cleanup_test.go @@ -11,6 +11,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/db" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/repository" + analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics" "github.com/glebarez/sqlite" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/require" @@ -36,13 +37,41 @@ func setupDatabaseCleanupTestDB(t *testing.T) context.Context { return context.Background() } -func TestCleanupDatabaseObservabilityDeletesTargetedRows(t *testing.T) { +func TestCleanupDatabaseObservabilityRejectsRetentionShorterThanTableTTL(t *testing.T) { + ctx := setupDatabaseCleanupTestDB(t) + + retentionDays := 7 // metric snapshots DDL TTL is 30 days + result, err := CleanupDatabaseObservability(ctx, DatabaseCleanupInput{ + Target: DatabaseCleanupTargetMetricSnapshots, + RetentionDays: &retentionDays, + }) + require.Error(t, err) + assert.Nil(t, result) + assert.Contains(t, err.Error(), "不能小于表 TTL") + assert.Contains(t, err.Error(), "30") +} + +func TestCleanupDatabaseObservabilityRejectsAccessLogRetentionShorterThanTableTTL(t *testing.T) { + ctx := setupDatabaseCleanupTestDB(t) + + retentionDays := 30 // access logs DDL TTL is 90 days + result, err := CleanupDatabaseObservability(ctx, DatabaseCleanupInput{ + Target: DatabaseCleanupTargetAccessLogs, + RetentionDays: &retentionDays, + }) + require.Error(t, err) + assert.Nil(t, result) + assert.Contains(t, err.Error(), "90") +} + +func TestCleanupDatabaseObservabilityMaterializeDoesNotClaimHardDelete(t *testing.T) { ctx := setupDatabaseCleanupTestDB(t) now := time.Now().UTC() + // One row past metric table TTL (30d), one still inside the window. require.NoError(t, model.InsertOpenFlareMetricSnapshot(ctx, &model.OpenFlareMetricSnapshot{ NodeID: "node-a", - CapturedAt: now.Add(-10 * 24 * time.Hour), + CapturedAt: now.Add(-40 * 24 * time.Hour), CPUUsagePercent: 10, })) require.NoError(t, model.InsertOpenFlareMetricSnapshot(ctx, &model.OpenFlareMetricSnapshot{ @@ -51,15 +80,22 @@ func TestCleanupDatabaseObservabilityDeletesTargetedRows(t *testing.T) { CPUUsagePercent: 20, })) - retentionDays := 7 + retentionDays := analyticsrepo.TableTTLDaysNodeMetricSnapshots result, err := CleanupDatabaseObservability(ctx, DatabaseCleanupInput{ Target: DatabaseCleanupTargetMetricSnapshots, RetentionDays: &retentionDays, }) require.NoError(t, err) assert.False(t, result.DeleteAll) - assert.Equal(t, int64(1), result.DeletedCount) + assert.Equal(t, analyticsrepo.CleanupModeTTLMaterialize, result.CleanupMode) + assert.Equal(t, analyticsrepo.TableTTLDaysNodeMetricSnapshots, result.TableTTLDays) + // MATERIALIZE is not a counted hard delete. + assert.Equal(t, int64(0), result.DeletedCount) + assert.Equal(t, int64(1), result.EligibleCount) + require.NotNil(t, result.Cutoff) + assert.True(t, result.Cutoff.Before(now.Add(-29*24*time.Hour))) + // Memory store applies the table-TTL cutoff for tests; only the recent row remains. rows, err := model.ListOpenFlareMetricSnapshotsSince(ctx, "", time.Time{}, 0) require.NoError(t, err) require.Len(t, rows, 1) @@ -94,20 +130,23 @@ func TestCleanupDatabaseObservabilityDeletesAllRowsWhenRetentionMissing(t *testi }) require.NoError(t, err) assert.True(t, result.DeleteAll) + assert.Equal(t, analyticsrepo.CleanupModeTruncate, result.CleanupMode) assert.Equal(t, int64(2), result.DeletedCount) + assert.Equal(t, int64(2), result.EligibleCount) rows, err := model.ListOpenFlareAccessLogs(ctx, model.OpenFlareAccessLogQuery{Page: 0, PageSize: 10}) require.NoError(t, err) assert.Empty(t, rows) } -func TestRunDatabaseAutoCleanupOnceDeletesAllObservabilityTargets(t *testing.T) { +func TestRunDatabaseAutoCleanupOnceClampsRetentionToTableTTL(t *testing.T) { ctx := setupDatabaseCleanupTestDB(t) now := time.Now().UTC() + // Access logs TTL=90d, metrics TTL=30d. Config retention=1 must clamp, not reject. require.NoError(t, model.InsertOpenFlareAccessLogsBatch(ctx, []*model.OpenFlareAccessLog{{ NodeID: "node-a", - LoggedAt: now.Add(-48 * time.Hour), + LoggedAt: now.Add(-100 * 24 * time.Hour), RemoteAddr: "203.0.113.10", Host: "example.com", Path: "/access", @@ -115,13 +154,13 @@ func TestRunDatabaseAutoCleanupOnceDeletesAllObservabilityTargets(t *testing.T) }})) require.NoError(t, model.InsertOpenFlareMetricSnapshot(ctx, &model.OpenFlareMetricSnapshot{ NodeID: "node-a", - CapturedAt: now.Add(-48 * time.Hour), + CapturedAt: now.Add(-40 * 24 * time.Hour), CPUUsagePercent: 10, })) require.NoError(t, model.InsertOpenFlareRequestReport(ctx, &model.OpenFlareRequestReport{ NodeID: "node-a", - WindowStartedAt: now.Add(-49 * time.Hour), - WindowEndedAt: now.Add(-48 * time.Hour), + WindowStartedAt: now.Add(-41 * 24 * time.Hour), + WindowEndedAt: now.Add(-40 * 24 * time.Hour), RequestCount: 15, })) @@ -132,6 +171,15 @@ func TestRunDatabaseAutoCleanupOnceDeletesAllObservabilityTargets(t *testing.T) require.NoError(t, err) require.NotNil(t, summary) require.Len(t, summary.Results, 6) + assert.Equal(t, 1, summary.RetentionDays) + + for _, result := range summary.Results { + assert.Equal(t, analyticsrepo.CleanupModeTTLMaterialize, result.CleanupMode) + assert.Equal(t, int64(0), result.DeletedCount, "target %s must not claim hard delete", result.Target) + assert.GreaterOrEqual(t, result.TableTTLDays, 30) + require.NotNil(t, result.RetentionDays) + assert.GreaterOrEqual(t, *result.RetentionDays, result.TableTTLDays) + } accessLogs, err := model.ListOpenFlareAccessLogs(ctx, model.OpenFlareAccessLogQuery{Page: 0, PageSize: 10}) require.NoError(t, err) @@ -145,3 +193,16 @@ func TestRunDatabaseAutoCleanupOnceDeletesAllObservabilityTargets(t *testing.T) require.NoError(t, err) assert.Empty(t, requestReports) } + +func TestTableTTLDaysForCleanupTarget(t *testing.T) { + days, ok := TableTTLDaysForCleanupTarget(DatabaseCleanupTargetAccessLogs) + require.True(t, ok) + assert.Equal(t, 90, days) + + days, ok = TableTTLDaysForCleanupTarget(DatabaseCleanupTargetMetricSnapshots) + require.True(t, ok) + assert.Equal(t, 30, days) + + _, ok = TableTTLDaysForCleanupTarget("unknown") + assert.False(t, ok) +} diff --git a/internal/apps/risk_control/logics.go b/internal/apps/risk_control/logics.go index 8cba0ade..84b1d083 100644 --- a/internal/apps/risk_control/logics.go +++ b/internal/apps/risk_control/logics.go @@ -6,6 +6,7 @@ package risk_control import ( "context" "sync" + "time" "github.com/Rain-kl/Wavelet/internal/config" "github.com/Rain-kl/Wavelet/internal/db/batchwriter" @@ -15,6 +16,11 @@ import ( "github.com/Rain-kl/Wavelet/pkg/logger" ) +const ( + // Bound visibility lag for sparse access-log traffic when MinBatchSize is not met. + accessLogMaxFlushWait = 3 * time.Second +) + var ( logWriterMu sync.RWMutex logWriter *batchwriter.Writer[*analytics.UserAccessLog] @@ -33,6 +39,8 @@ func InitLogWriter(ctx context.Context) { } cfg := batchwriter.DefaultConfig() + cfg.Name = "user_access_logs" + cfg.MaxFlushWait = accessLogMaxFlushWait writer, err := batchwriter.New[*analytics.UserAccessLog](cfg, func(ctx context.Context, items []*analytics.UserAccessLog) error { rows := make([]analytics.UserAccessLog, 0, len(items)) for _, item := range items { @@ -50,8 +58,8 @@ func InitLogWriter(ctx context.Context) { } logger.WarnF(context.Background(), "[RiskControl] Log queue full, dropping log item for path: %s", path) }), - batchwriter.WithFlushErrorHandler[*analytics.UserAccessLog](func(ctx context.Context, batchSize int, err error) { - logger.ErrorF(ctx, "[RiskControl] Send ClickHouse batch failed (batch=%d): %v", batchSize, err) + batchwriter.WithFlushErrorHandler[*analytics.UserAccessLog](func(ctx context.Context, items []*analytics.UserAccessLog, err error) { + logger.ErrorF(ctx, "[RiskControl] Send ClickHouse batch failed (batch=%d): %v", len(items), err) }), ) if err != nil { @@ -82,6 +90,16 @@ func IsBufferFull() bool { return writer.IsFull() } +// LogWriterStats returns queue depth and failure counters for the access-log writer. +// When the writer is not initialized, it returns a zero-value Stats with the expected name. +func LogWriterStats() batchwriter.Stats { + writer := currentLogWriter() + if writer == nil { + return batchwriter.Stats{Name: "user_access_logs"} + } + return writer.Stats() +} + // QueueAccessLog enqueues an access log without blocking. func QueueAccessLog(logItem *analytics.UserAccessLog) { writer := currentLogWriter() @@ -108,4 +126,4 @@ func currentLogWriter() *batchwriter.Writer[*analytics.UserAccessLog] { logWriterMu.RLock() defer logWriterMu.RUnlock() return logWriter -} \ No newline at end of file +} diff --git a/internal/apps/risk_control/logics_test.go b/internal/apps/risk_control/logics_test.go new file mode 100644 index 00000000..69fe2490 --- /dev/null +++ b/internal/apps/risk_control/logics_test.go @@ -0,0 +1,30 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package risk_control + +import ( + "testing" + "time" +) + +func TestAccessLogMaxFlushWaitInRange(t *testing.T) { + t.Parallel() + if accessLogMaxFlushWait < 2*time.Second || accessLogMaxFlushWait > 5*time.Second { + t.Fatalf("accessLogMaxFlushWait = %v, want in [2s, 5s]", accessLogMaxFlushWait) + } +} + +func TestLogWriterStatsWhenNil(t *testing.T) { + t.Parallel() + reset := SetLogWriterForTest(nil) + t.Cleanup(reset) + + stats := LogWriterStats() + if stats.Name != "user_access_logs" { + t.Fatalf("LogWriterStats().Name = %q, want user_access_logs", stats.Name) + } + if stats.Running { + t.Fatal("LogWriterStats().Running = true for nil writer, want false") + } +} diff --git a/internal/config/config.go b/internal/config/config.go index df5ed8cf..edf354dc 100644 --- a/internal/config/config.go +++ b/internal/config/config.go @@ -134,11 +134,14 @@ func applyClickHouseDefaults(c *configModel) { if c.ClickHouse.Username == "" { c.ClickHouse.Username = "default" } + // Pool / buffer defaults target small control-plane hosts (e.g. 3c6g): + // oversized open/idle pools waste RAM and amplify concurrent CH pressure; + // large block buffers add client memory without helping our small batch inserts. if c.ClickHouse.MaxIdleConn <= 0 { - c.ClickHouse.MaxIdleConn = 20 + c.ClickHouse.MaxIdleConn = 8 } if c.ClickHouse.MaxOpenConn <= 0 { - c.ClickHouse.MaxOpenConn = 50 + c.ClickHouse.MaxOpenConn = 16 } if c.ClickHouse.ConnMaxLifetime <= 0 { c.ClickHouse.ConnMaxLifetime = 3600 @@ -147,7 +150,7 @@ func applyClickHouseDefaults(c *configModel) { c.ClickHouse.DialTimeout = 5 } if c.ClickHouse.BlockBufferSize == 0 { - c.ClickHouse.BlockBufferSize = 100 + c.ClickHouse.BlockBufferSize = 32 } } diff --git a/internal/config/model.go b/internal/config/model.go index 6af5910e..0c8722af 100644 --- a/internal/config/model.go +++ b/internal/config/model.go @@ -71,18 +71,19 @@ type databaseReplicaConfig struct { Password string `mapstructure:"password"` } -// clickhouse 配置 +// clickHouseConfig ClickHouse 原生客户端配置。 +// 连接池 / block_buffer 默认值按小型控制面主机(如 3c6g)收敛,见 applyClickHouseDefaults。 type clickHouseConfig struct { Enabled bool `mapstructure:"enabled"` Hosts []string `mapstructure:"hosts"` Username string `mapstructure:"username"` Password string `mapstructure:"password"` Database string `mapstructure:"database"` - MaxIdleConn int `mapstructure:"max_idle_conn"` - MaxOpenConn int `mapstructure:"max_open_conn"` - ConnMaxLifetime int `mapstructure:"conn_max_lifetime"` - DialTimeout int `mapstructure:"dial_timeout"` - BlockBufferSize uint8 `mapstructure:"block_buffer_size"` + MaxIdleConn int `mapstructure:"max_idle_conn"` // 默认 8 + MaxOpenConn int `mapstructure:"max_open_conn"` // 默认 16 + ConnMaxLifetime int `mapstructure:"conn_max_lifetime"` // 秒 + DialTimeout int `mapstructure:"dial_timeout"` // 秒 + BlockBufferSize uint8 `mapstructure:"block_buffer_size"` // 默认 32 } // redisConfig Redis配置 diff --git a/internal/db/batchwriter/writer.go b/internal/db/batchwriter/writer.go index 13890fb4..bad25115 100644 --- a/internal/db/batchwriter/writer.go +++ b/internal/db/batchwriter/writer.go @@ -9,22 +9,34 @@ package batchwriter import ( "context" "sync" + "sync/atomic" "time" ) // FlushFunc persists a batch of queued items. It is invoked from the worker goroutine. type FlushFunc[T any] func(ctx context.Context, items []T) error -// FlushErrorHandler is called when FlushFunc returns an error. The batch is discarded -// after the handler returns; the worker continues processing. -type FlushErrorHandler func(ctx context.Context, batchSize int, err error) +// FlushErrorHandler is called when FlushFunc returns an error after optional retries. +// The batch is discarded after the handler returns; the worker continues processing. +// Handlers receive the failed items so callers can release dedup keys or re-queue. +type FlushErrorHandler[T any] func(ctx context.Context, items []T, err error) + +// Stats is a point-in-time snapshot of Writer queue and failure counters. +type Stats struct { + Name string `json:"name"` + Depth int `json:"depth"` + Cap int `json:"cap"` + Drops int64 `json:"drops"` + FlushErrors int64 `json:"flush_errors"` + Running bool `json:"running"` +} // Writer buffers items and flushes them by size or interval. type Writer[T any] struct { cfg Config flush FlushFunc[T] - onFlushError FlushErrorHandler + onFlushError FlushErrorHandler[T] onDrop func(T) startOnce sync.Once @@ -34,13 +46,16 @@ type Writer[T any] struct { ch chan T workerCtx context.Context done chan struct{} + + drops atomic.Int64 + flushErrors atomic.Int64 } // Option configures optional Writer callbacks. type Option[T any] func(*Writer[T]) // WithFlushErrorHandler registers a callback for flush failures. -func WithFlushErrorHandler[T any](handler FlushErrorHandler) Option[T] { +func WithFlushErrorHandler[T any](handler FlushErrorHandler[T]) Option[T] { return func(w *Writer[T]) { w.onFlushError = handler } @@ -168,6 +183,18 @@ func (w *Writer[T]) Cap() int { return w.cfg.QueueSize } +// Stats returns a point-in-time snapshot of queue depth and failure counters. +func (w *Writer[T]) Stats() Stats { + return Stats{ + Name: w.cfg.Name, + Depth: w.Len(), + Cap: w.Cap(), + Drops: w.drops.Load(), + FlushErrors: w.flushErrors.Load(), + Running: w.Running(), + } +} + func (w *Writer[T]) run() { ticker := time.NewTicker(w.cfg.FlushInterval) defer ticker.Stop() @@ -180,8 +207,9 @@ func (w *Writer[T]) run() { } items := append([]T(nil), batch...) if err := w.flush(w.workerCtx, items); err != nil { + w.flushErrors.Add(1) if w.onFlushError != nil { - w.onFlushError(w.workerCtx, len(items), err) + w.onFlushError(w.workerCtx, items, err) } } batch = batch[:0] @@ -228,8 +256,9 @@ func (w *Writer[T]) shouldFlushOnInterval(batchLen int, batchStartedAt time.Time } func (w *Writer[T]) notifyDrop(item T) { + w.drops.Add(1) if w.onDrop == nil { return } w.onDrop(item) -} \ No newline at end of file +} diff --git a/internal/db/batchwriter/writer_test.go b/internal/db/batchwriter/writer_test.go index 13c0fb1d..1818444c 100644 --- a/internal/db/batchwriter/writer_test.go +++ b/internal/db/batchwriter/writer_test.go @@ -36,7 +36,7 @@ func TestWriterFlushesOnMaxBatchSize(t *testing.T) { t.Parallel() var ( - mu sync.Mutex + mu sync.Mutex batches [][]int ) cfg := DefaultConfig() @@ -385,23 +385,24 @@ func TestWriterInvokesFlushErrorHandler(t *testing.T) { t.Parallel() cfg := DefaultConfig() + cfg.Name = "test-flush-err" cfg.MaxBatchSize = 1 cfg.FlushInterval = time.Hour flushErr := errors.New("flush failed") var ( - mu sync.Mutex - errCount int - batchSize int + mu sync.Mutex + errCount int + gotItems []int ) writer, err := New[int](cfg, func(context.Context, []int) error { return flushErr - }, WithFlushErrorHandler[int](func(_ context.Context, size int, err error) { + }, WithFlushErrorHandler[int](func(_ context.Context, items []int, err error) { mu.Lock() defer mu.Unlock() errCount++ - batchSize = size + gotItems = append([]int(nil), items...) if !errors.Is(err, flushErr) { t.Errorf("flush error = %v, want %v", err, flushErr) } @@ -434,13 +435,61 @@ func TestWriterInvokesFlushErrorHandler(t *testing.T) { mu.Lock() gotCount := errCount - gotSize := batchSize + items := gotItems mu.Unlock() if gotCount != 1 { t.Fatalf("flush error handler count = %d, want 1", gotCount) } - if gotSize != 1 { - t.Fatalf("flush error handler batch size = %d, want 1", gotSize) + if diff := cmp.Diff([]int{7}, items); diff != "" { + t.Fatalf("flush error handler items mismatch (-want +got):\n%s", diff) } -} \ No newline at end of file + + stats := writer.Stats() + if stats.FlushErrors != 1 { + t.Fatalf("Stats().FlushErrors = %d, want 1", stats.FlushErrors) + } + if stats.Name != "test-flush-err" { + t.Fatalf("Stats().Name = %q, want test-flush-err", stats.Name) + } +} + +func TestWriterStatsTracksDrops(t *testing.T) { + t.Parallel() + + cfg := DefaultConfig() + cfg.Name = "test-drops" + cfg.QueueSize = 1 + cfg.MaxBatchSize = 10 + cfg.FlushInterval = time.Hour + + writer, err := New[int](cfg, func(context.Context, []int) error { return nil }) + if err != nil { + t.Fatalf("New() error = %v", err) + } + + writer.Start(context.Background()) + t.Cleanup(func() { + stopCtx, cancel := context.WithTimeout(context.Background(), time.Second) + defer cancel() + _ = writer.Stop(stopCtx) + }) + + if !writer.TryEnqueue(1) { + t.Fatal("TryEnqueue(1) = false, want true") + } + if writer.TryEnqueue(2) { + t.Fatal("TryEnqueue(2) = true, want false") + } + + stats := writer.Stats() + if stats.Drops != 1 { + t.Fatalf("Stats().Drops = %d, want 1", stats.Drops) + } + if stats.Cap != 1 { + t.Fatalf("Stats().Cap = %d, want 1", stats.Cap) + } + if !stats.Running { + t.Fatal("Stats().Running = false, want true") + } +} diff --git a/internal/db/clickhouse.go b/internal/db/clickhouse.go index cc087eb6..5855a1ae 100644 --- a/internal/db/clickhouse.go +++ b/internal/db/clickhouse.go @@ -16,10 +16,18 @@ import ( ) const ( - clickhouseMaxExecTime = 60 // ClickHouse 最大执行时间(秒) - clickhouseReadTimeoutFactor = 2 // ReadTimeout 为 DialTimeout 的倍数 + clickhouseMaxExecTime = 60 // ClickHouse 最大执行时间(秒) + clickhouseReadTimeoutFactor = 2 // ReadTimeout 为 DialTimeout 的倍数 + + // async_insert 仅挂在运行时 ChConn(写路径)上,不进入 migrator OpenDB: + // 迁移/DDL 需要同步可见结果,且不应走异步 insert 缓冲。 + // + // 为何启用:batchwriter 仍可能在短间隔内写出相对小的块;服务端 async_insert + // 把多次 INSERT 合并成更大 part,减轻 3c6g 上 background merge 的 CPU 压力。 + // wait_for_async_insert=1:调用方在 flush 返回前等待落盘,避免进程崩溃丢批。 + // max_data_size / busy_timeout:约 10MB 或 ~2s 触发刷出,在延迟与 part 数之间折中。 clickhouseAsyncInsertMaxDataSize = 10_000_000 - clickhouseAsyncInsertBusyTimeoutMs = 1000 + clickhouseAsyncInsertBusyTimeoutMs = 2000 ) var ( @@ -52,6 +60,8 @@ func init() { log.Println("[ClickHouse] connection established successfully") } +// buildClickHouseOptions builds the runtime native client options (queries + batch inserts). +// Migrator uses a separate clickhouse.OpenDB path without async_insert settings. func buildClickHouseOptions() *clickhouse.Options { cfg := config.Config.ClickHouse diff --git a/internal/db/migrator/goose/clickhouse/202607100001_create_node_metric_openresty_hourly.sql b/internal/db/migrator/goose/clickhouse/202607100001_create_node_metric_openresty_hourly.sql new file mode 100644 index 00000000..0ec5c1b7 --- /dev/null +++ b/internal/db/migrator/goose/clickhouse/202607100001_create_node_metric_openresty_hourly.sql @@ -0,0 +1,80 @@ +-- +goose Up +-- Hourly capacity rollups (avg CPU/memory + counter min/max for in-hour delta approximation). +-- Network/disk counters are cumulative; max-min within an hour approximates that hour's delta +-- (cross-hour continuity is intentionally approximate for dashboard trends). +CREATE TABLE IF NOT EXISTS of_node_metric_capacity_hourly +( + node_id String, + hour DateTime, + cpu_usage_sum SimpleAggregateFunction(sum, Float64), + cpu_usage_count SimpleAggregateFunction(sum, UInt64), + memory_usage_sum SimpleAggregateFunction(sum, Float64), + memory_usage_count SimpleAggregateFunction(sum, UInt64), + network_rx_min SimpleAggregateFunction(min, Int64), + network_rx_max SimpleAggregateFunction(max, Int64), + network_tx_min SimpleAggregateFunction(min, Int64), + network_tx_max SimpleAggregateFunction(max, Int64), + disk_read_min SimpleAggregateFunction(min, Int64), + disk_read_max SimpleAggregateFunction(max, Int64), + disk_write_min SimpleAggregateFunction(min, Int64), + disk_write_max SimpleAggregateFunction(max, Int64) +) +ENGINE = AggregatingMergeTree() +PARTITION BY toYYYYMM(hour) +ORDER BY (node_id, hour) +TTL hour + INTERVAL 30 DAY; + +CREATE MATERIALIZED VIEW IF NOT EXISTS of_node_metric_capacity_hourly_mv +TO of_node_metric_capacity_hourly +AS +SELECT + node_id, + toStartOfHour(captured_at) AS hour, + sum(cpu_usage_percent) AS cpu_usage_sum, + toUInt64(count()) AS cpu_usage_count, + sum(if(memory_total_bytes > 0, (memory_used_bytes * 100.0) / memory_total_bytes, 0)) AS memory_usage_sum, + toUInt64(countIf(memory_total_bytes > 0)) AS memory_usage_count, + min(network_rx_bytes) AS network_rx_min, + max(network_rx_bytes) AS network_rx_max, + min(network_tx_bytes) AS network_tx_min, + max(network_tx_bytes) AS network_tx_max, + min(disk_read_bytes) AS disk_read_min, + max(disk_read_bytes) AS disk_read_max, + min(disk_write_bytes) AS disk_write_min, + max(disk_write_bytes) AS disk_write_max +FROM of_node_metric_snapshots +GROUP BY node_id, hour; + +-- Hourly OpenResty counter rollups (min/max per node-hour for delta approximation). +CREATE TABLE IF NOT EXISTS of_node_openresty_hourly +( + node_id String, + hour DateTime, + openresty_rx_min SimpleAggregateFunction(min, Int64), + openresty_rx_max SimpleAggregateFunction(max, Int64), + openresty_tx_min SimpleAggregateFunction(min, Int64), + openresty_tx_max SimpleAggregateFunction(max, Int64) +) +ENGINE = AggregatingMergeTree() +PARTITION BY toYYYYMM(hour) +ORDER BY (node_id, hour) +TTL hour + INTERVAL 30 DAY; + +CREATE MATERIALIZED VIEW IF NOT EXISTS of_node_openresty_hourly_mv +TO of_node_openresty_hourly +AS +SELECT + node_id, + toStartOfHour(captured_at) AS hour, + min(openresty_rx_bytes) AS openresty_rx_min, + max(openresty_rx_bytes) AS openresty_rx_max, + min(openresty_tx_bytes) AS openresty_tx_min, + max(openresty_tx_bytes) AS openresty_tx_max +FROM of_node_obs_openresty +GROUP BY node_id, hour; + +-- +goose Down +DROP VIEW IF EXISTS of_node_openresty_hourly_mv; +DROP TABLE IF EXISTS of_node_openresty_hourly; +DROP VIEW IF EXISTS of_node_metric_capacity_hourly_mv; +DROP TABLE IF EXISTS of_node_metric_capacity_hourly; diff --git a/internal/db/migrator/goose/clickhouse/202607100002_node_traffic_hourly_ttl_uv.sql b/internal/db/migrator/goose/clickhouse/202607100002_node_traffic_hourly_ttl_uv.sql new file mode 100644 index 00000000..62e81862 --- /dev/null +++ b/internal/db/migrator/goose/clickhouse/202607100002_node_traffic_hourly_ttl_uv.sql @@ -0,0 +1,42 @@ +-- +goose Up +-- Hourly traffic rollups: 30d TTL + UV aggregation semantics. +-- +-- unique_visitor_count on of_node_request_reports is per short report window +-- (agent local distinct count for that window only). Summing those values in the +-- MV (and again via SummingMergeTree part merges) invents a "true UV" number that +-- double-counts visitors across windows. Prefer max() as a peak-window estimate; +-- still NOT distinct visitors across the hour — UI/API must not overclaim. + +ALTER TABLE of_node_traffic_hourly + MODIFY TTL toDateTime(hour) + INTERVAL 30 DAY; + +DROP VIEW IF EXISTS of_node_traffic_hourly_mv; + +CREATE MATERIALIZED VIEW of_node_traffic_hourly_mv +TO of_node_traffic_hourly +AS +SELECT + node_id, + toStartOfHour(window_ended_at) AS hour, + sum(request_count) AS request_count, + sum(error_count) AS error_count, + -- Peak per-window UV estimate for the hour; not true cross-window distinct UV. + max(unique_visitor_count) AS unique_visitor_count +FROM of_node_request_reports +GROUP BY node_id, hour; + +-- +goose Down +-- TTL reverse is not safe without table rewrite; restore prior MV definition only. +DROP VIEW IF EXISTS of_node_traffic_hourly_mv; + +CREATE MATERIALIZED VIEW of_node_traffic_hourly_mv +TO of_node_traffic_hourly +AS +SELECT + node_id, + toStartOfHour(window_ended_at) AS hour, + sum(request_count) AS request_count, + sum(error_count) AS error_count, + sum(unique_visitor_count) AS unique_visitor_count +FROM of_node_request_reports +GROUP BY node_id, hour; diff --git a/internal/model/openflare_access_log_store.go b/internal/model/openflare_access_log_store.go index a69ba5e3..b68b6725 100644 --- a/internal/model/openflare_access_log_store.go +++ b/internal/model/openflare_access_log_store.go @@ -8,11 +8,34 @@ import ( "sync" "time" - "github.com/Rain-kl/Wavelet/internal/apps/openflare/chwriter" analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics" ) +// AccessLogInsertHooks queues node access logs for async ClickHouse write. +// Wired from openflare/chwriter.Init so model never imports the apps layer. +type AccessLogInsertHooks struct { + QueueNodeAccessLogs func(logs []analyticsmodel.NodeAccessLog) +} + +var ( + accessLogInsertHooksMu sync.RWMutex + accessLogInsertHooks AccessLogInsertHooks +) + +// SetAccessLogInsertHooks registers async queue callbacks for access log inserts. +func SetAccessLogInsertHooks(hooks AccessLogInsertHooks) { + accessLogInsertHooksMu.Lock() + accessLogInsertHooks = hooks + accessLogInsertHooksMu.Unlock() +} + +func currentAccessLogInsertHooks() AccessLogInsertHooks { + accessLogInsertHooksMu.RLock() + defer accessLogInsertHooksMu.RUnlock() + return accessLogInsertHooks +} + type accessLogStore interface { InsertBatch(ctx context.Context, records []*OpenFlareAccessLog) error List(ctx context.Context, query OpenFlareAccessLogQuery) ([]*OpenFlareAccessLog, error) @@ -75,7 +98,9 @@ func (clickhouseAccessLogStore) InsertBatch(_ context.Context, records []*OpenFl } logs = append(logs, toAnalyticsNodeAccessLog(record)) } - chwriter.QueueNodeAccessLogs(logs) + if hook := currentAccessLogInsertHooks().QueueNodeAccessLogs; hook != nil { + hook(logs) + } return nil } diff --git a/internal/model/openflare_insert_hooks_test.go b/internal/model/openflare_insert_hooks_test.go new file mode 100644 index 00000000..fe6f45be --- /dev/null +++ b/internal/model/openflare_insert_hooks_test.go @@ -0,0 +1,75 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package model + +import ( + "context" + "testing" + "time" + + analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" +) + +// Hook setters are process-global; keep these tests serial. + +func TestObservabilityInsertHooksAreInvoked(t *testing.T) { + var gotSnapshot analyticsmodel.NodeMetricSnapshot + SetObservabilityInsertHooks(ObservabilityInsertHooks{ + QueueMetricSnapshot: func(s analyticsmodel.NodeMetricSnapshot) { + gotSnapshot = s + }, + }) + t.Cleanup(func() { + SetObservabilityInsertHooks(ObservabilityInsertHooks{}) + }) + + record := &OpenFlareMetricSnapshot{ + NodeID: "node-1", + CapturedAt: time.Unix(100, 0).UTC(), + } + if err := (clickhouseObservabilityStore{}).InsertMetricSnapshot(context.Background(), record); err != nil { + t.Fatalf("InsertMetricSnapshot error = %v", err) + } + if gotSnapshot.NodeID != "node-1" { + t.Fatalf("hook node id = %q, want node-1", gotSnapshot.NodeID) + } +} + +func TestAccessLogInsertHooksAreInvoked(t *testing.T) { + var got []analyticsmodel.NodeAccessLog + SetAccessLogInsertHooks(AccessLogInsertHooks{ + QueueNodeAccessLogs: func(logs []analyticsmodel.NodeAccessLog) { + got = append([]analyticsmodel.NodeAccessLog(nil), logs...) + }, + }) + t.Cleanup(func() { + SetAccessLogInsertHooks(AccessLogInsertHooks{}) + }) + + records := []*OpenFlareAccessLog{ + {NodeID: "n1", Path: "/a"}, + {NodeID: "n1", Path: "/b"}, + } + if err := (clickhouseAccessLogStore{}).InsertBatch(context.Background(), records); err != nil { + t.Fatalf("InsertBatch error = %v", err) + } + if len(got) != 2 { + t.Fatalf("hook logs = %d, want 2", len(got)) + } + if got[0].Path != "/a" || got[1].Path != "/b" { + t.Fatalf("hook paths = %q/%q, want /a /b", got[0].Path, got[1].Path) + } +} + +func TestInsertHooksNoopWhenUnset(t *testing.T) { + SetObservabilityInsertHooks(ObservabilityInsertHooks{}) + SetAccessLogInsertHooks(AccessLogInsertHooks{}) + + if err := (clickhouseObservabilityStore{}).InsertMetricSnapshot(context.Background(), &OpenFlareMetricSnapshot{NodeID: "x"}); err != nil { + t.Fatalf("InsertMetricSnapshot with nil hook error = %v", err) + } + if err := (clickhouseAccessLogStore{}).InsertBatch(context.Background(), []*OpenFlareAccessLog{{NodeID: "x"}}); err != nil { + t.Fatalf("InsertBatch with nil hook error = %v", err) + } +} diff --git a/internal/model/openflare_observability.go b/internal/model/openflare_observability.go index 9f769fb5..db31e6ef 100644 --- a/internal/model/openflare_observability.go +++ b/internal/model/openflare_observability.go @@ -333,11 +333,82 @@ func ListOpenFlareMetricSnapshotsSince(ctx context.Context, nodeID string, since return currentObservabilityStore().ListMetricSnapshots(ctx, nodeID, since, limit) } +// ListOpenFlareLatestMetricSnapshotsSince returns the latest metric snapshot per node. +// Prefer ClickHouse LIMIT 1 BY; on CH unavailability fall back to store list + reduce. +func ListOpenFlareLatestMetricSnapshotsSince(ctx context.Context, nodeID string, since time.Time) ([]*OpenFlareMetricSnapshot, error) { + rows, err := analyticsrepo.ListLatestNodeMetricSnapshots(ctx, analyticsrepo.NodeObservabilityFilter{ + NodeID: nodeID, + Since: since, + }) + if err == nil { + return fromAnalyticsNodeMetricSnapshots(rows), nil + } + // Fallback for unit tests (memory store) and environments without ClickHouse. + all, listErr := ListOpenFlareMetricSnapshotsSince(ctx, nodeID, since, 0) + if listErr != nil { + return nil, err + } + return openFlareLatestMetricSnapshots(all), nil +} + // ListOpenFlareRequestReportsSince returns request reports since the given time. func ListOpenFlareRequestReportsSince(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareRequestReport, error) { return currentObservabilityStore().ListRequestReports(ctx, nodeID, since, limit) } +// ListOpenFlareLatestRequestReportsSince returns the latest request report per node. +// Prefer ClickHouse LIMIT 1 BY; on CH unavailability fall back to store list + reduce. +func ListOpenFlareLatestRequestReportsSince(ctx context.Context, nodeID string, since time.Time) ([]*OpenFlareRequestReport, error) { + rows, err := analyticsrepo.ListLatestNodeRequestReports(ctx, analyticsrepo.NodeObservabilityFilter{ + NodeID: nodeID, + Since: since, + }) + if err == nil { + return fromAnalyticsNodeRequestReports(rows), nil + } + all, listErr := ListOpenFlareRequestReportsSince(ctx, nodeID, since, 0) + if listErr != nil { + return nil, err + } + return openFlareLatestRequestReports(all), nil +} + +func openFlareLatestMetricSnapshots(snapshots []*OpenFlareMetricSnapshot) []*OpenFlareMetricSnapshot { + latestByNode := make(map[string]*OpenFlareMetricSnapshot, len(snapshots)) + for _, snapshot := range snapshots { + if snapshot == nil || snapshot.NodeID == "" { + continue + } + if existing, ok := latestByNode[snapshot.NodeID]; ok && !snapshot.CapturedAt.After(existing.CapturedAt) { + continue + } + latestByNode[snapshot.NodeID] = snapshot + } + result := make([]*OpenFlareMetricSnapshot, 0, len(latestByNode)) + for _, snapshot := range latestByNode { + result = append(result, snapshot) + } + return result +} + +func openFlareLatestRequestReports(reports []*OpenFlareRequestReport) []*OpenFlareRequestReport { + latestByNode := make(map[string]*OpenFlareRequestReport, len(reports)) + for _, report := range reports { + if report == nil || report.NodeID == "" { + continue + } + if existing, ok := latestByNode[report.NodeID]; ok && !report.WindowEndedAt.After(existing.WindowEndedAt) { + continue + } + latestByNode[report.NodeID] = report + } + result := make([]*OpenFlareRequestReport, 0, len(latestByNode)) + for _, report := range latestByNode { + result = append(result, report) + } + return result +} + // OpenFlareTrafficHourly is an hourly traffic rollup row. type OpenFlareTrafficHourly struct { NodeID string `json:"node_id"` diff --git a/internal/model/openflare_observability_store.go b/internal/model/openflare_observability_store.go index 25136adf..56d60edd 100644 --- a/internal/model/openflare_observability_store.go +++ b/internal/model/openflare_observability_store.go @@ -9,11 +9,38 @@ import ( "sync" "time" - "github.com/Rain-kl/Wavelet/internal/apps/openflare/chwriter" analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics" ) +// ObservabilityInsertHooks queues observability rows for async ClickHouse write. +// Wired from openflare/chwriter.Init so model never imports the apps layer. +type ObservabilityInsertHooks struct { + QueueMetricSnapshot func(analyticsmodel.NodeMetricSnapshot) + QueueRequestReport func(analyticsmodel.NodeRequestReport) + QueueOpenrestyObservation func(analyticsmodel.NodeObsOpenresty) + QueueFrpsObservation func(analyticsmodel.NodeObsFrps) + QueueFrpcObservation func(analyticsmodel.NodeObsFrpc) +} + +var ( + observabilityInsertHooksMu sync.RWMutex + observabilityInsertHooks ObservabilityInsertHooks +) + +// SetObservabilityInsertHooks registers async queue callbacks for observability inserts. +func SetObservabilityInsertHooks(hooks ObservabilityInsertHooks) { + observabilityInsertHooksMu.Lock() + observabilityInsertHooks = hooks + observabilityInsertHooksMu.Unlock() +} + +func currentObservabilityInsertHooks() ObservabilityInsertHooks { + observabilityInsertHooksMu.RLock() + defer observabilityInsertHooksMu.RUnlock() + return observabilityInsertHooks +} + type observabilityStore interface { InsertMetricSnapshot(ctx context.Context, record *OpenFlareMetricSnapshot) error ListMetricSnapshots(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareMetricSnapshot, error) @@ -79,7 +106,9 @@ func (clickhouseObservabilityStore) InsertMetricSnapshot(_ context.Context, reco if record == nil { return nil } - chwriter.QueueMetricSnapshot(toAnalyticsNodeMetricSnapshot(record)) + if hook := currentObservabilityInsertHooks().QueueMetricSnapshot; hook != nil { + hook(toAnalyticsNodeMetricSnapshot(record)) + } return nil } @@ -103,7 +132,9 @@ func (clickhouseObservabilityStore) InsertRequestReport(_ context.Context, recor if record == nil { return nil } - chwriter.QueueRequestReport(toAnalyticsNodeRequestReport(record)) + if hook := currentObservabilityInsertHooks().QueueRequestReport; hook != nil { + hook(toAnalyticsNodeRequestReport(record)) + } return nil } @@ -127,7 +158,9 @@ func (clickhouseObservabilityStore) InsertNodeObservationOpenresty(_ context.Con if record == nil { return nil } - chwriter.QueueOpenrestyObservation(toAnalyticsNodeObsOpenresty(record)) + if hook := currentObservabilityInsertHooks().QueueOpenrestyObservation; hook != nil { + hook(toAnalyticsNodeObsOpenresty(record)) + } return nil } @@ -151,7 +184,9 @@ func (clickhouseObservabilityStore) InsertNodeObservationFrps(_ context.Context, if record == nil { return nil } - chwriter.QueueFrpsObservation(toAnalyticsNodeObsFrps(record)) + if hook := currentObservabilityInsertHooks().QueueFrpsObservation; hook != nil { + hook(toAnalyticsNodeObsFrps(record)) + } return nil } @@ -175,7 +210,9 @@ func (clickhouseObservabilityStore) InsertNodeObservationFrpc(_ context.Context, if record == nil { return nil } - chwriter.QueueFrpcObservation(toAnalyticsNodeObsFrpc(record)) + if hook := currentObservabilityInsertHooks().QueueFrpcObservation; hook != nil { + hook(toAnalyticsNodeObsFrpc(record)) + } return nil } diff --git a/internal/repository/analytics/access_log_test.go b/internal/repository/analytics/access_log_test.go index 4ba181bc..e248b01b 100644 --- a/internal/repository/analytics/access_log_test.go +++ b/internal/repository/analytics/access_log_test.go @@ -99,6 +99,9 @@ type mockConn struct { batchQuery string prepareCalled bool preparedQuery string + queries []string + queryArgs [][]any + queryFn func(ctx context.Context, query string, args ...any) (driver.Rows, error) } func (m *mockConn) Contributors() []string { return nil } @@ -107,8 +110,13 @@ func (m *mockConn) ServerVersion() (*driver.ServerVersion, error) { return nil, func (m *mockConn) Select(_ context.Context, _ any, _ string, _ ...any) error { return nil } -func (m *mockConn) Query(_ context.Context, _ string, _ ...any) (driver.Rows, error) { - return nil, nil +func (m *mockConn) Query(ctx context.Context, query string, args ...any) (driver.Rows, error) { + m.queries = append(m.queries, query) + m.queryArgs = append(m.queryArgs, args) + if m.queryFn != nil { + return m.queryFn(ctx, query, args...) + } + return &mockRows{}, nil } func (m *mockConn) QueryRow(_ context.Context, _ string, _ ...any) driver.Row { return nil } @@ -158,4 +166,96 @@ func (m *mockBatch) Rows() int { return len(m.rows) } func (m *mockBatch) Columns() []column.Interface { return nil } -func (m *mockBatch) Close() error { return nil } \ No newline at end of file +func (m *mockBatch) Close() error { return nil } + +// mockRows is an empty driver.Rows implementation for query-path unit tests. +type mockRows struct { + index int + data [][]any + err error +} + +func (m *mockRows) Next() bool { + if m.err != nil { + return false + } + if m.index >= len(m.data) { + return false + } + m.index++ + return true +} + +func (m *mockRows) Scan(dest ...any) error { + if m.err != nil { + return m.err + } + if m.index == 0 || m.index > len(m.data) { + return nil + } + row := m.data[m.index-1] + for i := range dest { + if i >= len(row) { + break + } + if err := assignMockScanValue(dest[i], row[i]); err != nil { + return err + } + } + return nil +} + +func (m *mockRows) ScanStruct(_ any) error { return nil } + +func (m *mockRows) ColumnTypes() []driver.ColumnType { return nil } + +func (m *mockRows) Totals(_ ...any) error { return nil } + +func (m *mockRows) Columns() []string { return nil } + +func (m *mockRows) Close() error { return nil } + +func (m *mockRows) Err() error { return m.err } + +func (m *mockRows) HasData() bool { return len(m.data) > 0 } + +func assignMockScanValue(dest any, value any) error { + switch d := dest.(type) { + case *string: + if v, ok := value.(string); ok { + *d = v + } + case *uint64: + switch v := value.(type) { + case uint64: + *d = v + case int: + *d = uint64(v) + case int64: + *d = uint64(v) + } + case *int64: + switch v := value.(type) { + case int64: + *d = v + case int: + *d = int64(v) + case uint64: + *d = int64(v) + } + case *float64: + switch v := value.(type) { + case float64: + *d = v + case float32: + *d = float64(v) + case int: + *d = float64(v) + } + case *time.Time: + if v, ok := value.(time.Time); ok { + *d = v + } + } + return nil +} \ No newline at end of file diff --git a/internal/repository/analytics/clickhouse_maintenance.go b/internal/repository/analytics/clickhouse_maintenance.go index 57224d3c..460ceb8d 100644 --- a/internal/repository/analytics/clickhouse_maintenance.go +++ b/internal/repository/analytics/clickhouse_maintenance.go @@ -6,21 +6,48 @@ package analytics import ( "context" "fmt" + "time" "github.com/ClickHouse/clickhouse-go/v2/lib/driver" ) +// DDL TTL days for analytics tables (must match goose ClickHouse migrations). +const ( + // TableTTLDaysNodeAccessLogs is the of_node_access_logs TTL (90 days). + TableTTLDaysNodeAccessLogs = 90 + // TableTTLDaysNodeMetricSnapshots is the of_node_metric_snapshots TTL (30 days). + TableTTLDaysNodeMetricSnapshots = 30 + // TableTTLDaysNodeRequestReports is the of_node_request_reports TTL (30 days). + TableTTLDaysNodeRequestReports = 30 + // TableTTLDaysNodeObs is the of_node_obs_* TTL (30 days). + TableTTLDaysNodeObs = 30 + // TableTTLDaysUserAccessLogs is the w_user_access_logs TTL (180 days). + TableTTLDaysUserAccessLogs = 180 +) + const ( // CleanupModeTTLMaterialize expires rows via table TTL instead of ALTER DELETE mutations. + // This is not a hard delete: deleted_count must stay 0; use EligibleCount as an estimate. CleanupModeTTLMaterialize = "ttl_materialize" - // CleanupModeTruncate removes all rows via TRUNCATE TABLE. + // CleanupModeTruncate removes all rows via TRUNCATE TABLE (hard delete). CleanupModeTruncate = "truncate" ) // CleanupOutcome describes a non-mutation ClickHouse cleanup operation. +// +// For CleanupModeTruncate: +// - DeletedCount and EligibleCount are the rows removed by TRUNCATE. +// +// For CleanupModeTTLMaterialize: +// - DeletedCount is always 0 (MATERIALIZE TTL is async / not a counted hard delete). +// - EligibleCount is an estimate of rows already past the table TTL policy (not an +// arbitrary user cutoff younger than the DDL TTL). +// - TableTTLDays is the DDL TTL used for the estimate and materialize. type CleanupOutcome struct { EligibleCount int64 + DeletedCount int64 Mode string + TableTTLDays int } func countClickHouseRows(ctx context.Context, conn driver.Conn, countSQL string, countArgs []any) (int64, error) { @@ -39,21 +66,41 @@ func materializeTableTTL(ctx context.Context, conn driver.Conn, tableName string return nil } -func expireRowsViaTTL(ctx context.Context, conn driver.Conn, tableName string, countSQL string, countArgs []any) (CleanupOutcome, error) { +// tableTTLCutoff returns the UTC instant at which rows become eligible under a fixed day TTL. +func tableTTLCutoff(tableTTLDays int, now time.Time) time.Time { + if tableTTLDays < 1 { + tableTTLDays = 1 + } + return now.UTC().Add(-time.Duration(tableTTLDays) * 24 * time.Hour) +} + +// materializeExpiredByTableTTL force-materializes table TTL and estimates rows past that policy. +// +// countSQL must count only rows older than the table TTL (callers pass tableTTLCutoff args). +// Node-scoped filters may be used for the estimate only; MATERIALIZE is always table-global. +func materializeExpiredByTableTTL( + ctx context.Context, + conn driver.Conn, + tableName string, + tableTTLDays int, + countSQL string, + countArgs []any, +) (CleanupOutcome, error) { + outcome := CleanupOutcome{ + Mode: CleanupModeTTLMaterialize, + TableTTLDays: tableTTLDays, + } count, err := countClickHouseRows(ctx, conn, countSQL, countArgs) if err != nil { return CleanupOutcome{}, err } - if count == 0 { - return CleanupOutcome{Mode: CleanupModeTTLMaterialize}, nil - } + outcome.EligibleCount = count + // Always force materialize so ClickHouse applies the DDL TTL policy promptly. + // EligibleCount is informational only; MATERIALIZE does not return a deleted row count. if err := materializeTableTTL(ctx, conn, tableName); err != nil { return CleanupOutcome{}, err } - return CleanupOutcome{ - EligibleCount: count, - Mode: CleanupModeTTLMaterialize, - }, nil + return outcome, nil } func truncateClickHouseTable(ctx context.Context, conn driver.Conn, tableName string) (CleanupOutcome, error) { @@ -69,6 +116,7 @@ func truncateClickHouseTable(ctx context.Context, conn driver.Conn, tableName st } return CleanupOutcome{ EligibleCount: count, + DeletedCount: count, Mode: CleanupModeTruncate, }, nil -} \ No newline at end of file +} diff --git a/internal/repository/analytics/clickhouse_maintenance_test.go b/internal/repository/analytics/clickhouse_maintenance_test.go new file mode 100644 index 00000000..c3898da0 --- /dev/null +++ b/internal/repository/analytics/clickhouse_maintenance_test.go @@ -0,0 +1,37 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package analytics + +import ( + "testing" + "time" + + "github.com/stretchr/testify/assert" +) + +func TestTableTTLCutoff(t *testing.T) { + now := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC) + got := tableTTLCutoff(30, now) + assert.Equal(t, now.Add(-30*24*time.Hour), got) + + got = tableTTLCutoff(90, now) + assert.Equal(t, now.Add(-90*24*time.Hour), got) + + // Invalid TTL floors to 1 day. + got = tableTTLCutoff(0, now) + assert.Equal(t, now.Add(-24*time.Hour), got) +} + +func TestCleanupModeConstants(t *testing.T) { + assert.Equal(t, "ttl_materialize", CleanupModeTTLMaterialize) + assert.Equal(t, "truncate", CleanupModeTruncate) +} + +func TestTableTTLDaysMatchDDL(t *testing.T) { + assert.Equal(t, 90, TableTTLDaysNodeAccessLogs) + assert.Equal(t, 30, TableTTLDaysNodeMetricSnapshots) + assert.Equal(t, 30, TableTTLDaysNodeRequestReports) + assert.Equal(t, 30, TableTTLDaysNodeObs) + assert.Equal(t, 180, TableTTLDaysUserAccessLogs) +} diff --git a/internal/repository/analytics/clickhouse_stats.go b/internal/repository/analytics/clickhouse_stats.go index ed4cd8f3..18e4ca1b 100644 --- a/internal/repository/analytics/clickhouse_stats.go +++ b/internal/repository/analytics/clickhouse_stats.go @@ -9,16 +9,20 @@ import ( "github.com/Rain-kl/Wavelet/internal/config" "github.com/Rain-kl/Wavelet/internal/db" + "github.com/Rain-kl/Wavelet/internal/db/batchwriter" ) -// ClickHouseOperationalStats summarizes ClickHouse merge/mutation pressure. +// ClickHouseOperationalStats summarizes ClickHouse merge/mutation pressure +// and in-process batch writer queue health. type ClickHouseOperationalStats struct { - Database string `json:"database"` - ActiveParts int64 `json:"active_parts"` - TotalRows int64 `json:"total_rows"` - PendingMutations int64 `json:"pending_mutations"` - AsyncInsertQueue int64 `json:"async_insert_queue"` - AsyncInsertBytes int64 `json:"async_insert_bytes"` + Database string `json:"database"` + ActiveParts int64 `json:"active_parts"` + TotalRows int64 `json:"total_rows"` + PendingMutations int64 `json:"pending_mutations"` + AsyncInsertQueue int64 `json:"async_insert_queue"` + AsyncInsertBytes int64 `json:"async_insert_bytes"` + // BatchWriters reports in-process queue depth/drops/flush errors for CH writers. + BatchWriters []batchwriter.Stats `json:"batch_writers,omitempty"` } // GetClickHouseOperationalStats returns operational metrics for the configured database. @@ -67,4 +71,4 @@ WHERE database = ?` } return stats, nil -} \ No newline at end of file +} diff --git a/internal/repository/analytics/node_access_log_delete.go b/internal/repository/analytics/node_access_log_delete.go index f910dd16..ae9020bb 100644 --- a/internal/repository/analytics/node_access_log_delete.go +++ b/internal/repository/analytics/node_access_log_delete.go @@ -9,7 +9,7 @@ import ( "time" ) -// DeleteAllNodeAccessLogs deletes all node access logs. +// DeleteAllNodeAccessLogs hard-deletes all node access logs via TRUNCATE. func DeleteAllNodeAccessLogs(ctx context.Context) (int64, error) { conn, err := nodeAccessLogConn() if err != nil { @@ -19,47 +19,69 @@ func DeleteAllNodeAccessLogs(ctx context.Context) (int64, error) { if err != nil { return 0, err } - return outcome.EligibleCount, nil + return outcome.DeletedCount, nil } -// DeleteNodeAccessLogsBefore expires logs older than cutoff via table TTL. -func DeleteNodeAccessLogsBefore(ctx context.Context, cutoff time.Time) (int64, error) { - conn, err := nodeAccessLogConn() +// DeleteNodeAccessLogsBefore force-materializes of_node_access_logs table TTL. +// +// The cutoff argument is kept for call-site compatibility and is not used to select rows: +// ClickHouse MATERIALIZE TTL only enforces the DDL policy (TableTTLDaysNodeAccessLogs). +// Returns an estimate of rows past table TTL as the int64 (not a hard-deleted count). +// Callers that need honest API fields should prefer MaterializeNodeAccessLogsTTL. +func DeleteNodeAccessLogsBefore(ctx context.Context, _ time.Time) (int64, error) { + outcome, err := MaterializeNodeAccessLogsTTL(ctx) if err != nil { return 0, err } + return outcome.EligibleCount, nil +} + +// MaterializeNodeAccessLogsTTL force-materializes table TTL and reports an honest outcome. +func MaterializeNodeAccessLogsTTL(ctx context.Context) (CleanupOutcome, error) { + conn, err := nodeAccessLogConn() + if err != nil { + return CleanupOutcome{}, err + } tableName := nodeAccessLogTableName() - cutoff = cutoff.UTC() - outcome, err := expireRowsViaTTL( + ttlDays := TableTTLDaysNodeAccessLogs + cutoff := tableTTLCutoff(ttlDays, time.Now()) + return materializeExpiredByTableTTL( ctx, conn, tableName, + ttlDays, fmt.Sprintf("SELECT count() FROM %s WHERE logged_at < ?", tableName), []any{cutoff}, ) +} + +// DeleteNodeAccessLogsByNodeBefore force-materializes table-global TTL. +// +// Node-scoped hard delete is not supported: MATERIALIZE TTL is table-global. +// The returned count is an estimate of rows for nodeID past table TTL only. +func DeleteNodeAccessLogsByNodeBefore(ctx context.Context, nodeID string, _ time.Time) (int64, error) { + outcome, err := MaterializeNodeAccessLogsTTLByNode(ctx, nodeID) if err != nil { return 0, err } return outcome.EligibleCount, nil } -// DeleteNodeAccessLogsByNodeBefore expires logs for a node older than cutoff via table TTL. -func DeleteNodeAccessLogsByNodeBefore(ctx context.Context, nodeID string, before time.Time) (int64, error) { +// MaterializeNodeAccessLogsTTLByNode materializes table-global TTL and estimates node-scoped rows past TTL. +func MaterializeNodeAccessLogsTTLByNode(ctx context.Context, nodeID string) (CleanupOutcome, error) { conn, err := nodeAccessLogConn() if err != nil { - return 0, err + return CleanupOutcome{}, err } tableName := nodeAccessLogTableName() - before = before.UTC() - outcome, err := expireRowsViaTTL( + ttlDays := TableTTLDaysNodeAccessLogs + cutoff := tableTTLCutoff(ttlDays, time.Now()) + return materializeExpiredByTableTTL( ctx, conn, tableName, + ttlDays, fmt.Sprintf("SELECT count() FROM %s WHERE node_id = ? AND logged_at < ?", tableName), - []any{nodeID, before}, + []any{nodeID, cutoff}, ) - if err != nil { - return 0, err - } - return outcome.EligibleCount, nil -} \ No newline at end of file +} diff --git a/internal/repository/analytics/node_observability.go b/internal/repository/analytics/node_observability.go index 1b53b4fd..8adc8c4b 100644 --- a/internal/repository/analytics/node_observability.go +++ b/internal/repository/analytics/node_observability.go @@ -45,6 +45,27 @@ ORDER BY %s`, tableName, clause, nodeObservabilityCapturedAtOrderClause()) return scanNodeMetricSnapshotRows(rows) } +// ListLatestNodeMetricSnapshots returns the latest metric snapshot per node_id. +// Uses ClickHouse LIMIT 1 BY so dashboard health does not depend on a global raw LIMIT. +func ListLatestNodeMetricSnapshots(ctx context.Context, filter NodeObservabilityFilter) ([]analyticsmodel.NodeMetricSnapshot, error) { + conn, err := observabilityConn() + if err != nil { + return nil, err + } + clause, args := buildNodeObservabilityFilterClause(filter, "captured_at") + sql := fmt.Sprintf(` +SELECT id, node_id, captured_at, cpu_usage_percent, memory_used_bytes, memory_total_bytes, storage_used_bytes, storage_total_bytes, disk_read_bytes, disk_write_bytes, network_rx_bytes, network_tx_bytes, created_at +FROM %s +WHERE %s +ORDER BY %s%s`, nodeMetricSnapshotTableName(), clause, nodeObservabilityCapturedAtOrderClause(), clickHouseLimit1ByNodeIDClause) + rows, err := conn.Query(ctx, sql, args...) + if err != nil { + return nil, fmt.Errorf("list latest node metric snapshots: %w", err) + } + defer func() { _ = rows.Close() }() + return scanNodeMetricSnapshotRows(rows) +} + // ListNodeRequestReports returns request reports matching filter. func ListNodeRequestReports(ctx context.Context, filter NodeObservabilityFilter) ([]analyticsmodel.NodeRequestReport, error) { conn, err := observabilityConn() @@ -70,6 +91,27 @@ ORDER BY %s`, tableName, clause, nodeObservabilityWindowEndedAtOrderClause()) return scanNodeRequestReportRows(rows) } +// ListLatestNodeRequestReports returns the latest request report per node_id. +// Uses ClickHouse LIMIT 1 BY so dashboard traffic health is not skewed by a global raw LIMIT. +func ListLatestNodeRequestReports(ctx context.Context, filter NodeObservabilityFilter) ([]analyticsmodel.NodeRequestReport, error) { + conn, err := observabilityConn() + if err != nil { + return nil, err + } + clause, args := buildNodeObservabilityFilterClause(filter, "window_ended_at") + sql := fmt.Sprintf(` +SELECT id, node_id, window_started_at, window_ended_at, request_count, error_count, unique_visitor_count, status_codes_json, top_domains_json, source_countries_json, created_at +FROM %s +WHERE %s +ORDER BY %s%s`, nodeRequestReportTableName(), clause, nodeObservabilityWindowEndedAtOrderClause(), clickHouseLimit1ByNodeIDClause) + rows, err := conn.Query(ctx, sql, args...) + if err != nil { + return nil, fmt.Errorf("list latest node request reports: %w", err) + } + defer func() { _ = rows.Close() }() + return scanNodeRequestReportRows(rows) +} + // ListNodeObsOpenresty returns OpenResty observations matching filter. func ListNodeObsOpenresty(ctx context.Context, filter NodeObservabilityFilter) ([]analyticsmodel.NodeObsOpenresty, error) { conn, err := observabilityConn() @@ -248,6 +290,10 @@ func scanNodeObsFrpsRows(rows driver.Rows) ([]analyticsmodel.NodeObsFrps, error) const nodeTrafficHourlyTableName = "of_node_traffic_hourly" // NodeTrafficHourly is an hourly traffic rollup row. +// +// UniqueVisitorCount is a peak per-window estimate from short request reports +// (MV uses max()), not true distinct visitors across the hour. SummingMergeTree +// may still inflate residual unmerged parts; do not present as exact UV. type NodeTrafficHourly struct { NodeID string Hour time.Time @@ -258,7 +304,8 @@ type NodeTrafficHourly struct { // NodeMetricHourly is an hourly metric snapshot aggregation row. // -// Disk and host network counters are cumulative; deltas use consecutive +// Disk and host network counters are cumulative. Prefer pre-aggregated min/max +// deltas from of_node_metric_capacity_hourly; raw fallback uses consecutive // lagInFrame samples per node (negative deltas after counter reset are dropped). type NodeMetricHourly struct { Hour time.Time @@ -292,7 +339,7 @@ SELECT hour, sum(request_count) AS request_count, sum(error_count) AS error_count, - sum(unique_visitor_count) AS unique_visitor_count + max(unique_visitor_count) AS unique_visitor_count FROM %s WHERE %s GROUP BY node_id, hour @@ -322,8 +369,43 @@ ORDER BY hour ASC`, nodeTrafficHourlyTableName, clause) } // ListNodeMetricHourly returns hourly metric snapshot aggregates matching filter. -// Capacity uses sample averages; network/disk use consecutive counter deltas via lagInFrame. +// Prefers of_node_metric_capacity_hourly rollups; falls back to raw lagInFrame on empty/error. func ListNodeMetricHourly(ctx context.Context, filter NodeObservabilityFilter) ([]NodeMetricHourly, error) { + if rows, err := listNodeMetricHourlyFromRollup(ctx, filter); err == nil && len(rows) > 0 { + return rows, nil + } + return listNodeMetricHourlyFromRaw(ctx, filter) +} + +func listNodeMetricHourlyFromRollup(ctx context.Context, filter NodeObservabilityFilter) ([]NodeMetricHourly, error) { + conn, err := observabilityConn() + if err != nil { + return nil, err + } + clause, args := buildNodeObservabilityFilterClause(filter, "hour") + sql := fmt.Sprintf(` +SELECT + hour, + if(sum(cpu_usage_count) > 0, sum(cpu_usage_sum) / sum(cpu_usage_count), 0) AS average_cpu_usage_percent, + if(sum(memory_usage_count) > 0, sum(memory_usage_sum) / sum(memory_usage_count), 0) AS average_memory_usage_percent, + sum(greatest(network_rx_max - network_rx_min, 0)) AS network_rx_bytes, + sum(greatest(network_tx_max - network_tx_min, 0)) AS network_tx_bytes, + sum(greatest(disk_read_max - disk_read_min, 0)) AS disk_read_bytes, + sum(greatest(disk_write_max - disk_write_min, 0)) AS disk_write_bytes, + toUInt64(uniqExact(node_id)) AS reported_nodes +FROM %s +WHERE %s +GROUP BY hour +ORDER BY hour ASC`, nodeMetricCapacityHourlyTableName(), clause) + rows, err := conn.Query(ctx, sql, args...) + if err != nil { + return nil, fmt.Errorf("list node metric hourly from rollup: %w", err) + } + defer func() { _ = rows.Close() }() + return scanNodeMetricHourlyRows(rows) +} + +func listNodeMetricHourlyFromRaw(ctx context.Context, filter NodeObservabilityFilter) ([]NodeMetricHourly, error) { conn, err := observabilityConn() if err != nil { return nil, err @@ -372,7 +454,10 @@ ORDER BY hour ASC`, tableName, clause) return nil, fmt.Errorf("list node metric hourly: %w", err) } defer func() { _ = rows.Close() }() + return scanNodeMetricHourlyRows(rows) +} +func scanNodeMetricHourlyRows(rows driver.Rows) ([]NodeMetricHourly, error) { result := make([]NodeMetricHourly, 0) for rows.Next() { var ( @@ -407,7 +492,39 @@ ORDER BY hour ASC`, tableName, clause) } // ListNodeOpenrestyHourly returns hourly OpenResty observation aggregates matching filter. +// Prefers of_node_openresty_hourly rollups; falls back to raw lagInFrame on empty/error. func ListNodeOpenrestyHourly(ctx context.Context, filter NodeObservabilityFilter) ([]NodeOpenrestyHourly, error) { + if rows, err := listNodeOpenrestyHourlyFromRollup(ctx, filter); err == nil && len(rows) > 0 { + return rows, nil + } + return listNodeOpenrestyHourlyFromRaw(ctx, filter) +} + +func listNodeOpenrestyHourlyFromRollup(ctx context.Context, filter NodeObservabilityFilter) ([]NodeOpenrestyHourly, error) { + conn, err := observabilityConn() + if err != nil { + return nil, err + } + clause, args := buildNodeObservabilityFilterClause(filter, "hour") + sql := fmt.Sprintf(` +SELECT + hour, + sum(greatest(openresty_rx_max - openresty_rx_min, 0)) AS openresty_rx_bytes, + sum(greatest(openresty_tx_max - openresty_tx_min, 0)) AS openresty_tx_bytes, + toUInt64(uniqExact(node_id)) AS reported_nodes +FROM %s +WHERE %s +GROUP BY hour +ORDER BY hour ASC`, nodeOpenrestyHourlyTableName(), clause) + rows, err := conn.Query(ctx, sql, args...) + if err != nil { + return nil, fmt.Errorf("list node openresty hourly from rollup: %w", err) + } + defer func() { _ = rows.Close() }() + return scanNodeOpenrestyHourlyRows(rows) +} + +func listNodeOpenrestyHourlyFromRaw(ctx context.Context, filter NodeObservabilityFilter) ([]NodeOpenrestyHourly, error) { conn, err := observabilityConn() if err != nil { return nil, err @@ -442,7 +559,10 @@ ORDER BY hour ASC`, tableName, clause) return nil, fmt.Errorf("list node openresty hourly: %w", err) } defer func() { _ = rows.Close() }() + return scanNodeOpenrestyHourlyRows(rows) +} +func scanNodeOpenrestyHourlyRows(rows driver.Rows) ([]NodeOpenrestyHourly, error) { result := make([]NodeOpenrestyHourly, 0) for rows.Next() { var ( diff --git a/internal/repository/analytics/node_observability_delete.go b/internal/repository/analytics/node_observability_delete.go index 978d9427..b3387f68 100644 --- a/internal/repository/analytics/node_observability_delete.go +++ b/internal/repository/analytics/node_observability_delete.go @@ -9,7 +9,7 @@ import ( "time" ) -// DeleteAllNodeMetricSnapshots deletes all node metric snapshots. +// DeleteAllNodeMetricSnapshots hard-deletes all node metric snapshots via TRUNCATE. func DeleteAllNodeMetricSnapshots(ctx context.Context) (int64, error) { conn, err := observabilityConn() if err != nil { @@ -19,31 +19,39 @@ func DeleteAllNodeMetricSnapshots(ctx context.Context) (int64, error) { if err != nil { return 0, err } - return outcome.EligibleCount, nil + return outcome.DeletedCount, nil } -// DeleteNodeMetricSnapshotsBefore expires metric snapshots captured before cutoff via table TTL. -func DeleteNodeMetricSnapshotsBefore(ctx context.Context, cutoff time.Time) (int64, error) { - conn, err := observabilityConn() +// DeleteNodeMetricSnapshotsBefore force-materializes of_node_metric_snapshots table TTL. +// cutoff is ignored; see MaterializeNodeMetricSnapshotsTTL. +func DeleteNodeMetricSnapshotsBefore(ctx context.Context, _ time.Time) (int64, error) { + outcome, err := MaterializeNodeMetricSnapshotsTTL(ctx) if err != nil { return 0, err } + return outcome.EligibleCount, nil +} + +// MaterializeNodeMetricSnapshotsTTL force-materializes table TTL and reports an honest outcome. +func MaterializeNodeMetricSnapshotsTTL(ctx context.Context) (CleanupOutcome, error) { + conn, err := observabilityConn() + if err != nil { + return CleanupOutcome{}, err + } tableName := nodeMetricSnapshotTableName() - cutoff = cutoff.UTC() - outcome, err := expireRowsViaTTL( + ttlDays := TableTTLDaysNodeMetricSnapshots + cutoff := tableTTLCutoff(ttlDays, time.Now()) + return materializeExpiredByTableTTL( ctx, conn, tableName, + ttlDays, fmt.Sprintf("SELECT count() FROM %s WHERE captured_at < ?", tableName), []any{cutoff}, ) - if err != nil { - return 0, err - } - return outcome.EligibleCount, nil } -// DeleteAllNodeRequestReports deletes all node request reports. +// DeleteAllNodeRequestReports hard-deletes all node request reports via TRUNCATE. func DeleteAllNodeRequestReports(ctx context.Context) (int64, error) { conn, err := observabilityConn() if err != nil { @@ -53,31 +61,39 @@ func DeleteAllNodeRequestReports(ctx context.Context) (int64, error) { if err != nil { return 0, err } - return outcome.EligibleCount, nil + return outcome.DeletedCount, nil } -// DeleteNodeRequestReportsBefore expires request reports ending before cutoff via table TTL. -func DeleteNodeRequestReportsBefore(ctx context.Context, cutoff time.Time) (int64, error) { - conn, err := observabilityConn() +// DeleteNodeRequestReportsBefore force-materializes of_node_request_reports table TTL. +// cutoff is ignored; see MaterializeNodeRequestReportsTTL. +func DeleteNodeRequestReportsBefore(ctx context.Context, _ time.Time) (int64, error) { + outcome, err := MaterializeNodeRequestReportsTTL(ctx) if err != nil { return 0, err } + return outcome.EligibleCount, nil +} + +// MaterializeNodeRequestReportsTTL force-materializes table TTL and reports an honest outcome. +func MaterializeNodeRequestReportsTTL(ctx context.Context) (CleanupOutcome, error) { + conn, err := observabilityConn() + if err != nil { + return CleanupOutcome{}, err + } tableName := nodeRequestReportTableName() - cutoff = cutoff.UTC() - outcome, err := expireRowsViaTTL( + ttlDays := TableTTLDaysNodeRequestReports + cutoff := tableTTLCutoff(ttlDays, time.Now()) + return materializeExpiredByTableTTL( ctx, conn, tableName, + ttlDays, fmt.Sprintf("SELECT count() FROM %s WHERE window_ended_at < ?", tableName), []any{cutoff}, ) - if err != nil { - return 0, err - } - return outcome.EligibleCount, nil } -// DeleteAllNodeObsOpenresty deletes all OpenResty observations. +// DeleteAllNodeObsOpenresty hard-deletes all OpenResty observations via TRUNCATE. func DeleteAllNodeObsOpenresty(ctx context.Context) (int64, error) { conn, err := observabilityConn() if err != nil { @@ -87,31 +103,39 @@ func DeleteAllNodeObsOpenresty(ctx context.Context) (int64, error) { if err != nil { return 0, err } - return outcome.EligibleCount, nil + return outcome.DeletedCount, nil } -// DeleteNodeObsOpenrestyBefore expires OpenResty observations captured before cutoff via table TTL. -func DeleteNodeObsOpenrestyBefore(ctx context.Context, cutoff time.Time) (int64, error) { - conn, err := observabilityConn() +// DeleteNodeObsOpenrestyBefore force-materializes of_node_obs_openresty table TTL. +// cutoff is ignored; see MaterializeNodeObsOpenrestyTTL. +func DeleteNodeObsOpenrestyBefore(ctx context.Context, _ time.Time) (int64, error) { + outcome, err := MaterializeNodeObsOpenrestyTTL(ctx) if err != nil { return 0, err } + return outcome.EligibleCount, nil +} + +// MaterializeNodeObsOpenrestyTTL force-materializes table TTL and reports an honest outcome. +func MaterializeNodeObsOpenrestyTTL(ctx context.Context) (CleanupOutcome, error) { + conn, err := observabilityConn() + if err != nil { + return CleanupOutcome{}, err + } tableName := nodeObsOpenrestyTableName() - cutoff = cutoff.UTC() - outcome, err := expireRowsViaTTL( + ttlDays := TableTTLDaysNodeObs + cutoff := tableTTLCutoff(ttlDays, time.Now()) + return materializeExpiredByTableTTL( ctx, conn, tableName, + ttlDays, fmt.Sprintf("SELECT count() FROM %s WHERE captured_at < ?", tableName), []any{cutoff}, ) - if err != nil { - return 0, err - } - return outcome.EligibleCount, nil } -// DeleteAllNodeObsFrps deletes all FRPS observations. +// DeleteAllNodeObsFrps hard-deletes all FRPS observations via TRUNCATE. func DeleteAllNodeObsFrps(ctx context.Context) (int64, error) { conn, err := observabilityConn() if err != nil { @@ -121,31 +145,39 @@ func DeleteAllNodeObsFrps(ctx context.Context) (int64, error) { if err != nil { return 0, err } - return outcome.EligibleCount, nil + return outcome.DeletedCount, nil } -// DeleteNodeObsFrpsBefore expires FRPS observations captured before cutoff via table TTL. -func DeleteNodeObsFrpsBefore(ctx context.Context, cutoff time.Time) (int64, error) { - conn, err := observabilityConn() +// DeleteNodeObsFrpsBefore force-materializes of_node_obs_frps table TTL. +// cutoff is ignored; see MaterializeNodeObsFrpsTTL. +func DeleteNodeObsFrpsBefore(ctx context.Context, _ time.Time) (int64, error) { + outcome, err := MaterializeNodeObsFrpsTTL(ctx) if err != nil { return 0, err } + return outcome.EligibleCount, nil +} + +// MaterializeNodeObsFrpsTTL force-materializes table TTL and reports an honest outcome. +func MaterializeNodeObsFrpsTTL(ctx context.Context) (CleanupOutcome, error) { + conn, err := observabilityConn() + if err != nil { + return CleanupOutcome{}, err + } tableName := nodeObsFrpsTableName() - cutoff = cutoff.UTC() - outcome, err := expireRowsViaTTL( + ttlDays := TableTTLDaysNodeObs + cutoff := tableTTLCutoff(ttlDays, time.Now()) + return materializeExpiredByTableTTL( ctx, conn, tableName, + ttlDays, fmt.Sprintf("SELECT count() FROM %s WHERE captured_at < ?", tableName), []any{cutoff}, ) - if err != nil { - return 0, err - } - return outcome.EligibleCount, nil } -// DeleteAllNodeObsFrpc deletes all FRPC observations. +// DeleteAllNodeObsFrpc hard-deletes all FRPC observations via TRUNCATE. func DeleteAllNodeObsFrpc(ctx context.Context) (int64, error) { conn, err := observabilityConn() if err != nil { @@ -155,26 +187,34 @@ func DeleteAllNodeObsFrpc(ctx context.Context) (int64, error) { if err != nil { return 0, err } + return outcome.DeletedCount, nil +} + +// DeleteNodeObsFrpcBefore force-materializes of_node_obs_frpc table TTL. +// cutoff is ignored; see MaterializeNodeObsFrpcTTL. +func DeleteNodeObsFrpcBefore(ctx context.Context, _ time.Time) (int64, error) { + outcome, err := MaterializeNodeObsFrpcTTL(ctx) + if err != nil { + return 0, err + } return outcome.EligibleCount, nil } -// DeleteNodeObsFrpcBefore expires FRPC observations captured before cutoff via table TTL. -func DeleteNodeObsFrpcBefore(ctx context.Context, cutoff time.Time) (int64, error) { +// MaterializeNodeObsFrpcTTL force-materializes table TTL and reports an honest outcome. +func MaterializeNodeObsFrpcTTL(ctx context.Context) (CleanupOutcome, error) { conn, err := observabilityConn() if err != nil { - return 0, err + return CleanupOutcome{}, err } tableName := nodeObsFrpcTableName() - cutoff = cutoff.UTC() - outcome, err := expireRowsViaTTL( + ttlDays := TableTTLDaysNodeObs + cutoff := tableTTLCutoff(ttlDays, time.Now()) + return materializeExpiredByTableTTL( ctx, conn, tableName, + ttlDays, fmt.Sprintf("SELECT count() FROM %s WHERE captured_at < ?", tableName), []any{cutoff}, ) - if err != nil { - return 0, err - } - return outcome.EligibleCount, nil -} \ No newline at end of file +} diff --git a/internal/repository/analytics/node_observability_filter.go b/internal/repository/analytics/node_observability_filter.go index 6f0c674f..bb8e6a15 100644 --- a/internal/repository/analytics/node_observability_filter.go +++ b/internal/repository/analytics/node_observability_filter.go @@ -61,3 +61,14 @@ func nodeObsFrpsTableName() string { func nodeObsFrpcTableName() string { return "of_node_obs_frpc" } + +func nodeMetricCapacityHourlyTableName() string { + return "of_node_metric_capacity_hourly" +} + +func nodeOpenrestyHourlyTableName() string { + return "of_node_openresty_hourly" +} + +// clickHouseLimit1ByNodeIDClause selects the first row per node_id after ORDER BY. +const clickHouseLimit1ByNodeIDClause = " LIMIT 1 BY node_id" diff --git a/internal/repository/analytics/node_observability_latest_test.go b/internal/repository/analytics/node_observability_latest_test.go new file mode 100644 index 00000000..8478e4fb --- /dev/null +++ b/internal/repository/analytics/node_observability_latest_test.go @@ -0,0 +1,166 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package analytics + +import ( + "context" + "errors" + "strings" + "testing" + "time" + + "github.com/ClickHouse/clickhouse-go/v2/lib/driver" + "github.com/Rain-kl/Wavelet/internal/db" + "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" +) + +func TestListLatestNodeMetricSnapshots_UsesLimit1ByNodeID(t *testing.T) { + ctx := context.Background() + mock := &mockConn{} + db.SetChConnForTest(mock) + t.Cleanup(func() { db.SetChConnForTest(nil) }) + + since := time.Date(2026, 7, 10, 0, 0, 0, 0, time.UTC) + _, err := ListLatestNodeMetricSnapshots(ctx, NodeObservabilityFilter{Since: since}) + require.NoError(t, err) + require.Len(t, mock.queries, 1) + assert.Contains(t, mock.queries[0], "LIMIT 1 BY node_id") + assert.Contains(t, mock.queries[0], nodeMetricSnapshotTableName()) + assert.Contains(t, mock.queries[0], "captured_at DESC") + assert.NotContains(t, mock.queries[0], "LIMIT ?") + require.Len(t, mock.queryArgs, 1) + require.Len(t, mock.queryArgs[0], 1) + assert.Equal(t, since, mock.queryArgs[0][0]) +} + +func TestListLatestNodeRequestReports_UsesLimit1ByNodeID(t *testing.T) { + ctx := context.Background() + mock := &mockConn{} + db.SetChConnForTest(mock) + t.Cleanup(func() { db.SetChConnForTest(nil) }) + + _, err := ListLatestNodeRequestReports(ctx, NodeObservabilityFilter{NodeID: "node-a"}) + require.NoError(t, err) + require.Len(t, mock.queries, 1) + assert.Contains(t, mock.queries[0], "LIMIT 1 BY node_id") + assert.Contains(t, mock.queries[0], nodeRequestReportTableName()) + assert.Contains(t, mock.queries[0], "window_ended_at DESC") + require.Len(t, mock.queryArgs, 1) + require.Len(t, mock.queryArgs[0], 1) + assert.Equal(t, "node-a", mock.queryArgs[0][0]) +} + +func TestListNodeMetricHourly_PrefersRollup(t *testing.T) { + ctx := context.Background() + hour := time.Date(2026, 7, 10, 12, 0, 0, 0, time.UTC) + mock := &mockConn{ + queryFn: func(_ context.Context, query string, _ ...any) (driver.Rows, error) { + if strings.Contains(query, nodeMetricCapacityHourlyTableName()) { + return &mockRows{data: [][]any{{ + hour, 42.5, 60.0, int64(100), int64(200), int64(10), int64(20), uint64(2), + }}}, nil + } + return nil, errors.New("raw path should not be used when rollup has rows") + }, + } + db.SetChConnForTest(mock) + t.Cleanup(func() { db.SetChConnForTest(nil) }) + + rows, err := ListNodeMetricHourly(ctx, NodeObservabilityFilter{}) + require.NoError(t, err) + require.Len(t, rows, 1) + assert.Equal(t, 42.5, rows[0].AverageCPUUsagePercent) + assert.Equal(t, 60.0, rows[0].AverageMemoryUsagePercent) + assert.Equal(t, int64(100), rows[0].NetworkRxBytes) + assert.Equal(t, 2, rows[0].ReportedNodes) + require.Len(t, mock.queries, 1) + assert.Contains(t, mock.queries[0], nodeMetricCapacityHourlyTableName()) +} + +func TestListNodeMetricHourly_FallsBackToRawOnRollupError(t *testing.T) { + ctx := context.Background() + hour := time.Date(2026, 7, 10, 13, 0, 0, 0, time.UTC) + mock := &mockConn{ + queryFn: func(_ context.Context, query string, _ ...any) (driver.Rows, error) { + if strings.Contains(query, nodeMetricCapacityHourlyTableName()) { + return nil, errors.New("rollup missing") + } + if strings.Contains(query, nodeMetricSnapshotTableName()) { + return &mockRows{data: [][]any{{ + hour, 10.0, 20.0, int64(1), int64(2), int64(3), int64(4), uint64(1), + }}}, nil + } + return &mockRows{}, nil + }, + } + db.SetChConnForTest(mock) + t.Cleanup(func() { db.SetChConnForTest(nil) }) + + rows, err := ListNodeMetricHourly(ctx, NodeObservabilityFilter{}) + require.NoError(t, err) + require.Len(t, rows, 1) + assert.Equal(t, 10.0, rows[0].AverageCPUUsagePercent) + assert.Equal(t, int64(3), rows[0].DiskReadBytes) + require.GreaterOrEqual(t, len(mock.queries), 2) + assert.Contains(t, mock.queries[0], nodeMetricCapacityHourlyTableName()) + assert.Contains(t, mock.queries[1], "lagInFrame") +} + +func TestListNodeOpenrestyHourly_PrefersRollup(t *testing.T) { + ctx := context.Background() + hour := time.Date(2026, 7, 10, 14, 0, 0, 0, time.UTC) + mock := &mockConn{ + queryFn: func(_ context.Context, query string, _ ...any) (driver.Rows, error) { + if strings.Contains(query, nodeOpenrestyHourlyTableName()) { + return &mockRows{data: [][]any{{ + hour, int64(50), int64(70), uint64(3), + }}}, nil + } + return nil, errors.New("raw path should not be used when rollup has rows") + }, + } + db.SetChConnForTest(mock) + t.Cleanup(func() { db.SetChConnForTest(nil) }) + + rows, err := ListNodeOpenrestyHourly(ctx, NodeObservabilityFilter{}) + require.NoError(t, err) + require.Len(t, rows, 1) + assert.Equal(t, int64(50), rows[0].OpenrestyRxBytes) + assert.Equal(t, int64(70), rows[0].OpenrestyTxBytes) + assert.Equal(t, 3, rows[0].ReportedNodes) +} + +func TestListNodeOpenrestyHourly_FallsBackToRawOnEmptyRollup(t *testing.T) { + ctx := context.Background() + hour := time.Date(2026, 7, 10, 15, 0, 0, 0, time.UTC) + mock := &mockConn{ + queryFn: func(_ context.Context, query string, _ ...any) (driver.Rows, error) { + if strings.Contains(query, nodeOpenrestyHourlyTableName()) { + return &mockRows{}, nil + } + if strings.Contains(query, nodeObsOpenrestyTableName()) { + return &mockRows{data: [][]any{{ + hour, int64(9), int64(8), uint64(1), + }}}, nil + } + return &mockRows{}, nil + }, + } + db.SetChConnForTest(mock) + t.Cleanup(func() { db.SetChConnForTest(nil) }) + + rows, err := ListNodeOpenrestyHourly(ctx, NodeObservabilityFilter{}) + require.NoError(t, err) + require.Len(t, rows, 1) + assert.Equal(t, int64(9), rows[0].OpenrestyRxBytes) + require.GreaterOrEqual(t, len(mock.queries), 2) + assert.Contains(t, mock.queries[1], "lagInFrame") +} + +func TestNonNegativeCounterRange(t *testing.T) { + assert.Equal(t, int64(0), nonNegativeCounterRange(10, 10)) + assert.Equal(t, int64(0), nonNegativeCounterRange(5, 10)) + assert.Equal(t, int64(15), nonNegativeCounterRange(25, 10)) +}