feat(zone): add Cloudflare-style traffic overview charts

Expose Zone stats API with multi-host access-log aggregates and time
series, and render unique visitors, requests and data served on the
Zone overview with 24h/7d/30d range controls.
This commit is contained in:
ryan
2026-07-12 16:38:35 +08:00
parent efcf61e32d
commit 31b4886a14
20 changed files with 1026 additions and 32 deletions
+1
View File
@@ -15,4 +15,5 @@ const (
errCertificateNotFound = "所选证书不存在"
errDomainBoundToRoute = "域名已绑定反代路由,请先解除绑定"
errZoneHasDomains = "根域下仍有域名,请先删除全部域名"
errStatsRangeInvalid = "时间范围无效,请选择 24h、7d 或 30d"
)
@@ -6,6 +6,7 @@ package zone
import (
"context"
"testing"
"time"
"github.com/Rain-kl/Wavelet/internal/db"
"github.com/Rain-kl/Wavelet/internal/model"
@@ -55,3 +56,49 @@ func TestLegacyImportUsesEffectiveTLDPlusOne(t *testing.T) {
require.NoError(t, err)
require.Equal(t, "example.co.uk", root)
}
func TestGetStatsAggregatesZoneHosts(t *testing.T) {
ctx := setupZoneDB(t)
reset := model.SetAccessLogStoreForTest(model.NewMemoryAccessLogStore())
t.Cleanup(reset)
zone, err := Create(ctx, Input{Domain: "example.com"})
require.NoError(t, err)
_, err = CreateDomain(ctx, zone.ID, DomainInput{Domain: "api.example.com"})
require.NoError(t, err)
_, err = CreateDomain(ctx, zone.ID, DomainInput{Domain: "www.example.com"})
require.NoError(t, err)
now := time.Now().UTC()
require.NoError(t, model.InsertOpenFlareAccessLogsBatch(ctx, []*model.OpenFlareAccessLog{
{NodeID: "n1", LoggedAt: now.Add(-1 * time.Hour), RemoteAddr: "1.1.1.1", Host: "api.example.com", Path: "/", StatusCode: 200},
{NodeID: "n1", LoggedAt: now.Add(-2 * time.Hour), RemoteAddr: "1.1.1.1", Host: "www.example.com", Path: "/", StatusCode: 200},
{NodeID: "n1", LoggedAt: now.Add(-3 * time.Hour), RemoteAddr: "2.2.2.2", Host: "api.example.com", Path: "/x", StatusCode: 404},
{NodeID: "n1", LoggedAt: now.Add(-3 * time.Hour), RemoteAddr: "3.3.3.3", Host: "other.com", Path: "/", StatusCode: 200},
{NodeID: "n1", LoggedAt: now.Add(-48 * time.Hour), RemoteAddr: "4.4.4.4", Host: "api.example.com", Path: "/", StatusCode: 200},
}))
stats, err := GetStats(ctx, zone.ID, "24h")
require.NoError(t, err)
require.Equal(t, StatsRange24h, stats.Range)
require.Equal(t, int64(3), stats.RequestCount)
require.Equal(t, int64(2), stats.UniqueVisitors)
require.Equal(t, 2, stats.DomainCount)
require.True(t, stats.Available)
require.NotEmpty(t, stats.Series)
require.Equal(t, 60, stats.BucketMinutes)
var seriesRequests int64
for _, point := range stats.Series {
seriesRequests += point.RequestCount
}
require.Equal(t, int64(3), seriesRequests)
stats7d, err := GetStats(ctx, zone.ID, "7d")
require.NoError(t, err)
require.Equal(t, int64(4), stats7d.RequestCount)
require.Equal(t, int64(3), stats7d.UniqueVisitors)
require.NotEmpty(t, stats7d.Series)
_, err = GetStats(ctx, zone.ID, "1h")
require.EqualError(t, err, errStatsRangeInvalid)
}
+24
View File
@@ -87,6 +87,30 @@ func GetOverviewHandler(c *gin.Context) {
c.JSON(http.StatusOK, response.OK(item))
}
// GetStatsHandler returns Zone traffic metrics for a time range.
// @Summary 获取 Zone 流量统计
// @Description 按 Zone 下全部域名聚合访问日志:唯一访问者、请求总数、已提供数据(字节)。range 支持 24h/7d/30d。
// @Tags openflare-zone
// @Produce json
// @Security SessionCookie
// @Param id path int true "Zone ID"
// @Param range query string false "时间范围:24h(默认)、7d、30d"
// @Success 200 {object} response.Any{data=zone.Stats}
// @Failure 400 {object} response.Any
// @Failure 404 {object} response.Any
// @Router /api/v1/d/zones/{id}/stats [get]
func GetStatsHandler(c *gin.Context) {
id, ok := apiutil.IDParam(c)
if !ok {
return
}
item, err := GetStats(c.Request.Context(), id, c.Query("range"))
if abort(c, err, errZoneNotFound) {
return
}
c.JSON(http.StatusOK, response.OK(item))
}
// UpdateHandler updates a Zone.
// @Summary 更新 Zone
// @Tags openflare-zone
+199
View File
@@ -0,0 +1,199 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package zone
import (
"context"
"errors"
"strings"
"time"
"github.com/Rain-kl/Wavelet/internal/db"
"github.com/Rain-kl/Wavelet/internal/model"
"gorm.io/gorm"
)
// StatsRange is a supported traffic window for Zone analytics.
type StatsRange string
const (
StatsRange24h StatsRange = "24h"
StatsRange7d StatsRange = "7d"
StatsRange30d StatsRange = "30d"
)
// StatsPoint is one bucket on a Zone traffic chart.
type StatsPoint struct {
BucketStartedAt time.Time `json:"bucket_started_at"`
RequestCount int64 `json:"request_count"`
UniqueVisitors int64 `json:"unique_visitors"`
BytesSent int64 `json:"bytes_sent"`
}
// Stats summarizes edge traffic for all domains under a Zone.
type Stats struct {
Range StatsRange `json:"range"`
RangeHours int `json:"range_hours"`
WindowStartedAt time.Time `json:"window_started_at"`
WindowEndedAt time.Time `json:"window_ended_at"`
BucketMinutes int `json:"bucket_minutes"`
UniqueVisitors int64 `json:"unique_visitors"`
RequestCount int64 `json:"request_count"`
BytesSent int64 `json:"bytes_sent"`
DomainCount int `json:"domain_count"`
Available bool `json:"available"`
Series []StatsPoint `json:"series"`
}
func parseStatsRange(raw string) (StatsRange, time.Duration, int, error) {
switch StatsRange(strings.TrimSpace(raw)) {
case "", StatsRange24h:
return StatsRange24h, 24 * time.Hour, 60, nil
case StatsRange7d:
return StatsRange7d, 7 * 24 * time.Hour, 6 * 60, nil
case StatsRange30d:
return StatsRange30d, 30 * 24 * time.Hour, 24 * 60, nil
default:
return "", 0, 0, errors.New(errStatsRangeInvalid)
}
}
// GetStats aggregates access-log traffic for a Zone over a time range.
func GetStats(ctx context.Context, id uint, rangeRaw string) (*Stats, error) {
statsRange, window, bucketMinutes, err := parseStatsRange(rangeRaw)
if err != nil {
return nil, err
}
var zone model.Zone
if err := db.DB(ctx).First(&zone, id).Error; err != nil {
return nil, err
}
var domains []model.ZoneDomain
if err := db.DB(ctx).Where("zone_id = ?", id).Order("domain asc").Find(&domains).Error; err != nil {
return nil, err
}
now := time.Now().UTC().Truncate(time.Minute)
since := now.Add(-window)
// Align chart window start to bucket boundary for cleaner x-axis labels.
bucket := time.Duration(bucketMinutes) * time.Minute
since = since.Truncate(bucket)
result := &Stats{
Range: statsRange,
RangeHours: int(window / time.Hour),
WindowStartedAt: since,
WindowEndedAt: now,
BucketMinutes: bucketMinutes,
DomainCount: len(domains),
Available: true,
Series: emptyStatsSeries(since, now, bucketMinutes),
}
if len(domains) == 0 {
return result, nil
}
hosts := make([]string, 0, len(domains))
for _, domain := range domains {
if host := strings.TrimSpace(domain.Domain); host != "" {
hosts = append(hosts, host)
}
}
if len(hosts) == 0 {
return result, nil
}
requestCount, uniqueVisitors, err := model.CountOpenFlareAccessLogs(ctx, model.OpenFlareAccessLogQuery{
Hosts: hosts,
Since: since,
Until: now,
})
if err != nil {
if isAnalyticsUnavailable(err) {
result.Available = false
return result, nil
}
return nil, err
}
result.RequestCount = requestCount
result.UniqueVisitors = uniqueVisitors
// Bytes are not yet persisted on edge access logs; keep the field for UI compatibility.
result.BytesSent = 0
buckets, err := model.ListOpenFlareAccessLogBuckets(ctx, model.OpenFlareAccessLogBucketQuery{
Hosts: hosts,
Since: since,
Until: now,
FoldMinutes: bucketMinutes,
SortBy: "logged_at",
SortOrder: "asc",
})
if err != nil {
if isAnalyticsUnavailable(err) {
result.Available = false
return result, nil
}
return nil, err
}
byEpoch := make(map[int64]model.OpenFlareAccessLogBucketRow, len(buckets))
for _, bucketRow := range buckets {
if bucketRow == nil {
continue
}
byEpoch[bucketRow.BucketEpoch] = *bucketRow
}
series := emptyStatsSeries(since, now, bucketMinutes)
for index := range series {
epoch := series[index].BucketStartedAt.Unix()
if row, ok := byEpoch[epoch]; ok {
series[index].RequestCount = row.RequestCount
series[index].UniqueVisitors = row.UniqueIPCount
series[index].BytesSent = 0
}
}
result.Series = series
return result, nil
}
func emptyStatsSeries(since, until time.Time, bucketMinutes int) []StatsPoint {
if bucketMinutes <= 0 {
bucketMinutes = 60
}
bucket := time.Duration(bucketMinutes) * time.Minute
start := since.UTC().Truncate(bucket)
end := until.UTC()
if !end.After(start) {
return []StatsPoint{{BucketStartedAt: start}}
}
// Cap points to keep chart readable.
maxPoints := 120
capacity := int(end.Sub(start)/bucket) + 1
if capacity > maxPoints {
capacity = maxPoints
}
points := make([]StatsPoint, 0, capacity)
for cursor := start; !cursor.After(end) && len(points) < maxPoints; cursor = cursor.Add(bucket) {
points = append(points, StatsPoint{BucketStartedAt: cursor})
}
if len(points) == 0 {
points = append(points, StatsPoint{BucketStartedAt: start})
}
return points
}
func isAnalyticsUnavailable(err error) bool {
if err == nil {
return false
}
if errors.Is(err, gorm.ErrInvalidDB) {
return true
}
msg := strings.ToLower(err.Error())
return strings.Contains(msg, "clickhouse connection is not initialized") ||
strings.Contains(msg, "clickhouse is not") ||
strings.Contains(msg, "database is not initialized")
}
+2
View File
@@ -309,8 +309,10 @@ func openFlareAccessLogQueryFromBucket(query OpenFlareAccessLogBucketQuery) Open
NodeID: query.NodeID,
RemoteAddr: query.RemoteAddr,
Host: query.Host,
Hosts: query.Hosts,
Path: query.Path,
Since: query.Since,
Until: query.Until,
Page: query.Page,
PageSize: query.PageSize,
SortBy: query.SortBy,
@@ -263,6 +263,7 @@ func toNodeAccessLogFilter(query OpenFlareAccessLogQuery) analyticsrepo.NodeAcce
NodeID: query.NodeID,
RemoteAddr: query.RemoteAddr,
Host: query.Host,
Hosts: query.Hosts,
Path: query.Path,
Since: query.Since,
Until: query.Until,
@@ -450,7 +450,19 @@ func memoryAccessLogMatches(row *OpenFlareAccessLog, query OpenFlareAccessLogQue
if trimmed := strings.TrimSpace(query.RemoteAddr); trimmed != "" && !strings.HasPrefix(strings.TrimSpace(row.RemoteAddr), trimmed) {
return false
}
if trimmed := strings.TrimSpace(query.Host); trimmed != "" && !strings.HasPrefix(strings.TrimSpace(row.Host), trimmed) {
if len(query.Hosts) > 0 {
rowHost := strings.ToLower(strings.TrimSpace(row.Host))
matched := false
for _, host := range query.Hosts {
if strings.ToLower(strings.TrimSpace(host)) == rowHost {
matched = true
break
}
}
if !matched {
return false
}
} else if trimmed := strings.TrimSpace(query.Host); trimmed != "" && !strings.HasPrefix(strings.TrimSpace(row.Host), trimmed) {
return false
}
if trimmed := strings.TrimSpace(query.Path); trimmed != "" && !strings.HasPrefix(strings.TrimSpace(row.Path), trimmed) {
+11 -7
View File
@@ -185,13 +185,15 @@ type OpenFlareAccessLogQuery struct {
NodeID string
RemoteAddr string
Host string
Path string
Since time.Time
Until time.Time
Page int
PageSize int
SortBy string
SortOrder string
// Hosts exact-matches any host (case-insensitive). Prefer over Host for multi-domain scopes.
Hosts []string
Path string
Since time.Time
Until time.Time
Page int
PageSize int
SortBy string
SortOrder string
}
// OpenFlareAccessLogBucketQuery filters folded access log queries (v1 stub).
@@ -199,8 +201,10 @@ type OpenFlareAccessLogBucketQuery struct {
NodeID string
RemoteAddr string
Host string
Hosts []string
Path string
Since time.Time
Until time.Time
Page int
PageSize int
SortBy string
@@ -12,8 +12,8 @@ import (
const (
nodeAccessLogFilterClauseCapacity = 6
nodeAccessLogSortDesc = "DESC"
nodeAccessLogSortAsc = "ASC"
nodeAccessLogSortDesc = "DESC"
nodeAccessLogSortAsc = "ASC"
nodeAccessLogSortAscInput = "asc"
nodeAccessLogColumnRemoteAddr = "remote_addr"
@@ -24,13 +24,15 @@ type NodeAccessLogFilter struct {
NodeID string
RemoteAddr string
Host string
Path string
Since time.Time
Until time.Time
Page int
PageSize int
SortBy string
SortOrder string
// Hosts exact-matches any host (case-insensitive). Prefer over Host for multi-domain scopes.
Hosts []string
Path string
Since time.Time
Until time.Time
Page int
PageSize int
SortBy string
SortOrder string
}
func buildNodeAccessLogFilterClause(filter NodeAccessLogFilter) (string, []any) {
@@ -44,7 +46,15 @@ func buildNodeAccessLogFilterClause(filter NodeAccessLogFilter) (string, []any)
parts = append(parts, "remote_addr LIKE ?")
args = append(args, trimmed+"%")
}
if trimmed := strings.TrimSpace(filter.Host); trimmed != "" {
hosts := normalizeNodeAccessLogHosts(filter.Hosts)
if len(hosts) > 0 {
placeholders := make([]string, 0, len(hosts))
for _, host := range hosts {
placeholders = append(placeholders, "?")
args = append(args, host)
}
parts = append(parts, "lowerUTF8(trim(host)) IN ("+strings.Join(placeholders, ", ")+")")
} else if trimmed := strings.TrimSpace(filter.Host); trimmed != "" {
parts = append(parts, "host LIKE ?")
args = append(args, trimmed+"%")
}
@@ -99,6 +109,26 @@ func normalizeNodeAccessLogRemoteAddr(value string) string {
return strings.TrimSpace(value)
}
func normalizeNodeAccessLogHosts(hosts []string) []string {
if len(hosts) == 0 {
return nil
}
seen := make(map[string]struct{}, len(hosts))
result := make([]string, 0, len(hosts))
for _, host := range hosts {
trimmed := strings.ToLower(strings.TrimSpace(host))
if trimmed == "" {
continue
}
if _, ok := seen[trimmed]; ok {
continue
}
seen[trimmed] = struct{}{}
result = append(result, trimmed)
}
return result
}
func normalizeNodeAccessLogSortOrder(sortOrder string) string {
if strings.EqualFold(strings.TrimSpace(sortOrder), "asc") {
return "asc"
@@ -15,6 +15,7 @@ func registerZoneRoutes(apiGroup *gin.RouterGroup) {
apiutil.RegisterCollection(zoneGroup, "GET", zone.ListHandler)
apiutil.RegisterCollection(zoneGroup, "POST", zone.CreateHandler)
zoneGroup.GET("/:id/overview", zone.GetOverviewHandler)
zoneGroup.GET("/:id/stats", zone.GetStatsHandler)
zoneGroup.POST("/:id/update", zone.UpdateHandler)
zoneGroup.POST("/:id/delete", zone.DeleteHandler)
zoneGroup.POST("/:id/domains", zone.CreateDomainHandler)