diff --git a/docs/changelog/index.md b/docs/changelog/index.md index d9cca519..d309384c 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -23,15 +23,17 @@ sidebar: false ### 新增 +- 边缘访问日志 `of_node_access_logs` ClickHouse 表新增 `bytes_sent` 字段并打通 Agent 到 Server 的上报和 Zone 统计显示逻辑。 - Zone 概览页新增 Cloudflare 风格流量图:支持 24 小时 / 7 天 / 30 天,展示唯一访问者、请求总数与已提供数据趋势。 - 新增第一阶段 Zone 与正规化 Zone 域名数据库表及路由绑定模型,为后续以稳定 ID 管理网站与域名关联提供基础。 - 新增 Zone 管理 API 与显式历史域名导入命令,使用公共后缀列表验证注册根域和域名归属。 -- 管理端网站入口改为 Zone 列表与 `/websites/:zoneId` 详情(概览 / 域名 / 路由 / 证书 / 设置),反代路由通过 Zone 域名选择器绑定。 -- 新增 Zone 域名迁移指南(`docs/guide/zone-domain-migration.md`);历史域名在 Server 启动时由 goose 自动导入,无需单独命令。 ### 修改 -- 配置快照、OpenResty 渲染、Tunnel 与 Uptime Kuma 监控改为从 Zone 域名绑定读取域名和证书,移除对反代路由旧域名/证书字段的运行时回退。 +- 重构并清理了分析仓统计层 `node_access_log_stats.go` 和 `openflare_access_log.go` 之间的重复模型,使用底层 type aliases 简化了类型转换和拷贝逻辑。 +- 配置快照、OpenResty 渲染、Tunnel 与 Uptime Kuma 监控改为从 Zone 域名绑定读取域名 and 证书,移除对反代路由旧域名/证书字段的运行时回退。 +- 管理端网站入口改为 Zone 列表与 `/websites/:zoneId` 详情(概览 / 域名 / 路由 / 证书 / 设置),反代路由通过 Zone 域名选择器绑定。 +- 新增 Zone 域名迁移指南(`docs/guide/zone-domain-migration.md`);历史域名在 Server 启动时由 goose 自动导入,无需单独命令。 - 反代路由 API 以 `zone_domain_ids` / `zone_domains` 为唯一域名与证书关联来源,不再接受或返回路由内嵌域名/证书字段。 - 第二阶段迁移删除 `of_managed_domains` 及 `of_proxy_routes` 上的 `domain` / `domains` / `cert_id` / `cert_ids` / `domain_cert_ids` 冗余列。 diff --git a/docs/docs.go b/docs/docs.go index 02aa09f3..b3ecc49f 100644 --- a/docs/docs.go +++ b/docs/docs.go @@ -14041,6 +14041,9 @@ const docTemplate = `{ "github_com_Rain-kl_Wavelet_pkg_protocol.NodeAccessLog": { "type": "object", "properties": { + "bytes_sent": { + "type": "integer" + }, "host": { "type": "string" }, @@ -19023,15 +19026,15 @@ const docTemplate = `{ "type": "object", "properties": { "available": { - "description": "Available is false when analytics storage is unavailable (e.g. ClickHouse off).", "type": "boolean" }, + "bucket_minutes": { + "type": "integer" + }, "bytes_sent": { - "description": "BytesSent is total response bytes. Currently always 0 until access logs store body_bytes_sent.", "type": "integer" }, "domain_count": { - "description": "DomainCount is the number of explicit Zone domains used to scope hosts.", "type": "integer" }, "range": { @@ -19041,11 +19044,15 @@ const docTemplate = `{ "type": "integer" }, "request_count": { - "description": "RequestCount is total access-log requests for Zone domains.", "type": "integer" }, + "series": { + "type": "array", + "items": { + "$ref": "#/definitions/zone.StatsPoint" + } + }, "unique_visitors": { - "description": "UniqueVisitors is distinct client IPs (remote_addr) over the window.", "type": "integer" }, "window_ended_at": { @@ -19056,6 +19063,23 @@ const docTemplate = `{ } } }, + "zone.StatsPoint": { + "type": "object", + "properties": { + "bucket_started_at": { + "type": "string" + }, + "bytes_sent": { + "type": "integer" + }, + "request_count": { + "type": "integer" + }, + "unique_visitors": { + "type": "integer" + } + } + }, "zone.StatsRange": { "type": "string", "enum": [ diff --git a/docs/swagger.json b/docs/swagger.json index 3a7389e2..5bd7ede5 100644 --- a/docs/swagger.json +++ b/docs/swagger.json @@ -14034,6 +14034,9 @@ "github_com_Rain-kl_Wavelet_pkg_protocol.NodeAccessLog": { "type": "object", "properties": { + "bytes_sent": { + "type": "integer" + }, "host": { "type": "string" }, @@ -19016,15 +19019,15 @@ "type": "object", "properties": { "available": { - "description": "Available is false when analytics storage is unavailable (e.g. ClickHouse off).", "type": "boolean" }, + "bucket_minutes": { + "type": "integer" + }, "bytes_sent": { - "description": "BytesSent is total response bytes. Currently always 0 until access logs store body_bytes_sent.", "type": "integer" }, "domain_count": { - "description": "DomainCount is the number of explicit Zone domains used to scope hosts.", "type": "integer" }, "range": { @@ -19034,11 +19037,15 @@ "type": "integer" }, "request_count": { - "description": "RequestCount is total access-log requests for Zone domains.", "type": "integer" }, + "series": { + "type": "array", + "items": { + "$ref": "#/definitions/zone.StatsPoint" + } + }, "unique_visitors": { - "description": "UniqueVisitors is distinct client IPs (remote_addr) over the window.", "type": "integer" }, "window_ended_at": { @@ -19049,6 +19056,23 @@ } } }, + "zone.StatsPoint": { + "type": "object", + "properties": { + "bucket_started_at": { + "type": "string" + }, + "bytes_sent": { + "type": "integer" + }, + "request_count": { + "type": "integer" + }, + "unique_visitors": { + "type": "integer" + } + } + }, "zone.StatsRange": { "type": "string", "enum": [ diff --git a/docs/swagger.yaml b/docs/swagger.yaml index 8139dbc8..6aed4726 100644 --- a/docs/swagger.yaml +++ b/docs/swagger.yaml @@ -656,6 +656,8 @@ definitions: type: object github_com_Rain-kl_Wavelet_pkg_protocol.NodeAccessLog: properties: + bytes_sent: + type: integer host: type: string logged_at_unix: @@ -3964,33 +3966,41 @@ definitions: zone.Stats: properties: available: - description: Available is false when analytics storage is unavailable (e.g. - ClickHouse off). type: boolean + bucket_minutes: + type: integer bytes_sent: - description: BytesSent is total response bytes. Currently always 0 until access - logs store body_bytes_sent. type: integer domain_count: - description: DomainCount is the number of explicit Zone domains used to scope - hosts. type: integer range: $ref: '#/definitions/zone.StatsRange' range_hours: type: integer request_count: - description: RequestCount is total access-log requests for Zone domains. type: integer + series: + items: + $ref: '#/definitions/zone.StatsPoint' + type: array unique_visitors: - description: UniqueVisitors is distinct client IPs (remote_addr) over the - window. type: integer window_ended_at: type: string window_started_at: type: string type: object + zone.StatsPoint: + properties: + bucket_started_at: + type: string + bytes_sent: + type: integer + request_count: + type: integer + unique_visitors: + type: integer + type: object zone.StatsRange: enum: - 24h diff --git a/internal/apps/agent/observability/traffic.go b/internal/apps/agent/observability/traffic.go index 8dafd8ad..1b85a936 100644 --- a/internal/apps/agent/observability/traffic.go +++ b/internal/apps/agent/observability/traffic.go @@ -199,6 +199,7 @@ func (aggregate *trafficAggregate) consume(line []byte) { Host: strings.TrimSpace(record.Host), Path: normalizeAccessLogPath(record.Path), StatusCode: record.Status, + BytesSent: record.BytesSent, }) } diff --git a/internal/apps/openflare/observability/access_log_logics.go b/internal/apps/openflare/observability/access_log_logics.go index d4593e53..dce6f52f 100644 --- a/internal/apps/openflare/observability/access_log_logics.go +++ b/internal/apps/openflare/observability/access_log_logics.go @@ -197,7 +197,7 @@ func ListAccessLogs(ctx context.Context, input AccessLogQuery) (*AccessLogList, if err != nil { return nil, err } - totalRecords, totalIPs, err := model.CountOpenFlareAccessLogs(ctx, modelQuery) + totalRecords, totalIPs, _, err := model.CountOpenFlareAccessLogs(ctx, modelQuery) if err != nil { return nil, err } @@ -260,7 +260,7 @@ func ListFoldedAccessLogs(ctx context.Context, input AccessLogQuery) (*FoldedAcc if err != nil { return nil, err } - totalRecords, totalIPs, err := model.CountOpenFlareAccessLogs(ctx, modelQuery) + totalRecords, totalIPs, _, err := model.CountOpenFlareAccessLogs(ctx, modelQuery) if err != nil { return nil, err } diff --git a/internal/apps/openflare/zone/logics_test.go b/internal/apps/openflare/zone/logics_test.go index d0ed5a84..b58b8cd0 100644 --- a/internal/apps/openflare/zone/logics_test.go +++ b/internal/apps/openflare/zone/logics_test.go @@ -71,11 +71,11 @@ func TestGetStatsAggregatesZoneHosts(t *testing.T) { 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}, + {NodeID: "n1", LoggedAt: now.Add(-1 * time.Hour), RemoteAddr: "1.1.1.1", Host: "api.example.com", Path: "/", StatusCode: 200, BytesSent: 1000}, + {NodeID: "n1", LoggedAt: now.Add(-2 * time.Hour), RemoteAddr: "1.1.1.1", Host: "www.example.com", Path: "/", StatusCode: 200, BytesSent: 500}, + {NodeID: "n1", LoggedAt: now.Add(-3 * time.Hour), RemoteAddr: "2.2.2.2", Host: "api.example.com", Path: "/x", StatusCode: 404, BytesSent: 200}, + {NodeID: "n1", LoggedAt: now.Add(-3 * time.Hour), RemoteAddr: "3.3.3.3", Host: "other.com", Path: "/", StatusCode: 200, BytesSent: 100}, + {NodeID: "n1", LoggedAt: now.Add(-48 * time.Hour), RemoteAddr: "4.4.4.4", Host: "api.example.com", Path: "/", StatusCode: 200, BytesSent: 800}, })) stats, err := GetStats(ctx, zone.ID, "24h") @@ -83,20 +83,25 @@ func TestGetStatsAggregatesZoneHosts(t *testing.T) { require.Equal(t, StatsRange24h, stats.Range) require.Equal(t, int64(3), stats.RequestCount) require.Equal(t, int64(2), stats.UniqueVisitors) + require.Equal(t, int64(1700), stats.BytesSent) 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 + var seriesBytes int64 for _, point := range stats.Series { seriesRequests += point.RequestCount + seriesBytes += point.BytesSent } require.Equal(t, int64(3), seriesRequests) + require.Equal(t, int64(1700), seriesBytes) 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.Equal(t, int64(2500), stats7d.BytesSent) require.NotEmpty(t, stats7d.Series) _, err = GetStats(ctx, zone.ID, "1h") diff --git a/internal/apps/openflare/zone/stats.go b/internal/apps/openflare/zone/stats.go index 11127721..620bc143 100644 --- a/internal/apps/openflare/zone/stats.go +++ b/internal/apps/openflare/zone/stats.go @@ -22,19 +22,19 @@ const ( // StatsRange24h represents a 24-hour time window. StatsRange24h StatsRange = "24h" // StatsRange7d represents a 7-day time window. - StatsRange7d StatsRange = "7d" + StatsRange7d StatsRange = "7d" // StatsRange30d represents a 30-day time window. StatsRange30d StatsRange = "30d" ) const ( - hoursPerDay = 24 - daysPerWeek = 7 - daysPerMonth = 30 - minutesPerHour = 60 - bucketMinutes24h = 60 - bucketMinutes7d = 6 * minutesPerHour - bucketMinutes30d = 24 * minutesPerHour + hoursPerDay = 24 + daysPerWeek = 7 + daysPerMonth = 30 + minutesPerHour = 60 + bucketMinutes24h = 60 + bucketMinutes7d = 6 * minutesPerHour + bucketMinutes30d = 24 * minutesPerHour ) // StatsPoint is one bucket on a Zone traffic chart. @@ -120,7 +120,7 @@ func GetStats(ctx context.Context, id uint, rangeRaw string) (*Stats, error) { return result, nil } - requestCount, uniqueVisitors, err := model.CountOpenFlareAccessLogs(ctx, model.OpenFlareAccessLogQuery{ + requestCount, uniqueVisitors, totalBytesSent, err := model.CountOpenFlareAccessLogs(ctx, model.OpenFlareAccessLogQuery{ Hosts: hosts, Since: since, Until: now, @@ -134,8 +134,7 @@ func GetStats(ctx context.Context, id uint, rangeRaw string) (*Stats, error) { } result.RequestCount = requestCount result.UniqueVisitors = uniqueVisitors - // Bytes are not yet persisted on edge access logs; keep the field for UI compatibility. - result.BytesSent = 0 + result.BytesSent = totalBytesSent buckets, err := model.ListOpenFlareAccessLogBuckets(ctx, model.OpenFlareAccessLogBucketQuery{ Hosts: hosts, @@ -166,7 +165,7 @@ func GetStats(ctx context.Context, id uint, rangeRaw string) (*Stats, error) { if row, ok := byEpoch[epoch]; ok { series[index].RequestCount = row.RequestCount series[index].UniqueVisitors = row.UniqueIPCount - series[index].BytesSent = 0 + series[index].BytesSent = row.BytesSent } } result.Series = series diff --git a/internal/db/migrator/goose/clickhouse/202607120001_add_bytes_sent_to_node_access_logs.sql b/internal/db/migrator/goose/clickhouse/202607120001_add_bytes_sent_to_node_access_logs.sql new file mode 100644 index 00000000..fc517f7e --- /dev/null +++ b/internal/db/migrator/goose/clickhouse/202607120001_add_bytes_sent_to_node_access_logs.sql @@ -0,0 +1,5 @@ +-- +goose Up +ALTER TABLE of_node_access_logs ADD COLUMN IF NOT EXISTS bytes_sent UInt64; + +-- +goose Down +ALTER TABLE of_node_access_logs DROP COLUMN IF EXISTS bytes_sent; diff --git a/internal/model/analytics/node_access_log.go b/internal/model/analytics/node_access_log.go index 69db2075..03760b1a 100644 --- a/internal/model/analytics/node_access_log.go +++ b/internal/model/analytics/node_access_log.go @@ -10,7 +10,7 @@ import ( const ( nodeAccessLogTableName = "of_node_access_logs" - nodeAccessLogInsertColumns = "id, node_id, logged_at, remote_addr, region, host, path, status_code, created_at" + nodeAccessLogInsertColumns = "id, node_id, logged_at, remote_addr, region, host, path, status_code, bytes_sent, created_at" ) // NodeAccessLog stores OpenFlare edge node access records in ClickHouse. @@ -23,6 +23,7 @@ type NodeAccessLog struct { Host string `gorm:"column:host"` Path string `gorm:"column:path"` StatusCode int32 `gorm:"column:status_code"` + BytesSent uint64 `gorm:"column:bytes_sent"` CreatedAt time.Time `gorm:"column:created_at"` } @@ -39,4 +40,4 @@ func (NodeAccessLog) InsertColumns() string { // BatchInsertSQL returns the INSERT prefix used by native batch writers. func (NodeAccessLog) BatchInsertSQL() string { return fmt.Sprintf("INSERT INTO %s (%s)", nodeAccessLogTableName, nodeAccessLogInsertColumns) -} \ No newline at end of file +} diff --git a/internal/model/analytics/node_access_log_stats.go b/internal/model/analytics/node_access_log_stats.go new file mode 100644 index 00000000..0cbbb26e --- /dev/null +++ b/internal/model/analytics/node_access_log_stats.go @@ -0,0 +1,58 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package analytics + +// NodeAccessLogBucketAggregate is a folded bucket aggregate row. +type NodeAccessLogBucketAggregate struct { + BucketEpoch int64 `gorm:"column:bucket_epoch"` + RequestCount int64 `gorm:"column:request_count"` + SuccessCount int64 `gorm:"column:success_count"` + ClientErrorCount int64 `gorm:"column:client_error_count"` + ServerErrorCount int64 `gorm:"column:server_error_count"` + UniqueIPCount int64 `gorm:"column:unique_ip_count"` + UniqueHostCount int64 `gorm:"column:unique_host_count"` + BytesSent int64 `gorm:"column:bytes_sent"` +} + +// NodeAccessLogWAFIPAggregate is a per-IP aggregate row for WAF automatic rules. +type NodeAccessLogWAFIPAggregate struct { + RemoteAddr string + RequestCount int64 + Status404Count int64 + ClientErrorCount int64 + ServerErrorCount int64 + IPHostCount int64 + LastSeenEpoch int64 + StatusCounts map[int]int64 +} + +// NodeAccessLogBucketDimension is a bucket dimension value. +type NodeAccessLogBucketDimension struct { + BucketEpoch int64 `gorm:"column:bucket_epoch"` + Value string `gorm:"column:value"` +} + +// NodeAccessLogIPAggregate is an IP aggregate row. +type NodeAccessLogIPAggregate struct { + RemoteAddr string `gorm:"column:remote_addr"` + RequestCount int64 `gorm:"column:request_count"` + SuccessCount int64 `gorm:"column:success_count"` + ClientErrorCount int64 `gorm:"column:client_error_count"` + ServerErrorCount int64 `gorm:"column:server_error_count"` + LastSeenEpoch int64 `gorm:"column:last_seen_epoch"` +} + +// NodeAccessLogIPSummary is an IP summary row. +type NodeAccessLogIPSummary struct { + RemoteAddr string `gorm:"column:remote_addr"` + TotalRequests int64 `gorm:"column:total_requests"` + RecentRequests int64 `gorm:"column:recent_requests"` + LastSeenEpoch int64 `gorm:"column:last_seen_epoch"` +} + +// NodeAccessLogIPTrend is an IP trend bucket row. +type NodeAccessLogIPTrend struct { + BucketEpoch int64 `gorm:"column:bucket_epoch"` + RequestCount int64 `gorm:"column:request_count"` +} diff --git a/internal/model/openflare_access_log.go b/internal/model/openflare_access_log.go index f5c08ff6..a1a050a4 100644 --- a/internal/model/openflare_access_log.go +++ b/internal/model/openflare_access_log.go @@ -9,6 +9,8 @@ import ( "sort" "strings" "time" + + analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" ) const ( @@ -18,52 +20,12 @@ const ( secondsPerMinute = 60 ) -type openFlareAccessLogBucketAggregateRow struct { - BucketEpoch int64 `gorm:"column:bucket_epoch"` - RequestCount int64 `gorm:"column:request_count"` - SuccessCount int64 `gorm:"column:success_count"` - ClientErrorCount int64 `gorm:"column:client_error_count"` - ServerErrorCount int64 `gorm:"column:server_error_count"` - UniqueIPCount int64 `gorm:"column:unique_ip_count"` - UniqueHostCount int64 `gorm:"column:unique_host_count"` -} - -type openFlareAccessLogBucketDimensionRow struct { - BucketEpoch int64 `gorm:"column:bucket_epoch"` - Value string `gorm:"column:value"` -} - -type openFlareAccessLogIPAggregateRow struct { - RemoteAddr string `gorm:"column:remote_addr"` - RequestCount int64 `gorm:"column:request_count"` - SuccessCount int64 `gorm:"column:success_count"` - ClientErrorCount int64 `gorm:"column:client_error_count"` - ServerErrorCount int64 `gorm:"column:server_error_count"` - LastSeenEpoch int64 `gorm:"column:last_seen_epoch"` -} - -type openFlareAccessLogIPSummaryRow struct { - RemoteAddr string `gorm:"column:remote_addr"` - TotalRequests int64 `gorm:"column:total_requests"` - RecentRequests int64 `gorm:"column:recent_requests"` - LastSeenEpoch int64 `gorm:"column:last_seen_epoch"` -} - -type openFlareAccessLogIPTrendRow struct { - BucketEpoch int64 `gorm:"column:bucket_epoch"` - RequestCount int64 `gorm:"column:request_count"` -} - -type openFlareAccessLogWAFIPAggregateRow struct { - RemoteAddr string - RequestCount int64 - Status404Count int64 - ClientErrorCount int64 - ServerErrorCount int64 - IPHostCount int64 - LastSeenEpoch int64 - StatusCounts map[int]int64 -} +type openFlareAccessLogBucketAggregateRow = analyticsmodel.NodeAccessLogBucketAggregate +type openFlareAccessLogBucketDimensionRow = analyticsmodel.NodeAccessLogBucketDimension +type openFlareAccessLogIPAggregateRow = analyticsmodel.NodeAccessLogIPAggregate +type openFlareAccessLogIPSummaryRow = analyticsmodel.NodeAccessLogIPSummary +type openFlareAccessLogIPTrendRow = analyticsmodel.NodeAccessLogIPTrend +type openFlareAccessLogWAFIPAggregateRow = analyticsmodel.NodeAccessLogWAFIPAggregate // ListOpenFlareAccessLogWAFIPAggregates returns per-IP aggregates for WAF automatic rules. func ListOpenFlareAccessLogWAFIPAggregates(ctx context.Context, query OpenFlareAccessLogQuery) ([]*OpenFlareAccessLogWAFIPAggregate, error) { @@ -105,8 +67,8 @@ func ListOpenFlareAccessLogs(ctx context.Context, query OpenFlareAccessLogQuery) return currentAccessLogStore().List(ctx, query) } -// CountOpenFlareAccessLogs counts access logs and distinct IPs matching the query. -func CountOpenFlareAccessLogs(ctx context.Context, query OpenFlareAccessLogQuery) (int64, int64, error) { +// CountOpenFlareAccessLogs counts access logs, distinct IPs, and total bytes sent matching the query. +func CountOpenFlareAccessLogs(ctx context.Context, query OpenFlareAccessLogQuery) (int64, int64, int64, error) { return currentAccessLogStore().Count(ctx, query) } @@ -229,6 +191,7 @@ func buildOpenFlareAccessLogBucketRows(ctx context.Context, query OpenFlareAcces SuccessCount: partial.SuccessCount, ClientErrorCount: partial.ClientErrorCount, ServerErrorCount: partial.ServerErrorCount, + BytesSent: partial.BytesSent, }) } return rows, nil diff --git a/internal/model/openflare_access_log_store.go b/internal/model/openflare_access_log_store.go index 9cd1eca6..6bb44560 100644 --- a/internal/model/openflare_access_log_store.go +++ b/internal/model/openflare_access_log_store.go @@ -5,6 +5,7 @@ package model import ( "context" + "math" "sync" "time" @@ -39,7 +40,7 @@ func currentAccessLogInsertHooks() AccessLogInsertHooks { type accessLogStore interface { InsertBatch(ctx context.Context, records []*OpenFlareAccessLog) error List(ctx context.Context, query OpenFlareAccessLogQuery) ([]*OpenFlareAccessLog, error) - Count(ctx context.Context, query OpenFlareAccessLogQuery) (int64, int64, error) + Count(ctx context.Context, query OpenFlareAccessLogQuery) (int64, int64, int64, error) RegionCounts(ctx context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareAccessLogRegionCount, error) BucketAggregates(ctx context.Context, filter OpenFlareAccessLogQuery, bucketSeconds int64) ([]openFlareAccessLogBucketAggregateRow, error) CountBuckets(ctx context.Context, filter OpenFlareAccessLogQuery, bucketSeconds int64) (int64, error) @@ -112,7 +113,7 @@ func (clickhouseAccessLogStore) List(ctx context.Context, query OpenFlareAccessL return fromAnalyticsNodeAccessLogs(rows), nil } -func (clickhouseAccessLogStore) Count(ctx context.Context, query OpenFlareAccessLogQuery) (int64, int64, error) { +func (clickhouseAccessLogStore) Count(ctx context.Context, query OpenFlareAccessLogQuery) (int64, int64, int64, error) { return analyticsrepo.CountNodeAccessLogs(ctx, toNodeAccessLogFilter(query)) } @@ -132,23 +133,7 @@ func (clickhouseAccessLogStore) RegionCounts(ctx context.Context, nodeID string, } func (clickhouseAccessLogStore) BucketAggregates(ctx context.Context, filter OpenFlareAccessLogQuery, bucketSeconds int64) ([]openFlareAccessLogBucketAggregateRow, error) { - rows, err := analyticsrepo.BucketAggregatesNodeAccessLogs(ctx, toNodeAccessLogFilter(filter), bucketSeconds) - if err != nil { - return nil, err - } - result := make([]openFlareAccessLogBucketAggregateRow, len(rows)) - for index, row := range rows { - result[index] = openFlareAccessLogBucketAggregateRow{ - BucketEpoch: row.BucketEpoch, - RequestCount: row.RequestCount, - SuccessCount: row.SuccessCount, - ClientErrorCount: row.ClientErrorCount, - ServerErrorCount: row.ServerErrorCount, - UniqueIPCount: row.UniqueIPCount, - UniqueHostCount: row.UniqueHostCount, - } - } - return result, nil + return analyticsrepo.BucketAggregatesNodeAccessLogs(ctx, toNodeAccessLogFilter(filter), bucketSeconds) } func (clickhouseAccessLogStore) CountBuckets(ctx context.Context, filter OpenFlareAccessLogQuery, bucketSeconds int64) (int64, error) { @@ -156,54 +141,15 @@ func (clickhouseAccessLogStore) CountBuckets(ctx context.Context, filter OpenFla } func (clickhouseAccessLogStore) BucketDimensions(ctx context.Context, filter OpenFlareAccessLogQuery, column string, bucketSeconds int64) ([]openFlareAccessLogBucketDimensionRow, error) { - rows, err := analyticsrepo.BucketDimensionsNodeAccessLogs(ctx, toNodeAccessLogFilter(filter), column, bucketSeconds) - if err != nil { - return nil, err - } - result := make([]openFlareAccessLogBucketDimensionRow, len(rows)) - for index, row := range rows { - result[index] = openFlareAccessLogBucketDimensionRow{ - BucketEpoch: row.BucketEpoch, - Value: row.Value, - } - } - return result, nil + return analyticsrepo.BucketDimensionsNodeAccessLogs(ctx, toNodeAccessLogFilter(filter), column, bucketSeconds) } func (clickhouseAccessLogStore) IPAggregates(ctx context.Context, filter OpenFlareAccessLogQuery, exactRemoteAddr bool) ([]openFlareAccessLogIPAggregateRow, error) { - rows, err := analyticsrepo.IPAggregatesNodeAccessLogs(ctx, toNodeAccessLogFilter(filter), exactRemoteAddr) - if err != nil { - return nil, err - } - result := make([]openFlareAccessLogIPAggregateRow, len(rows)) - for index, row := range rows { - result[index] = openFlareAccessLogIPAggregateRow{ - RemoteAddr: row.RemoteAddr, - RequestCount: row.RequestCount, - SuccessCount: row.SuccessCount, - ClientErrorCount: row.ClientErrorCount, - ServerErrorCount: row.ServerErrorCount, - LastSeenEpoch: row.LastSeenEpoch, - } - } - return result, nil + return analyticsrepo.IPAggregatesNodeAccessLogs(ctx, toNodeAccessLogFilter(filter), exactRemoteAddr) } func (clickhouseAccessLogStore) IPSummaries(ctx context.Context, filter OpenFlareAccessLogQuery, recentSince time.Time) ([]openFlareAccessLogIPSummaryRow, error) { - rows, err := analyticsrepo.IPSummariesNodeAccessLogs(ctx, toNodeAccessLogFilter(filter), recentSince) - if err != nil { - return nil, err - } - result := make([]openFlareAccessLogIPSummaryRow, len(rows)) - for index, row := range rows { - result[index] = openFlareAccessLogIPSummaryRow{ - RemoteAddr: row.RemoteAddr, - TotalRequests: row.TotalRequests, - RecentRequests: row.RecentRequests, - LastSeenEpoch: row.LastSeenEpoch, - } - } - return result, nil + return analyticsrepo.IPSummariesNodeAccessLogs(ctx, toNodeAccessLogFilter(filter), recentSince) } func (clickhouseAccessLogStore) CountIPSummaries(ctx context.Context, filter OpenFlareAccessLogQuery) (int64, error) { @@ -211,39 +157,11 @@ func (clickhouseAccessLogStore) CountIPSummaries(ctx context.Context, filter Ope } func (clickhouseAccessLogStore) WAFIPAggregates(ctx context.Context, filter OpenFlareAccessLogQuery) ([]openFlareAccessLogWAFIPAggregateRow, error) { - rows, err := analyticsrepo.IPAggregatesForWAFNodeAccessLogs(ctx, toNodeAccessLogFilter(filter)) - if err != nil { - return nil, err - } - result := make([]openFlareAccessLogWAFIPAggregateRow, len(rows)) - for index, row := range rows { - result[index] = openFlareAccessLogWAFIPAggregateRow{ - RemoteAddr: row.RemoteAddr, - RequestCount: row.RequestCount, - Status404Count: row.Status404Count, - ClientErrorCount: row.ClientErrorCount, - ServerErrorCount: row.ServerErrorCount, - IPHostCount: row.IPHostCount, - LastSeenEpoch: row.LastSeenEpoch, - StatusCounts: row.StatusCounts, - } - } - return result, nil + return analyticsrepo.IPAggregatesForWAFNodeAccessLogs(ctx, toNodeAccessLogFilter(filter)) } func (clickhouseAccessLogStore) IPTrend(ctx context.Context, filter OpenFlareAccessLogQuery, bucketSeconds int64) ([]openFlareAccessLogIPTrendRow, error) { - rows, err := analyticsrepo.IPTrendNodeAccessLogs(ctx, toNodeAccessLogFilter(filter), bucketSeconds) - if err != nil { - return nil, err - } - result := make([]openFlareAccessLogIPTrendRow, len(rows)) - for index, row := range rows { - result[index] = openFlareAccessLogIPTrendRow{ - BucketEpoch: row.BucketEpoch, - RequestCount: row.RequestCount, - } - } - return result, nil + return analyticsrepo.IPTrendNodeAccessLogs(ctx, toNodeAccessLogFilter(filter), bucketSeconds) } func (clickhouseAccessLogStore) DeleteAll(ctx context.Context) (int64, error) { @@ -275,6 +193,10 @@ func toNodeAccessLogFilter(query OpenFlareAccessLogQuery) analyticsrepo.NodeAcce } func toAnalyticsNodeAccessLog(record *OpenFlareAccessLog) analyticsmodel.NodeAccessLog { + var bytesSent uint64 + if record.BytesSent > 0 { + bytesSent = uint64(record.BytesSent) + } return analyticsmodel.NodeAccessLog{ ID: record.ID, NodeID: record.NodeID, @@ -284,6 +206,7 @@ func toAnalyticsNodeAccessLog(record *OpenFlareAccessLog) analyticsmodel.NodeAcc Host: record.Host, Path: record.Path, StatusCode: openFlareAccessLogStatusCodeToInt32(record.StatusCode), + BytesSent: bytesSent, CreatedAt: record.CreatedAt, } } @@ -291,6 +214,12 @@ func toAnalyticsNodeAccessLog(record *OpenFlareAccessLog) analyticsmodel.NodeAcc func fromAnalyticsNodeAccessLogs(rows []analyticsmodel.NodeAccessLog) []*OpenFlareAccessLog { result := make([]*OpenFlareAccessLog, len(rows)) for index, row := range rows { + var bytesSent int64 + if row.BytesSent <= math.MaxInt64 { + bytesSent = int64(row.BytesSent) + } else { + bytesSent = math.MaxInt64 + } result[index] = &OpenFlareAccessLog{ ID: row.ID, NodeID: row.NodeID, @@ -300,6 +229,7 @@ func fromAnalyticsNodeAccessLogs(rows []analyticsmodel.NodeAccessLog) []*OpenFla Host: row.Host, Path: row.Path, StatusCode: int(row.StatusCode), + BytesSent: bytesSent, CreatedAt: row.CreatedAt, } } diff --git a/internal/model/openflare_access_log_store_memory.go b/internal/model/openflare_access_log_store_memory.go index cb282f86..2d30b5af 100644 --- a/internal/model/openflare_access_log_store_memory.go +++ b/internal/model/openflare_access_log_store_memory.go @@ -55,19 +55,21 @@ func (s *memoryAccessLogStore) List(_ context.Context, query OpenFlareAccessLogQ return cloneAccessLogSlice(rows), nil } -func (s *memoryAccessLogStore) Count(_ context.Context, query OpenFlareAccessLogQuery) (int64, int64, error) { +func (s *memoryAccessLogStore) Count(_ context.Context, query OpenFlareAccessLogQuery) (int64, int64, int64, error) { s.mu.RLock() defer s.mu.RUnlock() rows := s.filterRecords(query) ips := make(map[string]struct{}) + var totalBytes int64 for _, row := range rows { + totalBytes += row.BytesSent remoteAddr := strings.TrimSpace(row.RemoteAddr) if remoteAddr == "" { continue } ips[remoteAddr] = struct{}{} } - return int64(len(rows)), int64(len(ips)), nil + return int64(len(rows)), int64(len(ips)), totalBytes, nil } func (s *memoryAccessLogStore) RegionCounts(_ context.Context, nodeID string, since time.Time, limit int) ([]*OpenFlareAccessLogRegionCount, error) { @@ -120,6 +122,7 @@ func (s *memoryAccessLogStore) BucketAggregates(_ context.Context, filter OpenFl aggregates[bucketEpoch] = item } item.RequestCount++ + item.BytesSent += row.BytesSent switch { case row.StatusCode < 400: item.SuccessCount++ @@ -151,6 +154,7 @@ func (s *memoryAccessLogStore) BucketAggregates(_ context.Context, filter OpenFl SuccessCount: result[index].SuccessCount, ClientErrorCount: result[index].ClientErrorCount, ServerErrorCount: result[index].ServerErrorCount, + BytesSent: result[index].BytesSent, } } sortOpenFlareAccessLogBucketRows(bucketRows, filter.SortBy, filter.SortOrder) @@ -163,6 +167,7 @@ func (s *memoryAccessLogStore) BucketAggregates(_ context.Context, filter OpenFl SuccessCount: bucketRows[index].SuccessCount, ClientErrorCount: bucketRows[index].ClientErrorCount, ServerErrorCount: bucketRows[index].ServerErrorCount, + BytesSent: bucketRows[index].BytesSent, } } if filter.PageSize > 0 { diff --git a/internal/model/openflare_access_log_test.go b/internal/model/openflare_access_log_test.go index 920f1227..f1f6a0d1 100644 --- a/internal/model/openflare_access_log_test.go +++ b/internal/model/openflare_access_log_test.go @@ -76,7 +76,7 @@ func TestCountOpenFlareAccessLogs(t *testing.T) { query := OpenFlareAccessLogQuery{ Since: now.Add(-10 * time.Minute), } - totalRecords, totalIPs, err := CountOpenFlareAccessLogs(ctx, query) + totalRecords, totalIPs, _, err := CountOpenFlareAccessLogs(ctx, query) require.NoError(t, err) assert.Equal(t, int64(5), totalRecords) assert.Equal(t, int64(3), totalIPs) @@ -113,7 +113,7 @@ func TestDeleteOpenFlareAccessLogsBefore(t *testing.T) { require.NoError(t, err) assert.Equal(t, int64(3), deleted) - totalRecords, _, err := CountOpenFlareAccessLogs(ctx, OpenFlareAccessLogQuery{Since: now.Add(-10 * time.Minute)}) + totalRecords, _, _, err := CountOpenFlareAccessLogs(ctx, OpenFlareAccessLogQuery{Since: now.Add(-10 * time.Minute)}) require.NoError(t, err) assert.Equal(t, int64(2), totalRecords) -} \ No newline at end of file +} diff --git a/internal/model/openflare_observability.go b/internal/model/openflare_observability.go index b3680659..9e664744 100644 --- a/internal/model/openflare_observability.go +++ b/internal/model/openflare_observability.go @@ -69,6 +69,7 @@ type OpenFlareAccessLog struct { Host string `json:"host" gorm:"index;size:255"` Path string `json:"path" gorm:"size:2048"` StatusCode int `json:"status_code" gorm:"index"` + BytesSent int64 `json:"bytes_sent" gorm:"column:bytes_sent;not null;default:0"` CreatedAt time.Time `json:"created_at" gorm:"autoCreateTime"` } @@ -221,6 +222,7 @@ type OpenFlareAccessLogBucketRow struct { SuccessCount int64 `json:"success_count"` ClientErrorCount int64 `json:"client_error_count"` ServerErrorCount int64 `json:"server_error_count"` + BytesSent int64 `json:"bytes_sent"` } // OpenFlareAccessLogBucketIPQuery filters folded IP summary queries (v1 stub). diff --git a/internal/repository/analytics/node_access_log.go b/internal/repository/analytics/node_access_log.go index 333160f4..93fa5f2c 100644 --- a/internal/repository/analytics/node_access_log.go +++ b/internal/repository/analytics/node_access_log.go @@ -35,7 +35,7 @@ func ListNodeAccessLogs(ctx context.Context, filter NodeAccessLogFilter) ([]anal clause, args := buildNodeAccessLogFilterClause(filter) tableName := nodeAccessLogTableName() sql := fmt.Sprintf(` -SELECT id, node_id, logged_at, remote_addr, region, host, path, status_code, created_at +SELECT id, node_id, logged_at, remote_addr, region, host, path, status_code, bytes_sent, created_at FROM %s WHERE %s ORDER BY %s`, tableName, clause, nodeAccessLogOrderClause(filter.SortBy, filter.SortOrder)) @@ -67,6 +67,7 @@ func scanNodeAccessLogRows(rows driver.Rows) ([]analyticsmodel.NodeAccessLog, er &item.Host, &item.Path, &item.StatusCode, + &item.BytesSent, &item.CreatedAt, ); err != nil { return nil, fmt.Errorf("scan node access log row: %w", err) @@ -78,11 +79,11 @@ func scanNodeAccessLogRows(rows driver.Rows) ([]analyticsmodel.NodeAccessLog, er return result, nil } -// CountNodeAccessLogs returns total records and distinct IPs matching filter. -func CountNodeAccessLogs(ctx context.Context, filter NodeAccessLogFilter) (int64, int64, error) { +// CountNodeAccessLogs returns total records, distinct IPs, and total bytes sent matching filter. +func CountNodeAccessLogs(ctx context.Context, filter NodeAccessLogFilter) (int64, int64, int64, error) { conn, err := nodeAccessLogConn() if err != nil { - return 0, 0, err + return 0, 0, 0, err } clause, args := buildNodeAccessLogFilterClause(filter) tableName := nodeAccessLogTableName() @@ -90,14 +91,15 @@ func CountNodeAccessLogs(ctx context.Context, filter NodeAccessLogFilter) (int64 countSQL := fmt.Sprintf(` SELECT count() AS total_records, - uniqExactIf(remote_addr, remote_addr != '') AS total_ips + uniqExactIf(remote_addr, remote_addr != '') AS total_ips, + sum(bytes_sent) AS total_bytes FROM %s WHERE %s`, tableName, clause) - var totalRecords, totalIPs uint64 - if err := conn.QueryRow(ctx, countSQL, args...).Scan(&totalRecords, &totalIPs); err != nil { - return 0, 0, fmt.Errorf("count node access logs: %w", err) + var totalRecords, totalIPs, totalBytes uint64 + if err := conn.QueryRow(ctx, countSQL, args...).Scan(&totalRecords, &totalIPs, &totalBytes); err != nil { + return 0, 0, 0, fmt.Errorf("count node access logs: %w", err) } - return safeInt64Count(totalRecords), safeInt64Count(totalIPs), nil + return safeInt64Count(totalRecords), safeInt64Count(totalIPs), safeInt64Count(totalBytes), nil } // RegionCountsNodeAccessLogs returns region counts for a node since a time. diff --git a/internal/repository/analytics/node_access_log_stats.go b/internal/repository/analytics/node_access_log_stats.go index adb8e489..93e3d20f 100644 --- a/internal/repository/analytics/node_access_log_stats.go +++ b/internal/repository/analytics/node_access_log_stats.go @@ -8,60 +8,27 @@ import ( "fmt" "strings" "time" + + analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" ) // NodeAccessLogBucketAggregate is a folded bucket aggregate row. -type NodeAccessLogBucketAggregate struct { - BucketEpoch int64 - RequestCount int64 - SuccessCount int64 - ClientErrorCount int64 - ServerErrorCount int64 - UniqueIPCount int64 - UniqueHostCount int64 -} +type NodeAccessLogBucketAggregate = analyticsmodel.NodeAccessLogBucketAggregate // NodeAccessLogWAFIPAggregate is a per-IP aggregate row for WAF automatic rules. -type NodeAccessLogWAFIPAggregate struct { - RemoteAddr string - RequestCount int64 - Status404Count int64 - ClientErrorCount int64 - ServerErrorCount int64 - IPHostCount int64 - LastSeenEpoch int64 - StatusCounts map[int]int64 -} +type NodeAccessLogWAFIPAggregate = analyticsmodel.NodeAccessLogWAFIPAggregate // NodeAccessLogBucketDimension is a bucket dimension value. -type NodeAccessLogBucketDimension struct { - BucketEpoch int64 - Value string -} +type NodeAccessLogBucketDimension = analyticsmodel.NodeAccessLogBucketDimension // NodeAccessLogIPAggregate is an IP aggregate row. -type NodeAccessLogIPAggregate struct { - RemoteAddr string - RequestCount int64 - SuccessCount int64 - ClientErrorCount int64 - ServerErrorCount int64 - LastSeenEpoch int64 -} +type NodeAccessLogIPAggregate = analyticsmodel.NodeAccessLogIPAggregate // NodeAccessLogIPSummary is an IP summary row. -type NodeAccessLogIPSummary struct { - RemoteAddr string - TotalRequests int64 - RecentRequests int64 - LastSeenEpoch int64 -} +type NodeAccessLogIPSummary = analyticsmodel.NodeAccessLogIPSummary // NodeAccessLogIPTrend is an IP trend bucket row. -type NodeAccessLogIPTrend struct { - BucketEpoch int64 - RequestCount int64 -} +type NodeAccessLogIPTrend = analyticsmodel.NodeAccessLogIPTrend // BucketAggregatesNodeAccessLogs returns folded bucket aggregates with unique IP/host counts. func BucketAggregatesNodeAccessLogs(ctx context.Context, filter NodeAccessLogFilter, bucketSeconds int64) ([]NodeAccessLogBucketAggregate, error) { @@ -80,7 +47,8 @@ SELECT countIf(status_code >= 400 AND status_code < 500) AS client_error_count, countIf(status_code >= 500) AS server_error_count, uniqExactIf(remote_addr, remote_addr != '') AS unique_ip_count, - uniqExactIf(host, host != '') AS unique_host_count + uniqExactIf(host, host != '') AS unique_host_count, + sum(bytes_sent) AS bytes_sent FROM %s WHERE %s GROUP BY bucket_epoch @@ -101,10 +69,10 @@ ORDER BY %s`, bucketExpr, tableName, clause, nodeAccessLogBucketOrderClause(filt var result []NodeAccessLogBucketAggregate for rows.Next() { var ( - bucketEpoch int64 - requestCount, successCount, clientErrorCount, serverErrorCount, uniqueIPCount, uniqueHostCount uint64 + bucketEpoch int64 + requestCount, successCount, clientErrorCount, serverErrorCount, uniqueIPCount, uniqueHostCount, bytesSent uint64 ) - if err := rows.Scan(&bucketEpoch, &requestCount, &successCount, &clientErrorCount, &serverErrorCount, &uniqueIPCount, &uniqueHostCount); err != nil { + if err := rows.Scan(&bucketEpoch, &requestCount, &successCount, &clientErrorCount, &serverErrorCount, &uniqueIPCount, &uniqueHostCount, &bytesSent); err != nil { return nil, fmt.Errorf("scan bucket aggregate row: %w", err) } result = append(result, NodeAccessLogBucketAggregate{ @@ -115,6 +83,7 @@ ORDER BY %s`, bucketExpr, tableName, clause, nodeAccessLogBucketOrderClause(filt ServerErrorCount: safeInt64Count(serverErrorCount), UniqueIPCount: safeInt64Count(uniqueIPCount), UniqueHostCount: safeInt64Count(uniqueHostCount), + BytesSent: safeInt64Count(bytesSent), }) } return result, nil @@ -215,8 +184,8 @@ GROUP BY remote_addr`, lastSeenExpr, tableName, queryClause) var result []NodeAccessLogIPAggregate for rows.Next() { var ( - remoteAddr string - lastSeenEpoch int64 + remoteAddr string + lastSeenEpoch int64 requestCount, successCount, clientErrorCount, serverErrorCount uint64 ) if err := rows.Scan(&remoteAddr, &requestCount, &successCount, &clientErrorCount, &serverErrorCount, &lastSeenEpoch); err != nil { @@ -347,8 +316,8 @@ GROUP BY remote_addr`, hostIsIPExpr, lastSeenExpr, tableName, clause) order := make([]string, 0) for rows.Next() { var ( - remoteAddr string - lastSeenEpoch int64 + remoteAddr string + lastSeenEpoch int64 requestCount, status404Count, clientErrorCount, serverErrorCount, ipHostCount uint64 ) if err := rows.Scan(&remoteAddr, &requestCount, &status404Count, &clientErrorCount, &serverErrorCount, &ipHostCount, &lastSeenEpoch); err != nil { diff --git a/pkg/protocol/agent.go b/pkg/protocol/agent.go index 33f529e4..b79311e4 100644 --- a/pkg/protocol/agent.go +++ b/pkg/protocol/agent.go @@ -151,6 +151,7 @@ type NodeAccessLog struct { Host string `json:"host"` Path string `json:"path"` StatusCode int `json:"status_code"` + BytesSent int64 `json:"bytes_sent"` } // BufferedObservabilityRecord is a buffered observability record.