mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-28 05:46:36 +08:00
feat(clickhouse): integrate goose migrations and analytics repository
Merge upstream Wavelet patch to manage ClickHouse DDL via goose (goose_clickhouse_version), add GORM ChDB access, and centralize w_user_access_logs reads/writes in internal/repository/analytics. OpenFlare of_node_access_logs DDL moves to goose/clickhouse; remove manual clickhouse_schema startup and support-files SQL duplicates.
This commit is contained in:
@@ -1,6 +1,6 @@
|
||||
---
|
||||
name: "database-migration"
|
||||
description: "Wavelet 项目专用:当新增或修改数据库表结构、索引、初始化数据、系统配置 seed、模板 seed、默认管理员、goose SQL 迁移、internal/db/migrator 或数据库升级流程时必须使用。本技能指导在 internal/db/migrator/goose 下编写 PostgreSQL/SQLite 双方言 SQL 迁移,并完成验证。"
|
||||
description: "Wavelet 项目专用:当新增或修改数据库表结构、索引、初始化数据、系统配置 seed、模板 seed、默认管理员、goose SQL 迁移、internal/db/migrator、ClickHouse 分析库 DDL 或数据库升级流程时必须使用。本技能指导在 internal/db/migrator/goose 下编写 PostgreSQL/SQLite 双方言 SQL 迁移,以及在 goose/clickhouse 下编写 ClickHouse 单方言分析表迁移,并完成验证。"
|
||||
---
|
||||
|
||||
# Wavelet 数据库升级操作指南
|
||||
@@ -75,3 +75,64 @@ make code-check
|
||||
- `system_configs`、默认 `admin`、内置模板能按预期初始化。
|
||||
- 新增表/列与 Go model 的列名、类型和默认值兼容。
|
||||
- 前端或接口消费的公共配置值仍按字符串解析。
|
||||
|
||||
## ClickHouse 分析库(辅助 OLAP)
|
||||
|
||||
ClickHouse 是**辅助 OLAP 存储**,与 PostgreSQL/SQLite 主库**完全独立**的迁移与访问管线:
|
||||
|
||||
- 主库(PG/SQLite):业务事务数据、`goose_db_version`、双方言 SQL。
|
||||
- 分析库(ClickHouse):访问日志、统计聚合等分析型数据、`goose_clickhouse_version`、单方言 SQL。
|
||||
|
||||
**不要**把 ClickHouse 表结构混入 PG/SQLite 迁移目录,也**不要**在 `support-files/`、`internal/apps/` 或 `internal/repository/` 中手写 DDL。
|
||||
|
||||
### 目录与职责
|
||||
|
||||
| 路径 | 职责 |
|
||||
| :--- | :--- |
|
||||
| `internal/db/migrator/goose/clickhouse/` | **唯一** ClickHouse DDL 来源(goose SQL,嵌入二进制) |
|
||||
| `internal/model/analytics/` | 分析表 Go model,列名须与 goose DDL 一致 |
|
||||
| `internal/repository/analytics/` | 所有 ClickHouse 读写(批量写入、查询、聚合) |
|
||||
| `internal/db/clickhouse.go` | 连接初始化(`ChConn` 原生批量、`ChDB` GORM 查询) |
|
||||
|
||||
### 迁移入口与版本表
|
||||
|
||||
- 入口:`migrator.MigrateClickHouse()`,在 `cmd/root.go` 的 `PreRun` 中于 `migrator.Migrate()` 之后调用。
|
||||
- 仅当 `clickhouse.enabled: true` 时执行;禁用时直接跳过(见 `TestMigrateClickHouseSkipsWhenDisabled`)。
|
||||
- 版本表:`goose_clickhouse_version`,与主库 `goose_db_version` **分离**,互不影响。
|
||||
- 方言:仅 ClickHouse,**无** SQLite 镜像目录。
|
||||
|
||||
### ClickHouse 迁移规则
|
||||
|
||||
1. **DDL 只写 goose SQL**:`CREATE TABLE IF NOT EXISTS ...`,禁止 GORM `AutoMigrate`、禁止在 repository 或 handler 中建表。
|
||||
2. **无事务**:ClickHouse 不支持 goose 事务包装;每个 `Up`/`Down` 语句独立提交。
|
||||
3. **幂等 Up**:表用 `IF NOT EXISTS`;`Down` 用 `DROP TABLE IF EXISTS`。
|
||||
4. **Down 谨慎**:MergeTree 等引擎上 `DROP TABLE` 会立即删除数据,生产环境通常只前滚;仅在开发/测试需要回滚时编写 `Down`。
|
||||
5. **DDL 与 DML 分离**:与主库相同,表结构变更与数据初始化分文件、分版本号;分析表通常无 seed,批量写入由 repository 在运行时完成。
|
||||
6. **引擎与排序键**:在 SQL 中显式声明 `ENGINE`、`PARTITION BY`、`ORDER BY` 等,与查询模式对齐(例如按 `created_at` 分区)。
|
||||
7. **禁止重复 DDL**:不要在 `support-files/`、`apps` 初始化逻辑或 `repository/analytics` 中复制建表语句。
|
||||
|
||||
### 新增分析表工作流
|
||||
|
||||
按以下顺序落地,避免列名或类型漂移:
|
||||
|
||||
1. **Model**:在 `internal/model/analytics/` 定义 struct,`gorm:"column:..."` 与 DDL 列名一一对应;实现 `TableName()`,批量写入表可提供 `InsertColumns()` / `BatchInsertSQL()`。
|
||||
2. **Goose SQL**:在 `internal/db/migrator/goose/clickhouse/` 新增递增版本文件(格式同主库,如 `YYYYMMDDNNNN_create_xxx.sql`),编写 `-- +goose Up` / `-- +goose Down`。
|
||||
3. **Repository**:在 `internal/repository/analytics/` 实现写入(优先 `db.ChConn` 批量)与查询(`db.ChDB`);连接未初始化时返回明确错误,**不要**在 handler 写 SQL。
|
||||
4. **Apps**:在 `internal/apps/` 编排业务(如中间件采集、管理端统计 API),只调用 repository,不触达 DDL。
|
||||
|
||||
### ClickHouse 验证
|
||||
|
||||
至少运行:
|
||||
|
||||
```bash
|
||||
go test ./internal/db/migrator
|
||||
go test ./internal/repository/analytics
|
||||
make code-check
|
||||
```
|
||||
|
||||
验证重点:
|
||||
|
||||
- goose 能在空 ClickHouse 实例上完整执行 `Up`。
|
||||
- `internal/model/analytics` 列名、类型与 goose SQL 一致。
|
||||
- repository 读写路径不依赖 handler 内联 SQL。
|
||||
- `clickhouse.enabled: false` 时启动不报错、不执行迁移。
|
||||
|
||||
@@ -42,7 +42,7 @@
|
||||
| `new-api` | 添加或修改自定义业务 API、Handler、服务层逻辑、自定义路由注册 |
|
||||
| `new-async-task` | 添加或修改 Asynq 任务、定时任务、TaskHandler、任务元数据 |
|
||||
| `new-setting` | 添加或修改系统/业务/公开设置、`/admin/system` 参数或 `/admin/settings` 图形化设置 |
|
||||
| `database-migration` | 数据库表结构变更、goose SQL 迁移、seed 数据 |
|
||||
| `database-migration` | 数据库表结构变更、goose SQL 迁移(PG/SQLite/ClickHouse)、seed 数据 |
|
||||
| `file-upload` | 业务上传文件、Worker 程序化摄取、`upload.Ingest` 策略选型、文件访问与 `w_uploads` / 统计排查 |
|
||||
| `push-notification` | 系统通知推送事件、统一触发器投递、带消息推送的业务功能 |
|
||||
| `release-guide` | 根据自上一正式版本 Tag 以来的提交整理 Version Bump 提交信息以触发双语 Release |
|
||||
|
||||
@@ -18,7 +18,8 @@ sidebar: false
|
||||
|
||||
### 变更
|
||||
|
||||
- `of_node_access_logs` 从 PostgreSQL/SQLite 迁移至 ClickHouse(数据库 `openflare`);系统启动时强依赖 ClickHouse 连接与表结构自动初始化。
|
||||
- `of_node_access_logs` 从 PostgreSQL/SQLite 迁移至 ClickHouse(数据库 `openflare`);ClickHouse 表结构改由 goose 独立迁移管线管理(`goose_clickhouse_version`),启动时在 `migrator.MigrateClickHouse()` 中执行,移除 `clickhouse_schema.go` 手写 DDL。
|
||||
- ClickHouse 分析库接入 GORM(`db.ChDB`)与 `internal/repository/analytics/`:用户访问日志(`w_user_access_logs`)读写、风控批量写入与管理端日志统计 API 统一经 repository 访问;新增 `internal/model/analytics/` 分析表模型。
|
||||
- ClickHouse 默认启用(`clickhouse.enabled: true`),`docker-compose` 默认启动 `clickhouse` 服务并纳入 `wavelet` 健康依赖。
|
||||
|
||||
### 移除
|
||||
|
||||
@@ -51,6 +51,7 @@ require (
|
||||
golang.org/x/oauth2 v0.36.0
|
||||
golang.org/x/sync v0.21.0
|
||||
gopkg.in/natefinch/lumberjack.v2 v2.2.1
|
||||
gorm.io/driver/clickhouse v0.7.0
|
||||
gorm.io/driver/postgres v1.6.0
|
||||
gorm.io/driver/sqlite v1.6.0
|
||||
gorm.io/gorm v1.31.1
|
||||
@@ -180,7 +181,6 @@ require (
|
||||
google.golang.org/grpc v1.80.0 // indirect
|
||||
google.golang.org/protobuf v1.36.11 // indirect
|
||||
gopkg.in/yaml.v3 v3.0.1 // indirect
|
||||
gorm.io/driver/clickhouse v0.7.0 // indirect
|
||||
gorm.io/driver/mysql v1.6.0 // indirect
|
||||
modernc.org/libc v1.72.1 // indirect
|
||||
modernc.org/mathutil v1.7.1 // indirect
|
||||
|
||||
@@ -10,14 +10,14 @@ import (
|
||||
"encoding/json"
|
||||
"fmt"
|
||||
"net/http"
|
||||
"sort"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/admin"
|
||||
"github.com/Rain-kl/Wavelet/internal/config"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics"
|
||||
"github.com/Rain-kl/Wavelet/pkg/logger"
|
||||
"github.com/gin-gonic/gin"
|
||||
|
||||
@@ -30,7 +30,7 @@ const (
|
||||
maxPageSize = 100
|
||||
hoursInDay = 24
|
||||
analyticsDays = 7
|
||||
queryExtraArgs = 2 // pageSize + offset
|
||||
topActiveLimit = 10
|
||||
)
|
||||
|
||||
// logsResponse 历史日志查询响应
|
||||
@@ -158,115 +158,75 @@ type accessLogsResponse struct {
|
||||
List []accessLogItem `json:"list"`
|
||||
}
|
||||
|
||||
// buildAccessLogFilters 构建 ClickHouse 访问日志查询过滤条件
|
||||
func buildAccessLogFilters(ctx context.Context, c *gin.Context) ([]string, []interface{}, []uint64, error) {
|
||||
var conditions []string
|
||||
var args []interface{}
|
||||
var userIDs []uint64
|
||||
func buildAccessLogFilter(ctx context.Context, c *gin.Context) (analyticsrepo.AccessLogFilter, error) {
|
||||
filter := analyticsrepo.AccessLogFilter{}
|
||||
|
||||
// 按用户名过滤
|
||||
username := c.Query("username")
|
||||
if username != "" {
|
||||
var userIDs []uint64
|
||||
err := db.DB(ctx).Model(&model.User{}).
|
||||
Where("username LIKE ?", "%"+username+"%").
|
||||
Pluck("id", &userIDs).Error
|
||||
if err != nil {
|
||||
return nil, nil, nil, fmt.Errorf("查询用户信息失败: %w", err)
|
||||
return filter, fmt.Errorf("查询用户信息失败: %w", err)
|
||||
}
|
||||
if len(userIDs) == 0 {
|
||||
return nil, nil, nil, nil // 无匹配用户
|
||||
}
|
||||
}
|
||||
|
||||
if len(userIDs) > 0 {
|
||||
placeholders := make([]string, len(userIDs))
|
||||
for i := range userIDs {
|
||||
placeholders[i] = "?"
|
||||
args = append(args, userIDs[i])
|
||||
}
|
||||
conditions = append(conditions, fmt.Sprintf("user_id IN (%s)", strings.Join(placeholders, ",")))
|
||||
filter.UserIDs = userIDs
|
||||
}
|
||||
|
||||
if path := c.Query("path"); path != "" {
|
||||
conditions = append(conditions, "path LIKE ?")
|
||||
args = append(args, "%"+path+"%")
|
||||
filter.Path = path
|
||||
}
|
||||
|
||||
if startTime := c.Query("start_time"); startTime != "" {
|
||||
if t, err := time.Parse(time.RFC3339, startTime); err == nil {
|
||||
conditions = append(conditions, "created_at >= ?")
|
||||
args = append(args, t)
|
||||
} else if t, err := time.Parse("2006-01-02 15:04:05", startTime); err == nil {
|
||||
conditions = append(conditions, "created_at >= ?")
|
||||
args = append(args, t)
|
||||
if t, err := parseAccessLogTime(startTime); err == nil {
|
||||
filter.StartTime = &t
|
||||
}
|
||||
}
|
||||
|
||||
if endTime := c.Query("end_time"); endTime != "" {
|
||||
if t, err := time.Parse(time.RFC3339, endTime); err == nil {
|
||||
conditions = append(conditions, "created_at <= ?")
|
||||
args = append(args, t)
|
||||
} else if t, err := time.Parse("2006-01-02 15:04:05", endTime); err == nil {
|
||||
conditions = append(conditions, "created_at <= ?")
|
||||
args = append(args, t)
|
||||
if t, err := parseAccessLogTime(endTime); err == nil {
|
||||
filter.EndTime = &t
|
||||
}
|
||||
}
|
||||
|
||||
return conditions, args, userIDs, nil
|
||||
return filter, nil
|
||||
}
|
||||
|
||||
// fetchAccessLogDetails 查询 ClickHouse 访问日志明细并填充用户名
|
||||
func fetchAccessLogDetails(ctx context.Context, whereClause string, args []interface{}, pageSize int, offset int) ([]accessLogItem, error) {
|
||||
dataQuery := fmt.Sprintf(`
|
||||
SELECT id, user_id, path, method, ip, user_agent, headers, status, latency, created_at
|
||||
FROM w_user_access_logs
|
||||
%s
|
||||
ORDER BY created_at DESC, id DESC
|
||||
LIMIT ? OFFSET ?
|
||||
`, whereClause)
|
||||
|
||||
selectArgs := make([]interface{}, len(args), len(args)+queryExtraArgs)
|
||||
copy(selectArgs, args)
|
||||
selectArgs = append(selectArgs, pageSize, offset)
|
||||
|
||||
rows, err := db.ChConn.Query(ctx, dataQuery, selectArgs...)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("查询 ClickHouse 日志明细失败: %w", err)
|
||||
func parseAccessLogTime(value string) (time.Time, error) {
|
||||
if t, err := time.Parse(time.RFC3339, value); err == nil {
|
||||
return t, nil
|
||||
}
|
||||
defer func() { _ = rows.Close() }()
|
||||
return time.Parse("2006-01-02 15:04:05", value)
|
||||
}
|
||||
|
||||
var list []accessLogItem
|
||||
var fetchUserIDs []uint64
|
||||
|
||||
for rows.Next() {
|
||||
var item accessLogItem
|
||||
var createdAt time.Time
|
||||
if err := rows.Scan(&item.ID, &item.UserID, &item.Path, &item.Method, &item.IP, &item.UserAgent, &item.Headers, &item.Status, &item.Latency, &createdAt); err != nil {
|
||||
return nil, fmt.Errorf("读取 ClickHouse 结果失败: %w", err)
|
||||
}
|
||||
item.CreatedAt = createdAt.Format(time.RFC3339)
|
||||
list = append(list, item)
|
||||
fetchUserIDs = append(fetchUserIDs, item.UserID)
|
||||
func enrichAccessLogsWithUsers(ctx context.Context, list []accessLogItem) {
|
||||
if len(list) == 0 {
|
||||
return
|
||||
}
|
||||
|
||||
// 反查 Postgres 关联 Username 和 Nickname
|
||||
if len(fetchUserIDs) > 0 {
|
||||
userMap := make(map[uint64]struct{ Username, Nickname string })
|
||||
var users []model.User
|
||||
if err := db.DB(ctx).Where("id IN ?", fetchUserIDs).Find(&users).Error; err == nil {
|
||||
for _, u := range users {
|
||||
userMap[u.ID] = struct{ Username, Nickname string }{Username: u.Username, Nickname: u.Nickname}
|
||||
}
|
||||
}
|
||||
for i := range list {
|
||||
if info, ok := userMap[list[i].UserID]; ok {
|
||||
list[i].Username = info.Username
|
||||
list[i].Nickname = info.Nickname
|
||||
}
|
||||
userIDs := make([]uint64, 0, len(list))
|
||||
seen := make(map[uint64]struct{}, len(list))
|
||||
for _, item := range list {
|
||||
if _, ok := seen[item.UserID]; ok {
|
||||
continue
|
||||
}
|
||||
seen[item.UserID] = struct{}{}
|
||||
userIDs = append(userIDs, item.UserID)
|
||||
}
|
||||
|
||||
return list, nil
|
||||
userMap := make(map[uint64]struct{ Username, Nickname string })
|
||||
var users []model.User
|
||||
if err := db.DB(ctx).Where("id IN ?", userIDs).Find(&users).Error; err == nil {
|
||||
for _, u := range users {
|
||||
userMap[u.ID] = struct{ Username, Nickname string }{Username: u.Username, Nickname: u.Nickname}
|
||||
}
|
||||
}
|
||||
for i := range list {
|
||||
if info, ok := userMap[list[i].UserID]; ok {
|
||||
list[i].Username = info.Username
|
||||
list[i].Nickname = info.Nickname
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// GetAccessLogs 获取 ClickHouse 异步采集的访问日志
|
||||
@@ -287,12 +247,12 @@ func fetchAccessLogDetails(ctx context.Context, whereClause string, args []inter
|
||||
// @Failure 403 {object} response.Any "无管理员权限"
|
||||
// @Router /api/v1/admin/logs/access [get]
|
||||
func GetAccessLogs(c *gin.Context) {
|
||||
if db.ChConn == nil {
|
||||
response.AbortWithError(c, http.StatusInternalServerError, "ClickHouse 未初始化,无法检索访问日志")
|
||||
ctx := c.Request.Context()
|
||||
if !config.Config.ClickHouse.Enabled || db.ChDB(ctx) == nil {
|
||||
response.AbortWithError(c, http.StatusBadRequest, "ClickHouse 存储服务未启用,无法检索访问日志")
|
||||
return
|
||||
}
|
||||
|
||||
// 2. 解析分页参数
|
||||
page, _ := strconv.Atoi(c.DefaultQuery("page", "1"))
|
||||
if page < 1 {
|
||||
page = 1
|
||||
@@ -304,29 +264,20 @@ func GetAccessLogs(c *gin.Context) {
|
||||
if pageSize > maxPageSize {
|
||||
pageSize = maxPageSize
|
||||
}
|
||||
offset := (page - 1) * pageSize
|
||||
|
||||
// 3. 构建过滤条件
|
||||
conditions, args, userIDs, err := buildAccessLogFilters(c.Request.Context(), c)
|
||||
filter, err := buildAccessLogFilter(ctx, c)
|
||||
if err != nil {
|
||||
response.AbortWithError(c, http.StatusInternalServerError, err.Error())
|
||||
return
|
||||
}
|
||||
if userIDs != nil && len(userIDs) == 0 {
|
||||
if filter.UserIDs != nil && len(filter.UserIDs) == 0 {
|
||||
c.JSON(http.StatusOK, response.OK(accessLogsResponse{Total: 0, List: []accessLogItem{}}))
|
||||
return
|
||||
}
|
||||
|
||||
whereClause := ""
|
||||
if len(conditions) > 0 {
|
||||
whereClause = "WHERE " + strings.Join(conditions, " AND ")
|
||||
}
|
||||
|
||||
// 4. 查询日志总数
|
||||
var total uint64
|
||||
countQuery := fmt.Sprintf("SELECT count() FROM w_user_access_logs %s", whereClause)
|
||||
if err := db.ChConn.QueryRow(c.Request.Context(), countQuery, args...).Scan(&total); err != nil {
|
||||
response.AbortWithError(c, http.StatusInternalServerError, "查询 ClickHouse 日志统计失败: "+err.Error())
|
||||
logs, total, err := analyticsrepo.ListAccessLogs(ctx, filter, page, pageSize)
|
||||
if err != nil {
|
||||
response.AbortWithError(c, http.StatusInternalServerError, err.Error())
|
||||
return
|
||||
}
|
||||
if total == 0 {
|
||||
@@ -334,12 +285,22 @@ func GetAccessLogs(c *gin.Context) {
|
||||
return
|
||||
}
|
||||
|
||||
// 5. 分页查询明细数据
|
||||
list, err := fetchAccessLogDetails(c.Request.Context(), whereClause, args, pageSize, offset)
|
||||
if err != nil {
|
||||
response.AbortWithError(c, http.StatusInternalServerError, err.Error())
|
||||
return
|
||||
list := make([]accessLogItem, len(logs))
|
||||
for i, logItem := range logs {
|
||||
list[i] = accessLogItem{
|
||||
ID: logItem.ID,
|
||||
UserID: logItem.UserID,
|
||||
Path: logItem.Path,
|
||||
Method: logItem.Method,
|
||||
IP: logItem.IP,
|
||||
UserAgent: logItem.UserAgent,
|
||||
Headers: logItem.Headers,
|
||||
Status: logItem.Status,
|
||||
Latency: logItem.Latency,
|
||||
CreatedAt: logItem.CreatedAt.Format(time.RFC3339),
|
||||
}
|
||||
}
|
||||
enrichAccessLogsWithUsers(ctx, list)
|
||||
|
||||
c.JSON(http.StatusOK, response.OK(accessLogsResponse{
|
||||
Total: total,
|
||||
@@ -386,133 +347,61 @@ type logsAnalyticsResponse struct {
|
||||
// @Failure 403 {object} response.Any "无管理员权限"
|
||||
// @Router /api/v1/admin/logs/analytics [get]
|
||||
func GetLogsAnalytics(c *gin.Context) {
|
||||
if db.ChConn == nil {
|
||||
response.AbortWithError(c, http.StatusInternalServerError, "ClickHouse 未初始化,无法获取分析数据")
|
||||
ctx := c.Request.Context()
|
||||
if !config.Config.ClickHouse.Enabled || db.ChDB(ctx) == nil {
|
||||
response.AbortWithError(c, http.StatusBadRequest, "ClickHouse 存储服务未启用,无法获取分析数据")
|
||||
return
|
||||
}
|
||||
|
||||
ctx := c.Request.Context()
|
||||
// 7 天前 00:00:00
|
||||
startTime := time.Now().AddDate(0, 0, -(analyticsDays - 1)).Truncate(hoursInDay * time.Hour)
|
||||
|
||||
trendList := queryAccessTrend(ctx, startTime)
|
||||
browserList := queryBrowserDistribution(ctx, startTime)
|
||||
topUsers := queryTopActiveUsers(ctx, startTime)
|
||||
|
||||
c.JSON(http.StatusOK, response.OK(logsAnalyticsResponse{
|
||||
Trend: trendList,
|
||||
Browsers: browserList,
|
||||
TopUsers: topUsers,
|
||||
}))
|
||||
}
|
||||
|
||||
// queryAccessTrend 查询最近 7 天的访问趋势
|
||||
func queryAccessTrend(ctx context.Context, startTime time.Time) []trendItem {
|
||||
trendRows, err := db.ChConn.Query(ctx, `
|
||||
SELECT toDate(created_at) as date, count() as count
|
||||
FROM w_user_access_logs
|
||||
WHERE created_at >= ?
|
||||
GROUP BY date
|
||||
ORDER BY date ASC
|
||||
`, startTime)
|
||||
|
||||
trendMap := make(map[string]uint64)
|
||||
for i := 0; i < analyticsDays; i++ {
|
||||
dStr := time.Now().AddDate(0, 0, -i).Format("2006-01-02")
|
||||
trendMap[dStr] = 0
|
||||
trendPoints, err := analyticsrepo.GetDailyTrend(ctx, analyticsDays)
|
||||
if err != nil {
|
||||
response.AbortWithError(c, http.StatusInternalServerError, "查询访问趋势失败: "+err.Error())
|
||||
return
|
||||
}
|
||||
|
||||
if err == nil {
|
||||
defer func() { _ = trendRows.Close() }()
|
||||
for trendRows.Next() {
|
||||
var dt time.Time
|
||||
var cnt uint64
|
||||
if errScan := trendRows.Scan(&dt, &cnt); errScan == nil {
|
||||
dStr := dt.Format("2006-01-02")
|
||||
trendMap[dStr] = cnt
|
||||
}
|
||||
trendList := make([]trendItem, len(trendPoints))
|
||||
for i, point := range trendPoints {
|
||||
trendList[i] = trendItem{
|
||||
Date: point.Date,
|
||||
Count: point.Count,
|
||||
}
|
||||
}
|
||||
|
||||
var trendList []trendItem
|
||||
for i := analyticsDays - 1; i >= 0; i-- {
|
||||
dStr := time.Now().AddDate(0, 0, -i).Format("2006-01-02")
|
||||
trendList = append(trendList, trendItem{
|
||||
Date: dStr,
|
||||
Count: trendMap[dStr],
|
||||
})
|
||||
browserPoints, err := analyticsrepo.GetBrowserDistribution(ctx, startTime)
|
||||
if err != nil {
|
||||
response.AbortWithError(c, http.StatusInternalServerError, "查询浏览器分布失败: "+err.Error())
|
||||
return
|
||||
}
|
||||
return trendList
|
||||
}
|
||||
|
||||
// queryBrowserDistribution 查询浏览器分布排行
|
||||
func queryBrowserDistribution(ctx context.Context, startTime time.Time) []browserItem {
|
||||
uaRows, err := db.ChConn.Query(ctx, `
|
||||
SELECT user_agent, count() as count
|
||||
FROM w_user_access_logs
|
||||
WHERE created_at >= ?
|
||||
GROUP BY user_agent
|
||||
`, startTime)
|
||||
|
||||
browserCounts := make(map[string]uint64)
|
||||
if err == nil {
|
||||
defer func() { _ = uaRows.Close() }()
|
||||
for uaRows.Next() {
|
||||
var ua string
|
||||
var cnt uint64
|
||||
if errScan := uaRows.Scan(&ua, &cnt); errScan == nil {
|
||||
browser := parseBrowserName(ua)
|
||||
browserCounts[browser] += cnt
|
||||
}
|
||||
browserList := make([]browserItem, len(browserPoints))
|
||||
for i, point := range browserPoints {
|
||||
browserList[i] = browserItem{
|
||||
Browser: point.Browser,
|
||||
Count: point.Count,
|
||||
}
|
||||
}
|
||||
|
||||
var browserList []browserItem
|
||||
for b, cnt := range browserCounts {
|
||||
browserList = append(browserList, browserItem{
|
||||
Browser: b,
|
||||
Count: cnt,
|
||||
})
|
||||
topUserPoints, err := analyticsrepo.GetTopActiveUsers(ctx, startTime, topActiveLimit)
|
||||
if err != nil {
|
||||
response.AbortWithError(c, http.StatusInternalServerError, "查询活跃用户失败: "+err.Error())
|
||||
return
|
||||
}
|
||||
sort.Slice(browserList, func(i, j int) bool {
|
||||
return browserList[i].Count > browserList[j].Count
|
||||
})
|
||||
return browserList
|
||||
}
|
||||
|
||||
// queryTopActiveUsers 查询活跃用户 Top 10
|
||||
func queryTopActiveUsers(ctx context.Context, startTime time.Time) []topUserItem {
|
||||
userRows, err := db.ChConn.Query(ctx, `
|
||||
SELECT user_id, count() as count
|
||||
FROM w_user_access_logs
|
||||
WHERE created_at >= ? AND user_id > 0
|
||||
GROUP BY user_id
|
||||
ORDER BY count DESC
|
||||
LIMIT 10
|
||||
`, startTime)
|
||||
|
||||
var topUsers []topUserItem
|
||||
var userIDs []uint64
|
||||
userCountMap := make(map[uint64]uint64)
|
||||
|
||||
if err == nil {
|
||||
defer func() { _ = userRows.Close() }()
|
||||
for userRows.Next() {
|
||||
var uid uint64
|
||||
var cnt uint64
|
||||
if errScan := userRows.Scan(&uid, &cnt); errScan == nil {
|
||||
userIDs = append(userIDs, uid)
|
||||
userCountMap[uid] = cnt
|
||||
}
|
||||
topUsers := make([]topUserItem, len(topUserPoints))
|
||||
userIDs := make([]uint64, len(topUserPoints))
|
||||
for i, point := range topUserPoints {
|
||||
topUsers[i] = topUserItem{
|
||||
UserID: point.UserID,
|
||||
Count: point.Count,
|
||||
}
|
||||
userIDs[i] = point.UserID
|
||||
}
|
||||
|
||||
// 反查 Postgres 补全活跃用户的用户名和昵称
|
||||
userProfileMap := make(map[uint64]struct {
|
||||
Username string
|
||||
Nickname string
|
||||
})
|
||||
if len(userIDs) > 0 {
|
||||
userProfileMap := make(map[uint64]struct {
|
||||
Username string
|
||||
Nickname string
|
||||
})
|
||||
var users []model.User
|
||||
if errProfile := db.DB(ctx).Where("id IN ?", userIDs).Find(&users).Error; errProfile == nil {
|
||||
for _, u := range users {
|
||||
@@ -525,40 +414,17 @@ func queryTopActiveUsers(ctx context.Context, startTime time.Time) []topUserItem
|
||||
}
|
||||
}
|
||||
}
|
||||
for i := range topUsers {
|
||||
if profile, ok := userProfileMap[topUsers[i].UserID]; ok {
|
||||
topUsers[i].Username = profile.Username
|
||||
topUsers[i].Nickname = profile.Nickname
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
for _, uid := range userIDs {
|
||||
profile := userProfileMap[uid]
|
||||
topUsers = append(topUsers, topUserItem{
|
||||
UserID: uid,
|
||||
Username: profile.Username,
|
||||
Nickname: profile.Nickname,
|
||||
Count: userCountMap[uid],
|
||||
})
|
||||
}
|
||||
return topUsers
|
||||
}
|
||||
|
||||
// parseBrowserName 简易的 User-Agent 浏览器类型识别
|
||||
func parseBrowserName(ua string) string {
|
||||
uaLower := strings.ToLower(ua)
|
||||
if strings.Contains(uaLower, "micromessenger") {
|
||||
return "WeChat"
|
||||
}
|
||||
if strings.Contains(uaLower, "postman") {
|
||||
return "Postman"
|
||||
}
|
||||
if strings.Contains(uaLower, "edg/") || strings.Contains(uaLower, "edge") {
|
||||
return "Edge"
|
||||
}
|
||||
if strings.Contains(uaLower, "firefox") {
|
||||
return "Firefox"
|
||||
}
|
||||
if strings.Contains(uaLower, "chrome") {
|
||||
return "Chrome"
|
||||
}
|
||||
if strings.Contains(uaLower, "safari") {
|
||||
return "Safari"
|
||||
}
|
||||
return "Other"
|
||||
c.JSON(http.StatusOK, response.OK(logsAnalyticsResponse{
|
||||
Trend: trendList,
|
||||
Browsers: browserList,
|
||||
TopUsers: topUsers,
|
||||
}))
|
||||
}
|
||||
|
||||
@@ -8,11 +8,12 @@ import (
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/config"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
"github.com/Rain-kl/Wavelet/internal/model/analytics"
|
||||
analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics"
|
||||
"github.com/Rain-kl/Wavelet/pkg/logger"
|
||||
)
|
||||
|
||||
var logChan chan *UserAccessLog
|
||||
var logChan chan *analytics.UserAccessLog
|
||||
|
||||
const (
|
||||
defaultQueueSize = 10000
|
||||
@@ -26,7 +27,7 @@ func InitLogWriter(ctx context.Context) {
|
||||
return
|
||||
}
|
||||
|
||||
logChan = make(chan *UserAccessLog, defaultQueueSize)
|
||||
logChan = make(chan *analytics.UserAccessLog, defaultQueueSize)
|
||||
go startBatchWorker(context.WithoutCancel(ctx))
|
||||
}
|
||||
|
||||
@@ -40,7 +41,7 @@ func IsBufferFull() bool {
|
||||
}
|
||||
|
||||
// QueueAccessLog 异步非阻塞地将日志推入缓冲队列
|
||||
func QueueAccessLog(logItem *UserAccessLog) {
|
||||
func QueueAccessLog(logItem *analytics.UserAccessLog) {
|
||||
if !config.Config.ClickHouse.Enabled || logChan == nil {
|
||||
return
|
||||
}
|
||||
@@ -56,43 +57,18 @@ func startBatchWorker(ctx context.Context) {
|
||||
ticker := time.NewTicker(flushInterval)
|
||||
defer ticker.Stop()
|
||||
|
||||
var batch []*UserAccessLog
|
||||
var batch []*analytics.UserAccessLog
|
||||
|
||||
flush := func() {
|
||||
if len(batch) == 0 {
|
||||
return
|
||||
}
|
||||
if db.ChConn == nil {
|
||||
batch = nil
|
||||
return
|
||||
}
|
||||
|
||||
b, err := db.ChConn.PrepareBatch(ctx, "INSERT INTO w_user_access_logs (id, user_id, path, method, ip, user_agent, headers, status, latency, created_at)")
|
||||
if err != nil {
|
||||
logger.ErrorF(ctx, "[RiskControl] Prepare ClickHouse batch failed: %v", err)
|
||||
batch = nil
|
||||
return
|
||||
items := make([]analytics.UserAccessLog, len(batch))
|
||||
for i, item := range batch {
|
||||
items[i] = *item
|
||||
}
|
||||
|
||||
for _, item := range batch {
|
||||
err = b.Append(
|
||||
item.ID,
|
||||
item.UserID,
|
||||
item.Path,
|
||||
item.Method,
|
||||
item.IP,
|
||||
item.UserAgent,
|
||||
item.Headers,
|
||||
item.Status,
|
||||
item.Latency,
|
||||
item.CreatedAt,
|
||||
)
|
||||
if err != nil {
|
||||
logger.ErrorF(ctx, "[RiskControl] Append item to ClickHouse batch failed: %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
if err := b.Send(); err != nil {
|
||||
if err := analyticsrepo.BatchInsert(ctx, items); err != nil {
|
||||
logger.ErrorF(ctx, "[RiskControl] Send ClickHouse batch failed: %v", err)
|
||||
}
|
||||
batch = nil
|
||||
@@ -113,4 +89,4 @@ func startBatchWorker(ctx context.Context) {
|
||||
flush()
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -14,6 +14,7 @@ import (
|
||||
"github.com/Rain-kl/Wavelet/internal/config"
|
||||
"github.com/Rain-kl/Wavelet/internal/db/idgen"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/model/analytics"
|
||||
"github.com/gin-gonic/gin"
|
||||
)
|
||||
|
||||
@@ -68,7 +69,7 @@ func RiskControlMiddleware() gin.HandlerFunc {
|
||||
status = maxHTTPStatus
|
||||
}
|
||||
|
||||
logItem := &UserAccessLog{
|
||||
logItem := &analytics.UserAccessLog{
|
||||
ID: idgen.NextUint64ID(),
|
||||
UserID: userObj.ID, // 直接从 Context 获取已登录用户ID,避免数据库查询
|
||||
Path: c.Request.URL.Path,
|
||||
|
||||
@@ -13,6 +13,7 @@ import (
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/oauth"
|
||||
"github.com/Rain-kl/Wavelet/internal/config"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/model/analytics"
|
||||
"github.com/Rain-kl/Wavelet/internal/testhelper"
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/stretchr/testify/assert"
|
||||
@@ -40,7 +41,7 @@ func TestRiskControlMiddleware(t *testing.T) {
|
||||
|
||||
t.Run("ClickHouse enabled - Normal Authenticated Request", func(t *testing.T) {
|
||||
config.Config.ClickHouse.Enabled = true
|
||||
logChan = make(chan *UserAccessLog, defaultQueueSize)
|
||||
logChan = make(chan *analytics.UserAccessLog, defaultQueueSize)
|
||||
defer func() {
|
||||
config.Config.ClickHouse.Enabled = false
|
||||
logChan = nil
|
||||
@@ -84,7 +85,7 @@ func TestRiskControlMiddleware(t *testing.T) {
|
||||
|
||||
t.Run("ClickHouse enabled - Unauthenticated Request", func(t *testing.T) {
|
||||
config.Config.ClickHouse.Enabled = true
|
||||
logChan = make(chan *UserAccessLog, defaultQueueSize)
|
||||
logChan = make(chan *analytics.UserAccessLog, defaultQueueSize)
|
||||
defer func() {
|
||||
config.Config.ClickHouse.Enabled = false
|
||||
logChan = nil
|
||||
@@ -113,7 +114,7 @@ func TestRiskControlMiddleware(t *testing.T) {
|
||||
|
||||
t.Run("ClickHouse enabled - Buffer Full Rate Limiting", func(t *testing.T) {
|
||||
config.Config.ClickHouse.Enabled = true
|
||||
logChan = make(chan *UserAccessLog, 2) // small capacity for quick fill
|
||||
logChan = make(chan *analytics.UserAccessLog, 2) // small capacity for quick fill
|
||||
defer func() {
|
||||
config.Config.ClickHouse.Enabled = false
|
||||
logChan = nil
|
||||
@@ -121,7 +122,7 @@ func TestRiskControlMiddleware(t *testing.T) {
|
||||
|
||||
// fill logChan up to cap to simulate buffer full
|
||||
for len(logChan) < cap(logChan) {
|
||||
logChan <- &UserAccessLog{}
|
||||
logChan <- &analytics.UserAccessLog{}
|
||||
}
|
||||
|
||||
r := testhelper.NewTestGinEngine(RiskControlMiddleware())
|
||||
|
||||
@@ -1,22 +0,0 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package risk_control
|
||||
|
||||
import (
|
||||
"time"
|
||||
)
|
||||
|
||||
// UserAccessLog 用户访问记录
|
||||
type UserAccessLog struct {
|
||||
ID uint64 `json:"id,string"`
|
||||
UserID uint64 `json:"user_id,string"`
|
||||
Path string `json:"path"`
|
||||
Method string `json:"method"`
|
||||
IP string `json:"ip"`
|
||||
UserAgent string `json:"user_agent"`
|
||||
Headers string `json:"headers"`
|
||||
Status int32 `json:"status"`
|
||||
Latency int64 `json:"latency"` // 耗时毫秒
|
||||
CreatedAt time.Time `json:"created_at"`
|
||||
}
|
||||
@@ -40,6 +40,7 @@ var rootCmd = &cobra.Command{
|
||||
},
|
||||
PreRun: func(_ *cobra.Command, _ []string) {
|
||||
migrator.Migrate()
|
||||
migrator.MigrateClickHouse()
|
||||
},
|
||||
PersistentPostRun: func(_ *cobra.Command, _ []string) {
|
||||
shutdownTraceProvider()
|
||||
|
||||
@@ -7,12 +7,20 @@ package db
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log"
|
||||
"net/url"
|
||||
"strconv"
|
||||
"strings"
|
||||
"time"
|
||||
|
||||
"github.com/ClickHouse/clickhouse-go/v2"
|
||||
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
|
||||
"github.com/Rain-kl/Wavelet/internal/config"
|
||||
"go.opentelemetry.io/otel/attribute"
|
||||
clickhouseDriver "gorm.io/driver/clickhouse"
|
||||
"gorm.io/gorm"
|
||||
"gorm.io/plugin/opentelemetry/tracing"
|
||||
)
|
||||
|
||||
const (
|
||||
@@ -21,8 +29,10 @@ const (
|
||||
)
|
||||
|
||||
var (
|
||||
// ChConn ClickHouse 连接实例
|
||||
// ChConn ClickHouse 原生连接实例,用于批量写入
|
||||
ChConn driver.Conn
|
||||
|
||||
chDB *gorm.DB
|
||||
)
|
||||
|
||||
func init() {
|
||||
@@ -35,10 +45,55 @@ func init() {
|
||||
log.Fatalf("[ClickHouse] database name is required (expected: openflare)\n")
|
||||
}
|
||||
|
||||
var err error
|
||||
opts := buildClickHouseOptions()
|
||||
|
||||
// 配置 ClickHouse 连接
|
||||
ChConn, err = clickhouse.Open(&clickhouse.Options{
|
||||
var err error
|
||||
ChConn, err = clickhouse.Open(opts)
|
||||
if err != nil {
|
||||
log.Fatalf("[ClickHouse] init connection failed: %v\n", err)
|
||||
}
|
||||
|
||||
if err = ChConn.Ping(context.Background()); err != nil {
|
||||
log.Fatalf("[ClickHouse] ping failed: %v\n", err)
|
||||
}
|
||||
|
||||
chDB, err = gorm.Open(clickhouseDriver.New(clickhouseDriver.Config{
|
||||
DSN: buildClickHouseDSN(),
|
||||
}), &gorm.Config{
|
||||
SkipDefaultTransaction: true,
|
||||
})
|
||||
if err != nil {
|
||||
log.Fatalf("[ClickHouse] init gorm connection failed: %v\n", err)
|
||||
}
|
||||
|
||||
if err = chDB.Use(
|
||||
tracing.NewPlugin(
|
||||
tracing.WithoutMetrics(),
|
||||
tracing.WithAttributes(
|
||||
attribute.String("db.instance", cfg.Database),
|
||||
attribute.String("db.system", "ClickHouse"),
|
||||
),
|
||||
),
|
||||
); err != nil {
|
||||
log.Fatalf("[ClickHouse] init trace failed: %v\n", err)
|
||||
}
|
||||
|
||||
sqlDB, err := chDB.DB()
|
||||
if err != nil {
|
||||
log.Fatalf("[ClickHouse] load sql db failed: %v\n", err)
|
||||
}
|
||||
|
||||
sqlDB.SetMaxIdleConns(cfg.MaxIdleConn)
|
||||
sqlDB.SetMaxOpenConns(cfg.MaxOpenConn)
|
||||
sqlDB.SetConnMaxLifetime(time.Duration(cfg.ConnMaxLifetime) * time.Second)
|
||||
|
||||
log.Println("[ClickHouse] connection established successfully")
|
||||
}
|
||||
|
||||
func buildClickHouseOptions() *clickhouse.Options {
|
||||
cfg := config.Config.ClickHouse
|
||||
|
||||
return &clickhouse.Options{
|
||||
Addr: cfg.Hosts,
|
||||
Auth: clickhouse.Auth{
|
||||
Database: cfg.Database,
|
||||
@@ -57,17 +112,44 @@ func init() {
|
||||
ConnMaxLifetime: time.Duration(cfg.ConnMaxLifetime) * time.Second,
|
||||
ReadTimeout: time.Duration(cfg.DialTimeout*clickhouseReadTimeoutFactor) * time.Second,
|
||||
BlockBufferSize: cfg.BlockBufferSize,
|
||||
})
|
||||
|
||||
if err != nil {
|
||||
log.Fatalf("[ClickHouse] init connection failed: %v\n", err)
|
||||
}
|
||||
|
||||
// 测试连接
|
||||
if err = ChConn.Ping(context.Background()); err != nil {
|
||||
log.Fatalf("[ClickHouse] ping failed: %v\n", err)
|
||||
}
|
||||
|
||||
ensureClickHouseSchemaOnStartup()
|
||||
log.Println("[ClickHouse] connection established successfully")
|
||||
}
|
||||
|
||||
func buildClickHouseDSN() string {
|
||||
cfg := config.Config.ClickHouse
|
||||
|
||||
chURL := &url.URL{
|
||||
Scheme: "clickhouse",
|
||||
Host: strings.Join(cfg.Hosts, ","),
|
||||
Path: "/" + cfg.Database,
|
||||
}
|
||||
if cfg.Username != "" || cfg.Password != "" {
|
||||
chURL.User = url.UserPassword(cfg.Username, cfg.Password)
|
||||
}
|
||||
|
||||
query := chURL.Query()
|
||||
query.Set("dial_timeout", fmt.Sprintf("%ds", cfg.DialTimeout))
|
||||
query.Set("read_timeout", fmt.Sprintf("%ds", cfg.DialTimeout*clickhouseReadTimeoutFactor))
|
||||
query.Set("max_execution_time", strconv.Itoa(clickhouseMaxExecTime))
|
||||
chURL.RawQuery = query.Encode()
|
||||
|
||||
return chURL.String()
|
||||
}
|
||||
|
||||
// ChDB returns a context-aware GORM ClickHouse instance.
|
||||
func ChDB(ctx context.Context) *gorm.DB {
|
||||
if chDB == nil {
|
||||
return nil
|
||||
}
|
||||
return chDB.WithContext(ctx)
|
||||
}
|
||||
|
||||
// SetChDBForTest sets the package-level ClickHouse GORM instance for testing.
|
||||
func SetChDBForTest(d *gorm.DB) {
|
||||
chDB = d
|
||||
}
|
||||
|
||||
// SetChConnForTest sets the package-level native ClickHouse connection for testing.
|
||||
func SetChConnForTest(c driver.Conn) {
|
||||
ChConn = c
|
||||
}
|
||||
@@ -1,51 +0,0 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package db
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"log"
|
||||
)
|
||||
|
||||
const openFlareNodeAccessLogsDDL = `
|
||||
CREATE TABLE IF NOT EXISTS of_node_access_logs
|
||||
(
|
||||
id UInt64,
|
||||
node_id String,
|
||||
logged_at DateTime64(3, 'UTC'),
|
||||
remote_addr String,
|
||||
region String,
|
||||
host String,
|
||||
path String,
|
||||
status_code Int32,
|
||||
created_at DateTime64(3, 'UTC')
|
||||
)
|
||||
ENGINE = MergeTree()
|
||||
PARTITION BY toYYYYMM(logged_at)
|
||||
ORDER BY (node_id, logged_at, remote_addr, host, path, status_code)
|
||||
SETTINGS index_granularity = 8192`
|
||||
|
||||
// EnsureClickHouseSchema creates the openflare database and required tables.
|
||||
func EnsureClickHouseSchema(ctx context.Context) error {
|
||||
if ChConn == nil {
|
||||
return fmt.Errorf("clickhouse connection is not initialized")
|
||||
}
|
||||
|
||||
if err := ChConn.Exec(ctx, "CREATE DATABASE IF NOT EXISTS openflare"); err != nil {
|
||||
return fmt.Errorf("create database openflare: %w", err)
|
||||
}
|
||||
if err := ChConn.Exec(ctx, openFlareNodeAccessLogsDDL); err != nil {
|
||||
return fmt.Errorf("create table of_node_access_logs: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
func ensureClickHouseSchemaOnStartup() {
|
||||
ctx := context.Background()
|
||||
if err := EnsureClickHouseSchema(ctx); err != nil {
|
||||
log.Fatalf("[ClickHouse] ensure schema failed: %v\n", err)
|
||||
}
|
||||
log.Println("[ClickHouse] schema ready (database: openflare)")
|
||||
}
|
||||
@@ -0,0 +1,77 @@
|
||||
// Copyright 2025 linux.do
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package migrator
|
||||
|
||||
import (
|
||||
"database/sql"
|
||||
"embed"
|
||||
"log"
|
||||
"time"
|
||||
|
||||
"github.com/ClickHouse/clickhouse-go/v2"
|
||||
"github.com/Rain-kl/Wavelet/internal/config"
|
||||
"github.com/pressly/goose/v3"
|
||||
)
|
||||
|
||||
const (
|
||||
clickhouseMigrationDir = "goose/clickhouse"
|
||||
clickhouseGooseVersionTable = "goose_clickhouse_version"
|
||||
clickhouseMaxExecTime = 60
|
||||
clickhouseReadTimeoutFactor = 2
|
||||
)
|
||||
|
||||
// clickhouseMigrationFS contains SQL migrations under goose/clickhouse.
|
||||
//
|
||||
//go:embed goose/clickhouse/*.sql
|
||||
var clickhouseMigrationFS embed.FS
|
||||
|
||||
// MigrateClickHouse runs goose migrations against ClickHouse when enabled.
|
||||
func MigrateClickHouse() {
|
||||
if !config.Config.ClickHouse.Enabled {
|
||||
return
|
||||
}
|
||||
|
||||
cfg := config.Config.ClickHouse
|
||||
sqlDB := clickhouse.OpenDB(&clickhouse.Options{
|
||||
Addr: cfg.Hosts,
|
||||
Auth: clickhouse.Auth{
|
||||
Database: cfg.Database,
|
||||
Username: cfg.Username,
|
||||
Password: cfg.Password,
|
||||
},
|
||||
Settings: clickhouse.Settings{
|
||||
"max_execution_time": clickhouseMaxExecTime,
|
||||
},
|
||||
Compression: &clickhouse.Compression{
|
||||
Method: clickhouse.CompressionLZ4,
|
||||
},
|
||||
DialTimeout: time.Duration(cfg.DialTimeout) * time.Second,
|
||||
MaxOpenConns: cfg.MaxOpenConn,
|
||||
MaxIdleConns: cfg.MaxIdleConn,
|
||||
ConnMaxLifetime: time.Duration(cfg.ConnMaxLifetime) * time.Second,
|
||||
ReadTimeout: time.Duration(cfg.DialTimeout*clickhouseReadTimeoutFactor) * time.Second,
|
||||
BlockBufferSize: cfg.BlockBufferSize,
|
||||
})
|
||||
|
||||
goose.SetBaseFS(clickhouseMigrationFS)
|
||||
if err := goose.SetDialect("clickhouse"); err != nil {
|
||||
closeClickHouseDB(sqlDB)
|
||||
log.Fatalf("[ClickHouse] set goose dialect failed: %v\n", err)
|
||||
}
|
||||
goose.SetTableName(clickhouseGooseVersionTable)
|
||||
if err := goose.Up(sqlDB, clickhouseMigrationDir); err != nil {
|
||||
closeClickHouseDB(sqlDB)
|
||||
log.Fatalf("[ClickHouse] goose migrate failed: %v\n", err)
|
||||
}
|
||||
closeClickHouseDB(sqlDB)
|
||||
|
||||
log.Println("[ClickHouse] goose migrate success")
|
||||
}
|
||||
|
||||
func closeClickHouseDB(sqlDB *sql.DB) {
|
||||
if err := sqlDB.Close(); err != nil {
|
||||
log.Printf("[ClickHouse] close sql db failed: %v\n", err)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,55 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package migrator
|
||||
|
||||
import (
|
||||
"testing"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/config"
|
||||
"github.com/pressly/goose/v3"
|
||||
)
|
||||
|
||||
func TestClickHouseMigrationFilesEmbedded(t *testing.T) {
|
||||
entries, err := clickhouseMigrationFS.ReadDir(clickhouseMigrationDir)
|
||||
if err != nil {
|
||||
t.Fatalf("ReadDir(%q) error = %v", clickhouseMigrationDir, err)
|
||||
}
|
||||
if len(entries) == 0 {
|
||||
t.Fatal("expected embedded ClickHouse migrations, got none")
|
||||
}
|
||||
|
||||
expected := map[string]bool{
|
||||
"202606190001_create_user_access_logs.sql": false,
|
||||
"202606200001_create_node_access_logs.sql": false,
|
||||
}
|
||||
for _, entry := range entries {
|
||||
if entry.IsDir() {
|
||||
continue
|
||||
}
|
||||
if _, ok := expected[entry.Name()]; ok {
|
||||
expected[entry.Name()] = true
|
||||
}
|
||||
}
|
||||
for name, found := range expected {
|
||||
if !found {
|
||||
t.Fatalf("expected %s in embedded migrations", name)
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func TestClickHouseGooseDialect(t *testing.T) {
|
||||
if err := goose.SetDialect("clickhouse"); err != nil {
|
||||
t.Fatalf("SetDialect(clickhouse) error = %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
func TestMigrateClickHouseSkipsWhenDisabled(t *testing.T) {
|
||||
previousEnabled := config.Config.ClickHouse.Enabled
|
||||
config.Config.ClickHouse.Enabled = false
|
||||
t.Cleanup(func() {
|
||||
config.Config.ClickHouse.Enabled = previousEnabled
|
||||
})
|
||||
|
||||
MigrateClickHouse()
|
||||
}
|
||||
+4
-4
@@ -1,7 +1,4 @@
|
||||
CREATE DATABASE IF NOT EXISTS openflare;
|
||||
|
||||
USE openflare;
|
||||
|
||||
-- +goose Up
|
||||
CREATE TABLE IF NOT EXISTS w_user_access_logs
|
||||
(
|
||||
id UInt64,
|
||||
@@ -19,3 +16,6 @@ ENGINE = MergeTree()
|
||||
PARTITION BY toYYYYMM(created_at)
|
||||
ORDER BY (created_at, ip, user_id)
|
||||
SETTINGS index_granularity = 8192;
|
||||
|
||||
-- +goose Down
|
||||
DROP TABLE IF EXISTS w_user_access_logs;
|
||||
+20
@@ -0,0 +1,20 @@
|
||||
-- +goose Up
|
||||
CREATE TABLE IF NOT EXISTS of_node_access_logs
|
||||
(
|
||||
id UInt64,
|
||||
node_id String,
|
||||
logged_at DateTime64(3, 'UTC'),
|
||||
remote_addr String,
|
||||
region String,
|
||||
host String,
|
||||
path String,
|
||||
status_code Int32,
|
||||
created_at DateTime64(3, 'UTC')
|
||||
)
|
||||
ENGINE = MergeTree()
|
||||
PARTITION BY toYYYYMM(logged_at)
|
||||
ORDER BY (node_id, logged_at, remote_addr, host, path, status_code)
|
||||
SETTINGS index_granularity = 8192;
|
||||
|
||||
-- +goose Down
|
||||
DROP TABLE IF EXISTS of_node_access_logs;
|
||||
@@ -0,0 +1,42 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package analytics
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
const (
|
||||
nodeAccessLogTableName = "of_node_access_logs"
|
||||
nodeAccessLogInsertColumns = "id, node_id, logged_at, remote_addr, region, host, path, status_code, created_at"
|
||||
)
|
||||
|
||||
// NodeAccessLog stores OpenFlare edge node access records in ClickHouse.
|
||||
type NodeAccessLog struct {
|
||||
ID uint64 `gorm:"column:id"`
|
||||
NodeID string `gorm:"column:node_id"`
|
||||
LoggedAt time.Time `gorm:"column:logged_at"`
|
||||
RemoteAddr string `gorm:"column:remote_addr"`
|
||||
Region string `gorm:"column:region"`
|
||||
Host string `gorm:"column:host"`
|
||||
Path string `gorm:"column:path"`
|
||||
StatusCode int32 `gorm:"column:status_code"`
|
||||
CreatedAt time.Time `gorm:"column:created_at"`
|
||||
}
|
||||
|
||||
// TableName returns the ClickHouse table name.
|
||||
func (NodeAccessLog) TableName() string {
|
||||
return nodeAccessLogTableName
|
||||
}
|
||||
|
||||
// InsertColumns returns comma-separated column names for batch insert.
|
||||
func (NodeAccessLog) InsertColumns() string {
|
||||
return nodeAccessLogInsertColumns
|
||||
}
|
||||
|
||||
// BatchInsertSQL returns the INSERT prefix used by native batch writers.
|
||||
func (NodeAccessLog) BatchInsertSQL() string {
|
||||
return fmt.Sprintf("INSERT INTO %s (%s)", nodeAccessLogTableName, nodeAccessLogInsertColumns)
|
||||
}
|
||||
@@ -0,0 +1,44 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package analytics defines ClickHouse analytics domain models.
|
||||
package analytics
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"time"
|
||||
)
|
||||
|
||||
const (
|
||||
userAccessLogTableName = "w_user_access_logs"
|
||||
userAccessLogInsertColumns = "id, user_id, path, method, ip, user_agent, headers, status, latency, created_at"
|
||||
)
|
||||
|
||||
// UserAccessLog stores HTTP access records in ClickHouse.
|
||||
type UserAccessLog struct {
|
||||
ID uint64 `gorm:"column:id"`
|
||||
UserID uint64 `gorm:"column:user_id"`
|
||||
Path string `gorm:"column:path"`
|
||||
Method string `gorm:"column:method"`
|
||||
IP string `gorm:"column:ip"`
|
||||
UserAgent string `gorm:"column:user_agent"`
|
||||
Headers string `gorm:"column:headers"`
|
||||
Status int32 `gorm:"column:status"`
|
||||
Latency int64 `gorm:"column:latency"`
|
||||
CreatedAt time.Time `gorm:"column:created_at"`
|
||||
}
|
||||
|
||||
// TableName returns the ClickHouse table name.
|
||||
func (UserAccessLog) TableName() string {
|
||||
return userAccessLogTableName
|
||||
}
|
||||
|
||||
// InsertColumns returns comma-separated column names for batch insert.
|
||||
func (UserAccessLog) InsertColumns() string {
|
||||
return userAccessLogInsertColumns
|
||||
}
|
||||
|
||||
// BatchInsertSQL returns the INSERT prefix used by native batch writers.
|
||||
func (UserAccessLog) BatchInsertSQL() string {
|
||||
return fmt.Sprintf("INSERT INTO %s (%s)", userAccessLogTableName, userAccessLogInsertColumns)
|
||||
}
|
||||
@@ -0,0 +1,96 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package analytics provides ClickHouse data access for analytics tables.
|
||||
package analytics
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
// CountAccessLogs returns the number of access logs matching filter.
|
||||
func CountAccessLogs(ctx context.Context, filter AccessLogFilter) (uint64, error) {
|
||||
ch := db.ChDB(ctx)
|
||||
if ch == nil {
|
||||
return 0, fmt.Errorf("clickhouse gorm connection is not initialized")
|
||||
}
|
||||
|
||||
var count int64
|
||||
query := applyFilter(ch.Model(&analyticsmodel.UserAccessLog{}), filter)
|
||||
if err := query.Count(&count).Error; err != nil {
|
||||
return 0, fmt.Errorf("count access logs: %w", err)
|
||||
}
|
||||
return safeUint64Count(count), nil
|
||||
}
|
||||
|
||||
// ListAccessLogs returns paginated access logs and the total match count.
|
||||
func ListAccessLogs(ctx context.Context, filter AccessLogFilter, page, pageSize int) ([]analyticsmodel.UserAccessLog, uint64, error) {
|
||||
ch := db.ChDB(ctx)
|
||||
if ch == nil {
|
||||
return nil, 0, fmt.Errorf("clickhouse gorm connection is not initialized")
|
||||
}
|
||||
|
||||
if filter.UserIDs != nil && len(filter.UserIDs) == 0 {
|
||||
return []analyticsmodel.UserAccessLog{}, 0, nil
|
||||
}
|
||||
|
||||
var total int64
|
||||
baseQuery := applyFilter(ch.Model(&analyticsmodel.UserAccessLog{}), filter)
|
||||
if err := baseQuery.Count(&total).Error; err != nil {
|
||||
return nil, 0, fmt.Errorf("count access logs: %w", err)
|
||||
}
|
||||
if total == 0 {
|
||||
return []analyticsmodel.UserAccessLog{}, 0, nil
|
||||
}
|
||||
|
||||
if page < 1 {
|
||||
page = 1
|
||||
}
|
||||
if pageSize < 1 {
|
||||
pageSize = 20
|
||||
}
|
||||
offset := (page - 1) * pageSize
|
||||
|
||||
var logs []analyticsmodel.UserAccessLog
|
||||
err := applyFilter(ch.Model(&analyticsmodel.UserAccessLog{}), filter).
|
||||
Order("created_at DESC, id DESC").
|
||||
Limit(pageSize).
|
||||
Offset(offset).
|
||||
Find(&logs).Error
|
||||
if err != nil {
|
||||
return nil, 0, fmt.Errorf("list access logs: %w", err)
|
||||
}
|
||||
|
||||
return logs, safeUint64Count(total), nil
|
||||
}
|
||||
|
||||
func safeUint64Count(count int64) uint64 {
|
||||
if count < 0 {
|
||||
return 0
|
||||
}
|
||||
return uint64(count)
|
||||
}
|
||||
|
||||
func applyFilter(query *gorm.DB, filter AccessLogFilter) *gorm.DB {
|
||||
if filter.UserIDs != nil {
|
||||
if len(filter.UserIDs) == 0 {
|
||||
return query.Where("1 = 0")
|
||||
}
|
||||
query = query.Where("user_id IN ?", filter.UserIDs)
|
||||
}
|
||||
if filter.Path != "" {
|
||||
query = query.Where("path LIKE ?", "%"+filter.Path+"%")
|
||||
}
|
||||
if filter.StartTime != nil {
|
||||
query = query.Where("created_at >= ?", *filter.StartTime)
|
||||
}
|
||||
if filter.EndTime != nil {
|
||||
query = query.Where("created_at <= ?", *filter.EndTime)
|
||||
}
|
||||
return query
|
||||
}
|
||||
@@ -0,0 +1,17 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package analytics
|
||||
|
||||
import "time"
|
||||
|
||||
// AccessLogFilter scopes ClickHouse user access log queries.
|
||||
type AccessLogFilter struct {
|
||||
// UserIDs filters by user IDs. nil means no user filter; an empty slice means no matches.
|
||||
UserIDs []uint64
|
||||
Path string
|
||||
// StartTime filters created_at >= StartTime when non-nil.
|
||||
StartTime *time.Time
|
||||
// EndTime filters created_at <= EndTime when non-nil.
|
||||
EndTime *time.Time
|
||||
}
|
||||
@@ -0,0 +1,159 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package analytics
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
"sort"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics"
|
||||
)
|
||||
|
||||
const hoursInDay = 24
|
||||
|
||||
// DailyTrend is a single day's access count.
|
||||
type DailyTrend struct {
|
||||
Date string
|
||||
Count uint64
|
||||
}
|
||||
|
||||
// BrowserShare is a browser group's share of access logs.
|
||||
type BrowserShare struct {
|
||||
Browser string
|
||||
Count uint64
|
||||
}
|
||||
|
||||
// TopUser is an active user ranked by access count.
|
||||
type TopUser struct {
|
||||
UserID uint64
|
||||
Count uint64
|
||||
}
|
||||
|
||||
// GetDailyTrend returns per-day access counts for the last days days (inclusive of today).
|
||||
func GetDailyTrend(ctx context.Context, days int) ([]DailyTrend, error) {
|
||||
if days < 1 {
|
||||
days = 7
|
||||
}
|
||||
|
||||
ch := db.ChDB(ctx)
|
||||
if ch == nil {
|
||||
return nil, fmt.Errorf("clickhouse gorm connection is not initialized")
|
||||
}
|
||||
|
||||
startTime := time.Now().AddDate(0, 0, -(days - 1)).Truncate(hoursInDay * time.Hour)
|
||||
tableName := analyticsmodel.UserAccessLog{}.TableName()
|
||||
|
||||
query := fmt.Sprintf(`
|
||||
SELECT toDate(created_at) AS date, count() AS count
|
||||
FROM %s
|
||||
WHERE created_at >= ?
|
||||
GROUP BY date
|
||||
ORDER BY date ASC
|
||||
`, tableName)
|
||||
|
||||
type trendRow struct {
|
||||
Date time.Time
|
||||
Count uint64
|
||||
}
|
||||
|
||||
var rows []trendRow
|
||||
if err := ch.Raw(query, startTime).Scan(&rows).Error; err != nil {
|
||||
return nil, fmt.Errorf("get daily trend: %w", err)
|
||||
}
|
||||
|
||||
trendMap := make(map[string]uint64, days)
|
||||
for i := 0; i < days; i++ {
|
||||
dateStr := time.Now().AddDate(0, 0, -i).Format("2006-01-02")
|
||||
trendMap[dateStr] = 0
|
||||
}
|
||||
for _, row := range rows {
|
||||
dateStr := row.Date.Format("2006-01-02")
|
||||
trendMap[dateStr] = row.Count
|
||||
}
|
||||
|
||||
result := make([]DailyTrend, 0, days)
|
||||
for i := days - 1; i >= 0; i-- {
|
||||
dateStr := time.Now().AddDate(0, 0, -i).Format("2006-01-02")
|
||||
result = append(result, DailyTrend{
|
||||
Date: dateStr,
|
||||
Count: trendMap[dateStr],
|
||||
})
|
||||
}
|
||||
return result, nil
|
||||
}
|
||||
|
||||
// GetBrowserDistribution returns browser-grouped access counts since startTime.
|
||||
func GetBrowserDistribution(ctx context.Context, startTime time.Time) ([]BrowserShare, error) {
|
||||
ch := db.ChDB(ctx)
|
||||
if ch == nil {
|
||||
return nil, fmt.Errorf("clickhouse gorm connection is not initialized")
|
||||
}
|
||||
|
||||
tableName := analyticsmodel.UserAccessLog{}.TableName()
|
||||
query := fmt.Sprintf(`
|
||||
SELECT user_agent, count() AS count
|
||||
FROM %s
|
||||
WHERE created_at >= ?
|
||||
GROUP BY user_agent
|
||||
`, tableName)
|
||||
|
||||
type uaRow struct {
|
||||
UserAgent string
|
||||
Count uint64
|
||||
}
|
||||
|
||||
var rows []uaRow
|
||||
if err := ch.Raw(query, startTime).Scan(&rows).Error; err != nil {
|
||||
return nil, fmt.Errorf("get browser distribution: %w", err)
|
||||
}
|
||||
|
||||
browserCounts := make(map[string]uint64)
|
||||
for _, row := range rows {
|
||||
browser := ParseBrowserName(row.UserAgent)
|
||||
browserCounts[browser] += row.Count
|
||||
}
|
||||
|
||||
result := make([]BrowserShare, 0, len(browserCounts))
|
||||
for browser, count := range browserCounts {
|
||||
result = append(result, BrowserShare{
|
||||
Browser: browser,
|
||||
Count: count,
|
||||
})
|
||||
}
|
||||
sort.Slice(result, func(i, j int) bool {
|
||||
return result[i].Count > result[j].Count
|
||||
})
|
||||
return result, nil
|
||||
}
|
||||
|
||||
// GetTopActiveUsers returns the most active users since startTime.
|
||||
func GetTopActiveUsers(ctx context.Context, startTime time.Time, limit int) ([]TopUser, error) {
|
||||
if limit < 1 {
|
||||
limit = 10
|
||||
}
|
||||
|
||||
ch := db.ChDB(ctx)
|
||||
if ch == nil {
|
||||
return nil, fmt.Errorf("clickhouse gorm connection is not initialized")
|
||||
}
|
||||
|
||||
tableName := analyticsmodel.UserAccessLog{}.TableName()
|
||||
query := fmt.Sprintf(`
|
||||
SELECT user_id, count() AS count
|
||||
FROM %s
|
||||
WHERE created_at >= ? AND user_id > 0
|
||||
GROUP BY user_id
|
||||
ORDER BY count DESC
|
||||
LIMIT ?
|
||||
`, tableName)
|
||||
|
||||
var users []TopUser
|
||||
if err := ch.Raw(query, startTime, limit).Scan(&users).Error; err != nil {
|
||||
return nil, fmt.Errorf("get top active users: %w", err)
|
||||
}
|
||||
return users, nil
|
||||
}
|
||||
@@ -0,0 +1,207 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package analytics
|
||||
|
||||
import (
|
||||
"context"
|
||||
"testing"
|
||||
"time"
|
||||
|
||||
"github.com/ClickHouse/clickhouse-go/v2/lib/column"
|
||||
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics"
|
||||
"github.com/glebarez/sqlite"
|
||||
"github.com/stretchr/testify/assert"
|
||||
"github.com/stretchr/testify/require"
|
||||
"gorm.io/gorm"
|
||||
)
|
||||
|
||||
func setupChGormDB(t *testing.T) *gorm.DB {
|
||||
t.Helper()
|
||||
|
||||
gormDB, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{
|
||||
DisableForeignKeyConstraintWhenMigrating: true,
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.NoError(t, gormDB.AutoMigrate(&analyticsmodel.UserAccessLog{}))
|
||||
db.SetChDBForTest(gormDB)
|
||||
return gormDB
|
||||
}
|
||||
|
||||
func TestParseBrowserName(t *testing.T) {
|
||||
tests := []struct {
|
||||
name string
|
||||
ua string
|
||||
want string
|
||||
}{
|
||||
{name: "chrome", ua: "Mozilla/5.0 Chrome/120.0.0.0", want: "Chrome"},
|
||||
{name: "firefox", ua: "Mozilla/5.0 Firefox/121.0", want: "Firefox"},
|
||||
{name: "safari", ua: "Mozilla/5.0 Safari/605.1.15", want: "Safari"},
|
||||
{name: "edge", ua: "Mozilla/5.0 Edg/120.0.0.0", want: "Edge"},
|
||||
{name: "wechat", ua: "MicroMessenger/8.0", want: "WeChat"},
|
||||
{name: "postman", ua: "PostmanRuntime/7.36.0", want: "Postman"},
|
||||
{name: "other", ua: "curl/8.0", want: "Other"},
|
||||
}
|
||||
|
||||
for _, tt := range tests {
|
||||
t.Run(tt.name, func(t *testing.T) {
|
||||
assert.Equal(t, tt.want, ParseBrowserName(tt.ua))
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
func TestCountAccessLogs_EmptyUserIDs(t *testing.T) {
|
||||
setupChGormDB(t)
|
||||
t.Cleanup(func() { db.SetChDBForTest(nil) })
|
||||
|
||||
count, err := CountAccessLogs(context.Background(), AccessLogFilter{UserIDs: []uint64{}})
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, uint64(0), count)
|
||||
}
|
||||
|
||||
func TestListAccessLogs_EmptyUserIDs(t *testing.T) {
|
||||
setupChGormDB(t)
|
||||
t.Cleanup(func() { db.SetChDBForTest(nil) })
|
||||
|
||||
logs, total, err := ListAccessLogs(context.Background(), AccessLogFilter{UserIDs: []uint64{}}, 1, 20)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, uint64(0), total)
|
||||
assert.Empty(t, logs)
|
||||
}
|
||||
|
||||
func TestListAccessLogs_WithFilters(t *testing.T) {
|
||||
gormDB := setupChGormDB(t)
|
||||
t.Cleanup(func() { db.SetChDBForTest(nil) })
|
||||
|
||||
now := time.Now().UTC().Truncate(time.Second)
|
||||
logs := []analyticsmodel.UserAccessLog{
|
||||
{ID: 1, UserID: 10, Path: "/api/v1/users", Method: "GET", Status: 200, CreatedAt: now},
|
||||
{ID: 2, UserID: 20, Path: "/api/v1/admin/logs", Method: "GET", Status: 200, CreatedAt: now},
|
||||
{ID: 3, UserID: 10, Path: "/api/v1/other", Method: "POST", Status: 201, CreatedAt: now},
|
||||
}
|
||||
require.NoError(t, gormDB.Create(&logs).Error)
|
||||
|
||||
start := now.Add(-time.Hour)
|
||||
filter := AccessLogFilter{
|
||||
UserIDs: []uint64{10},
|
||||
Path: "users",
|
||||
StartTime: &start,
|
||||
}
|
||||
|
||||
count, err := CountAccessLogs(context.Background(), filter)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, uint64(1), count)
|
||||
|
||||
result, total, err := ListAccessLogs(context.Background(), filter, 1, 10)
|
||||
require.NoError(t, err)
|
||||
assert.Equal(t, uint64(1), total)
|
||||
require.Len(t, result, 1)
|
||||
assert.Equal(t, uint64(1), result[0].ID)
|
||||
assert.Equal(t, "/api/v1/users", result[0].Path)
|
||||
}
|
||||
|
||||
func TestBatchInsert_Empty(t *testing.T) {
|
||||
err := BatchInsert(context.Background(), nil)
|
||||
require.NoError(t, err)
|
||||
}
|
||||
|
||||
func TestBatchInsert_UsesModelBatchSQL(t *testing.T) {
|
||||
ctx := context.Background()
|
||||
mockBatch := &mockBatch{}
|
||||
mockConn := &mockConn{
|
||||
batch: mockBatch,
|
||||
batchQuery: analyticsmodel.UserAccessLog{}.BatchInsertSQL(),
|
||||
}
|
||||
db.SetChConnForTest(mockConn)
|
||||
t.Cleanup(func() { db.SetChConnForTest(nil) })
|
||||
|
||||
createdAt := time.Now().UTC()
|
||||
err := BatchInsert(ctx, []analyticsmodel.UserAccessLog{
|
||||
{
|
||||
ID: 1,
|
||||
UserID: 42,
|
||||
Path: "/api/v1/test",
|
||||
Method: "GET",
|
||||
IP: "127.0.0.1",
|
||||
UserAgent: "test-agent",
|
||||
Headers: "{}",
|
||||
Status: 200,
|
||||
Latency: 12,
|
||||
CreatedAt: createdAt,
|
||||
},
|
||||
})
|
||||
require.NoError(t, err)
|
||||
assert.True(t, mockConn.prepareCalled)
|
||||
assert.Equal(t, analyticsmodel.UserAccessLog{}.BatchInsertSQL(), mockConn.preparedQuery)
|
||||
assert.True(t, mockBatch.sendCalled)
|
||||
require.Len(t, mockBatch.rows, 1)
|
||||
assert.Equal(t, uint64(42), mockBatch.rows[0][1])
|
||||
}
|
||||
|
||||
type mockConn struct {
|
||||
batch driver.Batch
|
||||
batchQuery string
|
||||
prepareCalled bool
|
||||
preparedQuery string
|
||||
}
|
||||
|
||||
func (m *mockConn) Contributors() []string { return nil }
|
||||
|
||||
func (m *mockConn) ServerVersion() (*driver.ServerVersion, error) { return nil, 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) QueryRow(_ context.Context, _ string, _ ...any) driver.Row { return nil }
|
||||
|
||||
func (m *mockConn) PrepareBatch(_ context.Context, query string, _ ...driver.PrepareBatchOption) (driver.Batch, error) {
|
||||
m.prepareCalled = true
|
||||
m.preparedQuery = query
|
||||
return m.batch, nil
|
||||
}
|
||||
|
||||
func (m *mockConn) Exec(_ context.Context, _ string, _ ...any) error { return nil }
|
||||
|
||||
func (m *mockConn) AsyncInsert(_ context.Context, _ string, _ bool, _ ...any) error { return nil }
|
||||
|
||||
func (m *mockConn) Ping(_ context.Context) error { return nil }
|
||||
|
||||
func (m *mockConn) Stats() driver.Stats { return driver.Stats{} }
|
||||
|
||||
func (m *mockConn) Close() error { return nil }
|
||||
|
||||
type mockBatch struct {
|
||||
rows [][]any
|
||||
sendCalled bool
|
||||
}
|
||||
|
||||
func (m *mockBatch) Abort() error { return nil }
|
||||
|
||||
func (m *mockBatch) Append(v ...any) error {
|
||||
m.rows = append(m.rows, v)
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *mockBatch) AppendStruct(_ any) error { return nil }
|
||||
|
||||
func (m *mockBatch) Column(_ int) driver.BatchColumn { return nil }
|
||||
|
||||
func (m *mockBatch) Flush() error { return nil }
|
||||
|
||||
func (m *mockBatch) Send() error {
|
||||
m.sendCalled = true
|
||||
return nil
|
||||
}
|
||||
|
||||
func (m *mockBatch) IsSent() bool { return m.sendCalled }
|
||||
|
||||
func (m *mockBatch) Rows() int { return len(m.rows) }
|
||||
|
||||
func (m *mockBatch) Columns() []column.Interface { return nil }
|
||||
|
||||
func (m *mockBatch) Close() error { return nil }
|
||||
@@ -0,0 +1,49 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package analytics
|
||||
|
||||
import (
|
||||
"context"
|
||||
"fmt"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/db"
|
||||
analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics"
|
||||
)
|
||||
|
||||
// BatchInsert writes access logs to ClickHouse using the native batch API.
|
||||
func BatchInsert(ctx context.Context, logs []analyticsmodel.UserAccessLog) error {
|
||||
if len(logs) == 0 {
|
||||
return nil
|
||||
}
|
||||
if db.ChConn == nil {
|
||||
return fmt.Errorf("clickhouse connection is not initialized")
|
||||
}
|
||||
|
||||
batch, err := db.ChConn.PrepareBatch(ctx, analyticsmodel.UserAccessLog{}.BatchInsertSQL())
|
||||
if err != nil {
|
||||
return fmt.Errorf("prepare clickhouse batch: %w", err)
|
||||
}
|
||||
|
||||
for _, logItem := range logs {
|
||||
if err := batch.Append(
|
||||
logItem.ID,
|
||||
logItem.UserID,
|
||||
logItem.Path,
|
||||
logItem.Method,
|
||||
logItem.IP,
|
||||
logItem.UserAgent,
|
||||
logItem.Headers,
|
||||
logItem.Status,
|
||||
logItem.Latency,
|
||||
logItem.CreatedAt,
|
||||
); err != nil {
|
||||
return fmt.Errorf("append access log to batch: %w", err)
|
||||
}
|
||||
}
|
||||
|
||||
if err := batch.Send(); err != nil {
|
||||
return fmt.Errorf("send clickhouse batch: %w", err)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
@@ -0,0 +1,30 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package analytics
|
||||
|
||||
import "strings"
|
||||
|
||||
// ParseBrowserName performs lightweight User-Agent browser identification.
|
||||
func ParseBrowserName(ua string) string {
|
||||
uaLower := strings.ToLower(ua)
|
||||
if strings.Contains(uaLower, "micromessenger") {
|
||||
return "WeChat"
|
||||
}
|
||||
if strings.Contains(uaLower, "postman") {
|
||||
return "Postman"
|
||||
}
|
||||
if strings.Contains(uaLower, "edg/") || strings.Contains(uaLower, "edge") {
|
||||
return "Edge"
|
||||
}
|
||||
if strings.Contains(uaLower, "firefox") {
|
||||
return "Firefox"
|
||||
}
|
||||
if strings.Contains(uaLower, "chrome") {
|
||||
return "Chrome"
|
||||
}
|
||||
if strings.Contains(uaLower, "safari") {
|
||||
return "Safari"
|
||||
}
|
||||
return "Other"
|
||||
}
|
||||
@@ -1,38 +0,0 @@
|
||||
CREATE DATABASE IF NOT EXISTS openflare;
|
||||
|
||||
USE openflare;
|
||||
|
||||
CREATE TABLE IF NOT EXISTS of_node_access_logs
|
||||
(
|
||||
id UInt64,
|
||||
node_id String,
|
||||
logged_at DateTime64(3, 'UTC'),
|
||||
remote_addr String,
|
||||
region String,
|
||||
host String,
|
||||
path String,
|
||||
status_code Int32,
|
||||
created_at DateTime64(3, 'UTC')
|
||||
)
|
||||
ENGINE = MergeTree()
|
||||
PARTITION BY toYYYYMM(logged_at)
|
||||
ORDER BY (node_id, logged_at, remote_addr, host, path, status_code)
|
||||
SETTINGS index_granularity = 8192;
|
||||
|
||||
CREATE TABLE IF NOT EXISTS w_user_access_logs
|
||||
(
|
||||
id UInt64,
|
||||
user_id UInt64,
|
||||
path String,
|
||||
method String,
|
||||
ip String,
|
||||
user_agent String,
|
||||
headers String,
|
||||
status Int32,
|
||||
latency Int64,
|
||||
created_at DateTime
|
||||
)
|
||||
ENGINE = MergeTree()
|
||||
PARTITION BY toYYYYMM(created_at)
|
||||
ORDER BY (created_at, ip, user_id)
|
||||
SETTINGS index_granularity = 8192;
|
||||
Reference in New Issue
Block a user