perf(clickhouse): Phase 2 legacy governance — TTL cleanup, unified pool, MV, ops API

- Replace retention ALTER DELETE with MATERIALIZE TTL; use TRUNCATE for delete-all
- Remove GORM ClickHouse pool; migrate user access log reads to ChConn
- Drop query-side trim(remote_addr); enable wait_for_async_insert=1
- Add of_node_traffic_hourly MV and dashboard traffic trend fallback
- Add GET /admin/status/clickhouse operational metrics endpoint
This commit is contained in:
ryan
2026-07-02 15:47:38 +08:00
parent 38946d1af5
commit 58624db397
30 changed files with 1351 additions and 436 deletions
+1
View File
@@ -20,6 +20,7 @@ sidebar: false
### 修改
- ClickHouse 遗留治理 Phase 2:保留期清理改为 TTL `MATERIALIZE TTL`(全量清理使用 `TRUNCATE`),消除定时 `ALTER DELETE` mutation;移除 GORM 双连接池并统一 `ChConn` 读路径;查询侧去除 `trim(remote_addr)`;`wait_for_async_insert` 调整为 1;新增 `/admin/status/clickhouse` 运维指标与 `of_node_traffic_hourly` 预聚合 MV。
- ClickHouse 写入路径优化:移除 Agent 心跳路径中的同步 `ALTER DELETE` 保留清理;`batchwriter` 新增 `MinBatchSize` 抑制过小批次定时 flush;可观测 writer 批次提升至 500、flush 间隔 5s,并为 OpenResty/FRPS/FRPC 补全去重。
- ClickHouse 客户端启用 `async_insert` 异步写入缓冲,并调高 `block_buffer_size` 与连接池默认值,降低小 part 与连接争用。
- Dashboard 与节点可观测 API 消除无 `LIMIT` 全表扫描、增加短 TTL 内存缓存,前端轮询间隔分别调整为 60s/30s。
+219 -5
View File
@@ -309,6 +309,7 @@ const docTemplate = `{
"parameters": [
{
"type": "integer",
"format": "int64",
"description": "认证源 ID 或名称",
"name": "id",
"in": "path",
@@ -386,6 +387,7 @@ const docTemplate = `{
"parameters": [
{
"type": "integer",
"format": "int64",
"description": "认证源 ID 或名称",
"name": "id",
"in": "path",
@@ -453,6 +455,7 @@ const docTemplate = `{
"parameters": [
{
"type": "integer",
"format": "int64",
"description": "认证源 ID 或名称",
"name": "id",
"in": "path",
@@ -1363,6 +1366,7 @@ const docTemplate = `{
"parameters": [
{
"type": "integer",
"format": "int64",
"description": "通道ID",
"name": "id",
"in": "path",
@@ -1416,6 +1420,7 @@ const docTemplate = `{
"parameters": [
{
"type": "integer",
"format": "int64",
"description": "通道ID",
"name": "id",
"in": "path",
@@ -1872,6 +1877,67 @@ const docTemplate = `{
}
}
},
"/api/v1/admin/status/clickhouse": {
"get": {
"security": [
{
"SessionCookie": []
}
],
"description": "返回 ClickHouse parts、mutation、async_insert 队列等运维指标,需要管理员权限",
"produces": [
"application/json"
],
"tags": [
"admin"
],
"summary": "获取 ClickHouse 运行指标",
"responses": {
"200": {
"description": "获取成功",
"schema": {
"allOf": [
{
"$ref": "#/definitions/response.Any"
},
{
"type": "object",
"properties": {
"data": {
"$ref": "#/definitions/analytics.ClickHouseOperationalStats"
}
}
}
]
}
},
"400": {
"description": "ClickHouse 未启用",
"schema": {
"$ref": "#/definitions/response.Any"
}
},
"401": {
"description": "未登录",
"schema": {
"$ref": "#/definitions/response.Any"
}
},
"403": {
"description": "无管理员权限",
"schema": {
"$ref": "#/definitions/response.Any"
}
},
"500": {
"description": "内部错误",
"schema": {
"$ref": "#/definitions/response.Any"
}
}
}
}
},
"/api/v1/admin/system-configs": {
"get": {
"security": [
@@ -3368,6 +3434,7 @@ const docTemplate = `{
},
{
"type": "integer",
"format": "int64",
"description": "上传用户 ID",
"name": "user_id",
"in": "query"
@@ -3691,6 +3758,11 @@ const docTemplate = `{
],
"summary": "获取用户列表",
"parameters": [
{
"type": "string",
"name": "email",
"in": "query"
},
{
"minimum": 1,
"type": "integer",
@@ -3909,6 +3981,92 @@ const docTemplate = `{
}
}
},
"put": {
"security": [
{
"SessionCookie": []
}
],
"description": "更新指定用户的昵称、邮箱、管理员权限,并可选重置密码,需要管理员权限",
"consumes": [
"application/json"
],
"produces": [
"application/json"
],
"tags": [
"admin"
],
"summary": "更新用户信息",
"parameters": [
{
"type": "integer",
"description": "用户 ID",
"name": "id",
"in": "path",
"required": true
},
{
"description": "更新参数",
"name": "request",
"in": "body",
"required": true,
"schema": {
"$ref": "#/definitions/user.updateUserRequest"
}
}
],
"responses": {
"200": {
"description": "更新成功",
"schema": {
"allOf": [
{
"$ref": "#/definitions/response.Any"
},
{
"type": "object",
"properties": {
"data": {
"type": "string"
}
}
}
]
}
},
"400": {
"description": "参数错误",
"schema": {
"$ref": "#/definitions/response.Any"
}
},
"401": {
"description": "未登录",
"schema": {
"$ref": "#/definitions/response.Any"
}
},
"403": {
"description": "无管理员权限或尝试修改自身权限",
"schema": {
"$ref": "#/definitions/response.Any"
}
},
"404": {
"description": "用户不存在",
"schema": {
"$ref": "#/definitions/response.Any"
}
},
"500": {
"description": "内部错误",
"schema": {
"$ref": "#/definitions/response.Any"
}
}
}
},
"delete": {
"security": [
{
@@ -11144,6 +11302,7 @@ const docTemplate = `{
"parameters": [
{
"type": "integer",
"format": "int64",
"description": "外部帐号绑定记录 ID",
"name": "id",
"in": "path",
@@ -12948,6 +13107,29 @@ const docTemplate = `{
}
}
},
"analytics.ClickHouseOperationalStats": {
"type": "object",
"properties": {
"active_parts": {
"type": "integer"
},
"async_insert_bytes": {
"type": "integer"
},
"async_insert_queue": {
"type": "integer"
},
"database": {
"type": "string"
},
"pending_mutations": {
"type": "integer"
},
"total_rows": {
"type": "integer"
}
}
},
"apply_log.CleanupInput": {
"type": "object",
"properties": {
@@ -13775,19 +13957,22 @@ const docTemplate = `{
"source_countries": {
"type": "object",
"additionalProperties": {
"type": "integer"
"type": "integer",
"format": "int64"
}
},
"status_codes": {
"type": "object",
"additionalProperties": {
"type": "integer"
"type": "integer",
"format": "int64"
}
},
"top_domains": {
"type": "object",
"additionalProperties": {
"type": "integer"
"type": "integer",
"format": "int64"
}
},
"unique_visitor_count": {
@@ -14217,7 +14402,7 @@ const docTemplate = `{
"type": "string"
},
"id": {
"type": "integer"
"type": "string"
},
"is_active": {
"type": "boolean"
@@ -14252,7 +14437,7 @@ const docTemplate = `{
"type": "string"
},
"id": {
"type": "integer"
"type": "string"
},
"is_active": {
"type": "boolean"
@@ -15092,6 +15277,11 @@ const docTemplate = `{
"UploadStatusPending": "待使用",
"UploadStatusUsed": "已使用"
},
"x-enum-descriptions": [
"待使用",
"已使用",
"已删除"
],
"x-enum-varnames": [
"UploadStatusPending",
"UploadStatusUsed",
@@ -18085,6 +18275,30 @@ const docTemplate = `{
}
}
},
"user.updateUserRequest": {
"type": "object",
"required": [
"email"
],
"properties": {
"email": {
"type": "string",
"maxLength": 255
},
"is_admin": {
"type": "boolean"
},
"nickname": {
"type": "string",
"maxLength": 64
},
"password": {
"type": "string",
"maxLength": 64,
"minLength": 8
}
}
},
"user.updateUserStatusRequest": {
"type": "object",
"properties": {
+1 -1
View File
@@ -1,7 +1,7 @@
# ClickHouse CPU 性能优化计划
> PLAN_ID: `63ba981b`
> 状态: 已完成
> 状态: 已完成(含 Phase 2 遗留治理)
> 目标: 完成 P0–P2 优化,降低 ClickHouse CPU 占用
## 背景
+219 -5
View File
@@ -302,6 +302,7 @@
"parameters": [
{
"type": "integer",
"format": "int64",
"description": "认证源 ID 或名称",
"name": "id",
"in": "path",
@@ -379,6 +380,7 @@
"parameters": [
{
"type": "integer",
"format": "int64",
"description": "认证源 ID 或名称",
"name": "id",
"in": "path",
@@ -446,6 +448,7 @@
"parameters": [
{
"type": "integer",
"format": "int64",
"description": "认证源 ID 或名称",
"name": "id",
"in": "path",
@@ -1356,6 +1359,7 @@
"parameters": [
{
"type": "integer",
"format": "int64",
"description": "通道ID",
"name": "id",
"in": "path",
@@ -1409,6 +1413,7 @@
"parameters": [
{
"type": "integer",
"format": "int64",
"description": "通道ID",
"name": "id",
"in": "path",
@@ -1865,6 +1870,67 @@
}
}
},
"/api/v1/admin/status/clickhouse": {
"get": {
"security": [
{
"SessionCookie": []
}
],
"description": "返回 ClickHouse parts、mutation、async_insert 队列等运维指标,需要管理员权限",
"produces": [
"application/json"
],
"tags": [
"admin"
],
"summary": "获取 ClickHouse 运行指标",
"responses": {
"200": {
"description": "获取成功",
"schema": {
"allOf": [
{
"$ref": "#/definitions/response.Any"
},
{
"type": "object",
"properties": {
"data": {
"$ref": "#/definitions/analytics.ClickHouseOperationalStats"
}
}
}
]
}
},
"400": {
"description": "ClickHouse 未启用",
"schema": {
"$ref": "#/definitions/response.Any"
}
},
"401": {
"description": "未登录",
"schema": {
"$ref": "#/definitions/response.Any"
}
},
"403": {
"description": "无管理员权限",
"schema": {
"$ref": "#/definitions/response.Any"
}
},
"500": {
"description": "内部错误",
"schema": {
"$ref": "#/definitions/response.Any"
}
}
}
}
},
"/api/v1/admin/system-configs": {
"get": {
"security": [
@@ -3361,6 +3427,7 @@
},
{
"type": "integer",
"format": "int64",
"description": "上传用户 ID",
"name": "user_id",
"in": "query"
@@ -3684,6 +3751,11 @@
],
"summary": "获取用户列表",
"parameters": [
{
"type": "string",
"name": "email",
"in": "query"
},
{
"minimum": 1,
"type": "integer",
@@ -3902,6 +3974,92 @@
}
}
},
"put": {
"security": [
{
"SessionCookie": []
}
],
"description": "更新指定用户的昵称、邮箱、管理员权限,并可选重置密码,需要管理员权限",
"consumes": [
"application/json"
],
"produces": [
"application/json"
],
"tags": [
"admin"
],
"summary": "更新用户信息",
"parameters": [
{
"type": "integer",
"description": "用户 ID",
"name": "id",
"in": "path",
"required": true
},
{
"description": "更新参数",
"name": "request",
"in": "body",
"required": true,
"schema": {
"$ref": "#/definitions/user.updateUserRequest"
}
}
],
"responses": {
"200": {
"description": "更新成功",
"schema": {
"allOf": [
{
"$ref": "#/definitions/response.Any"
},
{
"type": "object",
"properties": {
"data": {
"type": "string"
}
}
}
]
}
},
"400": {
"description": "参数错误",
"schema": {
"$ref": "#/definitions/response.Any"
}
},
"401": {
"description": "未登录",
"schema": {
"$ref": "#/definitions/response.Any"
}
},
"403": {
"description": "无管理员权限或尝试修改自身权限",
"schema": {
"$ref": "#/definitions/response.Any"
}
},
"404": {
"description": "用户不存在",
"schema": {
"$ref": "#/definitions/response.Any"
}
},
"500": {
"description": "内部错误",
"schema": {
"$ref": "#/definitions/response.Any"
}
}
}
},
"delete": {
"security": [
{
@@ -11137,6 +11295,7 @@
"parameters": [
{
"type": "integer",
"format": "int64",
"description": "外部帐号绑定记录 ID",
"name": "id",
"in": "path",
@@ -12941,6 +13100,29 @@
}
}
},
"analytics.ClickHouseOperationalStats": {
"type": "object",
"properties": {
"active_parts": {
"type": "integer"
},
"async_insert_bytes": {
"type": "integer"
},
"async_insert_queue": {
"type": "integer"
},
"database": {
"type": "string"
},
"pending_mutations": {
"type": "integer"
},
"total_rows": {
"type": "integer"
}
}
},
"apply_log.CleanupInput": {
"type": "object",
"properties": {
@@ -13768,19 +13950,22 @@
"source_countries": {
"type": "object",
"additionalProperties": {
"type": "integer"
"type": "integer",
"format": "int64"
}
},
"status_codes": {
"type": "object",
"additionalProperties": {
"type": "integer"
"type": "integer",
"format": "int64"
}
},
"top_domains": {
"type": "object",
"additionalProperties": {
"type": "integer"
"type": "integer",
"format": "int64"
}
},
"unique_visitor_count": {
@@ -14210,7 +14395,7 @@
"type": "string"
},
"id": {
"type": "integer"
"type": "string"
},
"is_active": {
"type": "boolean"
@@ -14245,7 +14430,7 @@
"type": "string"
},
"id": {
"type": "integer"
"type": "string"
},
"is_active": {
"type": "boolean"
@@ -15085,6 +15270,11 @@
"UploadStatusPending": "待使用",
"UploadStatusUsed": "已使用"
},
"x-enum-descriptions": [
"待使用",
"已使用",
"已删除"
],
"x-enum-varnames": [
"UploadStatusPending",
"UploadStatusUsed",
@@ -18078,6 +18268,30 @@
}
}
},
"user.updateUserRequest": {
"type": "object",
"required": [
"email"
],
"properties": {
"email": {
"type": "string",
"maxLength": 255
},
"is_admin": {
"type": "boolean"
},
"nickname": {
"type": "string",
"maxLength": 64
},
"password": {
"type": "string",
"maxLength": 64,
"minLength": 8
}
}
},
"user.updateUserStatusRequest": {
"type": "object",
"properties": {
+140 -2
View File
@@ -169,6 +169,21 @@ definitions:
$ref: '#/definitions/github_com_Rain-kl_Wavelet_pkg_protocol.WAFIPGroup'
type: array
type: object
analytics.ClickHouseOperationalStats:
properties:
active_parts:
type: integer
async_insert_bytes:
type: integer
async_insert_queue:
type: integer
database:
type: string
pending_mutations:
type: integer
total_rows:
type: integer
type: object
apply_log.CleanupInput:
properties:
delete_all:
@@ -713,14 +728,17 @@ definitions:
type: integer
source_countries:
additionalProperties:
format: int64
type: integer
type: object
status_codes:
additionalProperties:
format: int64
type: integer
type: object
top_domains:
additionalProperties:
format: int64
type: integer
type: object
unique_visitor_count:
@@ -1004,7 +1022,7 @@ definitions:
created_by:
type: string
id:
type: integer
type: string
is_active:
type: boolean
main_config:
@@ -1027,7 +1045,7 @@ definitions:
created_by:
type: string
id:
type: integer
type: string
is_active:
type: boolean
version:
@@ -1593,6 +1611,10 @@ definitions:
UploadStatusDeleted: 已删除
UploadStatusPending: 待使用
UploadStatusUsed: 已使用
x-enum-descriptions:
- 待使用
- 已使用
- 已删除
x-enum-varnames:
- UploadStatusPending
- UploadStatusUsed
@@ -3577,6 +3599,23 @@ definitions:
website:
type: string
type: object
user.updateUserRequest:
properties:
email:
maxLength: 255
type: string
is_admin:
type: boolean
nickname:
maxLength: 64
type: string
password:
maxLength: 64
minLength: 8
type: string
required:
- email
type: object
user.updateUserStatusRequest:
properties:
is_active:
@@ -4085,6 +4124,7 @@ paths:
description: 删除指定认证源及其关联的所有外部帐号绑定记录,警告:删除后相关用户将无法通过该源登录,需要管理员权限
parameters:
- description: 认证源 ID 或名称
format: int64
in: path
name: id
required: true
@@ -4124,6 +4164,7 @@ paths:
description: 更新指定 ID 的认证源配置。若 client_secret 字段为空,则保留原有密钥不变,需要管理员权限
parameters:
- description: 认证源 ID 或名称
format: int64
in: path
name: id
required: true
@@ -4174,6 +4215,7 @@ paths:
description: 启用或禁用指定认证源。尝试启用时将验证 Client ID 和 Client Secret 是否已配置,需要管理员权限
parameters:
- description: 认证源 ID 或名称
format: int64
in: path
name: id
required: true
@@ -4670,6 +4712,7 @@ paths:
description: 根据ID删除消息通道,需要管理员权限
parameters:
- description: 通道ID
format: int64
in: path
name: id
required: true
@@ -4692,6 +4735,7 @@ paths:
description: 修改消息通道配置,需要管理员权限
parameters:
- description: 通道ID
format: int64
in: path
name: id
required: true
@@ -5016,6 +5060,42 @@ paths:
summary: 获取系统状态信息
tags:
- admin
/api/v1/admin/status/clickhouse:
get:
description: 返回 ClickHouse parts、mutation、async_insert 队列等运维指标,需要管理员权限
produces:
- application/json
responses:
"200":
description: 获取成功
schema:
allOf:
- $ref: '#/definitions/response.Any'
- properties:
data:
$ref: '#/definitions/analytics.ClickHouseOperationalStats'
type: object
"400":
description: ClickHouse 未启用
schema:
$ref: '#/definitions/response.Any'
"401":
description: 未登录
schema:
$ref: '#/definitions/response.Any'
"403":
description: 无管理员权限
schema:
$ref: '#/definitions/response.Any'
"500":
description: 内部错误
schema:
$ref: '#/definitions/response.Any'
security:
- SessionCookie: []
summary: 获取 ClickHouse 运行指标
tags:
- admin
/api/v1/admin/system-configs:
get:
description: 返回所有系统配置列表,支持按配置类型(system/business)过滤,需要管理员权限
@@ -5912,6 +5992,7 @@ paths:
name: extension
type: string
- description: 上传用户 ID
format: int64
in: query
name: user_id
type: integer
@@ -6108,6 +6189,9 @@ paths:
get:
description: 分页返回用户列表,支持按用户 ID 和用户名筛选,需要管理员权限
parameters:
- in: query
name: email
type: string
- in: query
minimum: 1
name: page
@@ -6291,6 +6375,59 @@ paths:
summary: 获取用户详情
tags:
- admin
put:
consumes:
- application/json
description: 更新指定用户的昵称、邮箱、管理员权限,并可选重置密码,需要管理员权限
parameters:
- description: 用户 ID
in: path
name: id
required: true
type: integer
- description: 更新参数
in: body
name: request
required: true
schema:
$ref: '#/definitions/user.updateUserRequest'
produces:
- application/json
responses:
"200":
description: 更新成功
schema:
allOf:
- $ref: '#/definitions/response.Any'
- properties:
data:
type: string
type: object
"400":
description: 参数错误
schema:
$ref: '#/definitions/response.Any'
"401":
description: 未登录
schema:
$ref: '#/definitions/response.Any'
"403":
description: 无管理员权限或尝试修改自身权限
schema:
$ref: '#/definitions/response.Any'
"404":
description: 用户不存在
schema:
$ref: '#/definitions/response.Any'
"500":
description: 内部错误
schema:
$ref: '#/definitions/response.Any'
security:
- SessionCookie: []
summary: 更新用户信息
tags:
- admin
/api/v1/admin/users/{id}/status:
put:
consumes:
@@ -10642,6 +10779,7 @@ paths:
description: 解除当前登录用户与指定外部帐号的绑定关系,需要登录
parameters:
- description: 外部帐号绑定记录 ID
format: int64
in: path
name: id
required: true
+2 -2
View File
@@ -248,7 +248,7 @@ func enrichAccessLogsWithUsers(ctx context.Context, list []accessLogItem) {
// @Router /api/v1/admin/logs/access [get]
func GetAccessLogs(c *gin.Context) {
ctx := c.Request.Context()
if !config.Config.ClickHouse.Enabled || db.ChDB(ctx) == nil {
if !config.Config.ClickHouse.Enabled || !db.ChConnReady() {
response.AbortWithError(c, http.StatusBadRequest, "ClickHouse 存储服务未启用,无法检索访问日志")
return
}
@@ -348,7 +348,7 @@ type logsAnalyticsResponse struct {
// @Router /api/v1/admin/logs/analytics [get]
func GetLogsAnalytics(c *gin.Context) {
ctx := c.Request.Context()
if !config.Config.ClickHouse.Enabled || db.ChDB(ctx) == nil {
if !config.Config.ClickHouse.Enabled || !db.ChConnReady() {
response.AbortWithError(c, http.StatusBadRequest, "ClickHouse 存储服务未启用,无法获取分析数据")
return
}
+40
View File
@@ -0,0 +1,40 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package status
import (
"net/http"
"github.com/Rain-kl/Wavelet/internal/common/response"
"github.com/Rain-kl/Wavelet/internal/config"
"github.com/Rain-kl/Wavelet/internal/db"
analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics"
"github.com/gin-gonic/gin"
)
// GetClickHouseStatus returns ClickHouse operational metrics for administrators.
// @Summary 获取 ClickHouse 运行指标
// @Description 返回 ClickHouse parts、mutation、async_insert 队列等运维指标,需要管理员权限
// @Tags admin
// @Produce json
// @Security SessionCookie
// @Success 200 {object} response.Any{data=analyticsrepo.ClickHouseOperationalStats} "获取成功"
// @Failure 400 {object} response.Any "ClickHouse 未启用"
// @Failure 401 {object} response.Any "未登录"
// @Failure 403 {object} response.Any "无管理员权限"
// @Failure 500 {object} response.Any "内部错误"
// @Router /api/v1/admin/status/clickhouse [get]
func GetClickHouseStatus(c *gin.Context) {
if !config.Config.ClickHouse.Enabled || !db.ChConnReady() {
response.AbortWithError(c, http.StatusBadRequest, "ClickHouse 存储服务未启用")
return
}
stats, err := analyticsrepo.GetClickHouseOperationalStats(c.Request.Context())
if err != nil {
response.AbortInternal(c, "获取 ClickHouse 运行指标失败")
return
}
c.JSON(http.StatusOK, response.OK(stats))
}
+5 -1
View File
@@ -137,13 +137,17 @@ func buildOverviewView(ctx context.Context) (*OverviewView, error) {
if err != nil {
return nil, err
}
trafficTrend := observability.BuildTrafficTrendPoints(now, reports)
if trafficHourly, hourlyErr := model.ListOpenFlareTrafficHourlySince(ctx, "", since); hourlyErr == nil && len(trafficHourly) > 0 {
trafficTrend = observability.BuildTrafficTrendPointsFromHourly(now, trafficHourly)
}
view := &OverviewView{
GeneratedAt: now,
Nodes: make([]NodeHealth, 0, len(nodes)),
Distributions: observability.BuildTrafficDistributions(reports, accessLogRegions, dashboardDistributionLimit),
Trends: observability.NodeTrends{
Traffic24h: observability.BuildTrafficTrendPoints(now, reports),
Traffic24h: trafficTrend,
Capacity24h: observability.BuildCapacityTrendPoints(now, snapshots),
Network24h: observability.BuildNetworkTrendPoints(now, snapshots, openrestySnapshots),
DiskIO24h: observability.BuildDiskIOTrendPoints(now, snapshots),
@@ -257,6 +257,28 @@ func buildHealthSummary(
return summary
}
// BuildTrafficTrendPointsFromHourly builds 24h traffic trend buckets from hourly rollups.
func BuildTrafficTrendPointsFromHourly(now time.Time, hourly []*model.OpenFlareTrafficHourly) []TrafficTrendPoint {
start := trendWindowStart(now)
points := make([]TrafficTrendPoint, observabilityTrendBuckets)
for index := range points {
points[index].BucketStartedAt = start.Add(time.Duration(index) * time.Hour)
}
for _, row := range hourly {
if row == nil {
continue
}
index, ok := trendBucketIndex(row.Hour, start)
if !ok {
continue
}
points[index].RequestCount += row.RequestCount
points[index].ErrorCount += row.ErrorCount
points[index].UniqueVisitorCount += row.UniqueVisitorCount
}
return points
}
// BuildTrafficTrendPoints builds 24h traffic trend buckets.
func BuildTrafficTrendPoints(now time.Time, reports []*model.OpenFlareRequestReport) []TrafficTrendPoint {
start := trendWindowStart(now)
@@ -10,6 +10,23 @@ import (
"github.com/Rain-kl/Wavelet/internal/model"
)
func TestBuildTrafficTrendPointsFromHourlyBucketsByHour(t *testing.T) {
now := time.Date(2026, 7, 2, 15, 30, 0, 0, time.UTC)
hourly := []*model.OpenFlareTrafficHourly{
{
NodeID: "node-a",
Hour: now.Add(-2 * time.Hour).Truncate(time.Hour),
RequestCount: 12,
ErrorCount: 1,
UniqueVisitorCount: 4,
},
}
points := BuildTrafficTrendPointsFromHourly(now, hourly)
if len(points) != observabilityTrendBuckets {
t.Fatalf("BuildTrafficTrendPointsFromHourly() len = %d, want %d", len(points), observabilityTrendBuckets)
}
}
func TestBuildTrafficTrendPointsBucketsByHour(t *testing.T) {
t.Parallel()
@@ -134,6 +134,10 @@ func GetNodeObservability(ctx context.Context, id uint, query NodeQuery) (*NodeV
if err != nil {
return nil, err
}
trafficTrend := BuildTrafficTrendPoints(now, reports)
if trafficHourly, hourlyErr := model.ListOpenFlareTrafficHourlySince(ctx, node.NodeID, now.Add(-24*time.Hour)); hourlyErr == nil && len(trafficHourly) > 0 {
trafficTrend = BuildTrafficTrendPointsFromHourly(now, trafficHourly)
}
view := &NodeView{
NodeID: node.NodeID,
@@ -147,7 +151,7 @@ func GetNodeObservability(ctx context.Context, id uint, query NodeQuery) (*NodeV
Health: buildHealthSummary(latestMetricSnapshot(snapshots), latestTrafficReport(reports), events),
},
Trends: NodeTrends{
Traffic24h: BuildTrafficTrendPoints(now, reports),
Traffic24h: trafficTrend,
Capacity24h: BuildCapacityTrendPoints(now, snapshots),
Network24h: BuildNetworkTrendPoints(now, snapshots, openrestyObs),
DiskIO24h: BuildDiskIOTrendPoints(now, snapshots),
@@ -12,6 +12,7 @@ import (
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/internal/repository"
analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics"
)
const (
@@ -49,6 +50,7 @@ type DatabaseCleanupResult struct {
Target string `json:"target"`
TargetLabel string `json:"target_label"`
DeletedCount int64 `json:"deleted_count"`
CleanupMode string `json:"cleanup_mode,omitempty"`
DeleteAll bool `json:"delete_all"`
RetentionDays *int `json:"retention_days,omitempty"`
Cutoff *time.Time `json:"cutoff,omitempty"`
@@ -79,21 +81,23 @@ func CleanupDatabaseObservability(ctx context.Context, input DatabaseCleanupInpu
}
if input.RetentionDays == nil {
deleted, err := deleteAllObservabilityRows(ctx, target)
deleted, mode, err := deleteAllObservabilityRows(ctx, target)
if err != nil {
return nil, err
}
result.DeletedCount = deleted
result.CleanupMode = mode
return result, nil
}
retentionDays := *input.RetentionDays
cutoff := time.Now().UTC().Add(-time.Duration(retentionDays) * 24 * time.Hour)
deleted, err := deleteObservabilityRowsBefore(ctx, target, cutoff)
deleted, mode, err := deleteObservabilityRowsBefore(ctx, target, cutoff)
if err != nil {
return nil, err
}
result.DeletedCount = deleted
result.CleanupMode = mode
result.RetentionDays = &retentionDays
result.Cutoff = &cutoff
return result, nil
@@ -141,40 +145,56 @@ func RunDatabaseAutoCleanupOnce(ctx context.Context, now time.Time) (*DatabaseAu
}, nil
}
func deleteAllObservabilityRows(ctx context.Context, target string) (int64, error) {
func deleteAllObservabilityRows(ctx context.Context, target string) (int64, string, error) {
var (
deleted int64
err error
)
switch target {
case DatabaseCleanupTargetAccessLogs:
return model.DeleteAllOpenFlareAccessLogs(ctx)
deleted, err = model.DeleteAllOpenFlareAccessLogs(ctx)
case DatabaseCleanupTargetMetricSnapshots:
return model.DeleteAllOpenFlareMetricSnapshots(ctx)
deleted, err = model.DeleteAllOpenFlareMetricSnapshots(ctx)
case DatabaseCleanupTargetRequestReports:
return model.DeleteAllOpenFlareRequestReports(ctx)
deleted, err = model.DeleteAllOpenFlareRequestReports(ctx)
case DatabaseCleanupTargetObsOpenresty:
return model.DeleteAllOpenFlareNodeObservationOpenresty(ctx)
deleted, err = model.DeleteAllOpenFlareNodeObservationOpenresty(ctx)
case DatabaseCleanupTargetObsFrps:
return model.DeleteAllOpenFlareNodeObservationFrps(ctx)
deleted, err = model.DeleteAllOpenFlareNodeObservationFrps(ctx)
case DatabaseCleanupTargetObsFrpc:
return model.DeleteAllOpenFlareNodeObservationFrpc(ctx)
deleted, err = model.DeleteAllOpenFlareNodeObservationFrpc(ctx)
default:
return 0, errors.New("unsupported cleanup target")
return 0, "", errors.New("unsupported cleanup target")
}
if err != nil {
return 0, "", err
}
return deleted, analyticsrepo.CleanupModeTruncate, nil
}
func deleteObservabilityRowsBefore(ctx context.Context, target string, cutoff time.Time) (int64, error) {
func deleteObservabilityRowsBefore(ctx context.Context, target string, cutoff time.Time) (int64, string, error) {
var (
deleted int64
err error
)
switch target {
case DatabaseCleanupTargetAccessLogs:
return model.DeleteOpenFlareAccessLogsBefore(ctx, cutoff)
deleted, err = model.DeleteOpenFlareAccessLogsBefore(ctx, cutoff)
case DatabaseCleanupTargetMetricSnapshots:
return model.DeleteOpenFlareMetricSnapshotsBefore(ctx, cutoff)
deleted, err = model.DeleteOpenFlareMetricSnapshotsBefore(ctx, cutoff)
case DatabaseCleanupTargetRequestReports:
return model.DeleteOpenFlareRequestReportsBefore(ctx, cutoff)
deleted, err = model.DeleteOpenFlareRequestReportsBefore(ctx, cutoff)
case DatabaseCleanupTargetObsOpenresty:
return model.DeleteOpenFlareNodeObservationOpenrestyBefore(ctx, cutoff)
deleted, err = model.DeleteOpenFlareNodeObservationOpenrestyBefore(ctx, cutoff)
case DatabaseCleanupTargetObsFrps:
return model.DeleteOpenFlareNodeObservationFrpsBefore(ctx, cutoff)
deleted, err = model.DeleteOpenFlareNodeObservationFrpsBefore(ctx, cutoff)
case DatabaseCleanupTargetObsFrpc:
return model.DeleteOpenFlareNodeObservationFrpcBefore(ctx, cutoff)
deleted, err = model.DeleteOpenFlareNodeObservationFrpcBefore(ctx, cutoff)
default:
return 0, errors.New("unsupported cleanup target")
return 0, "", errors.New("unsupported cleanup target")
}
if err != nil {
return 0, "", err
}
return deleted, analyticsrepo.CleanupModeTTLMaterialize, nil
}
+5 -80
View File
@@ -7,20 +7,12 @@ 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 (
@@ -31,10 +23,8 @@ const (
)
var (
// ChConn ClickHouse 原生连接实例,用于批量写入
// ChConn ClickHouse 原生连接实例,用于批量写入与查询
ChConn driver.Conn
chDB *gorm.DB
)
func init() {
@@ -59,36 +49,6 @@ func init() {
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")
}
@@ -105,7 +65,7 @@ func buildClickHouseOptions() *clickhouse.Options {
Settings: clickhouse.Settings{
"max_execution_time": clickhouseMaxExecTime,
"async_insert": 1,
"wait_for_async_insert": 0,
"wait_for_async_insert": 1,
"async_insert_max_data_size": clickhouseAsyncInsertMaxDataSize,
"async_insert_busy_timeout_ms": clickhouseAsyncInsertBusyTimeoutMs,
},
@@ -121,44 +81,9 @@ func buildClickHouseOptions() *clickhouse.Options {
}
}
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))
// GORM uses clickhouse-go ParseDSN; unknown params are passed as server Settings.
query.Set("async_insert", "1")
query.Set("wait_for_async_insert", "0")
query.Set("async_insert_max_data_size", strconv.Itoa(clickhouseAsyncInsertMaxDataSize))
query.Set("async_insert_busy_timeout_ms", strconv.Itoa(clickhouseAsyncInsertBusyTimeoutMs))
query.Set("max_insert_block_size", strconv.Itoa(clickhouseAsyncInsertMaxDataSize))
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
// ChConnReady reports whether the native ClickHouse connection is initialized.
func ChConnReady() bool {
return ChConn != nil
}
// SetChConnForTest sets the package-level native ClickHouse connection for testing.
@@ -0,0 +1,28 @@
-- +goose Up
CREATE TABLE IF NOT EXISTS of_node_traffic_hourly
(
node_id String,
hour DateTime,
request_count UInt64,
error_count UInt64,
unique_visitor_count UInt64
)
ENGINE = SummingMergeTree()
PARTITION BY toYYYYMM(hour)
ORDER BY (node_id, hour);
CREATE MATERIALIZED VIEW IF NOT EXISTS of_node_traffic_hourly_mv
TO of_node_traffic_hourly
AS
SELECT
node_id,
toStartOfHour(window_ended_at) AS hour,
sum(request_count) AS request_count,
sum(error_count) AS error_count,
sum(unique_visitor_count) AS unique_visitor_count
FROM of_node_request_reports
GROUP BY node_id, hour;
-- +goose Down
DROP VIEW IF EXISTS of_node_traffic_hourly_mv;
DROP TABLE IF EXISTS of_node_traffic_hourly;
+32
View File
@@ -10,6 +10,7 @@ import (
"time"
"github.com/Rain-kl/Wavelet/internal/db"
analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics"
"gorm.io/gorm"
)
@@ -337,6 +338,37 @@ func ListOpenFlareRequestReportsSince(ctx context.Context, nodeID string, since
return currentObservabilityStore().ListRequestReports(ctx, nodeID, since, limit)
}
// OpenFlareTrafficHourly is an hourly traffic rollup row.
type OpenFlareTrafficHourly struct {
NodeID string `json:"node_id"`
Hour time.Time `json:"hour"`
RequestCount int64 `json:"request_count"`
ErrorCount int64 `json:"error_count"`
UniqueVisitorCount int64 `json:"unique_visitor_count"`
}
// ListOpenFlareTrafficHourlySince returns hourly traffic rollup rows since the given time.
func ListOpenFlareTrafficHourlySince(ctx context.Context, nodeID string, since time.Time) ([]*OpenFlareTrafficHourly, error) {
rows, err := analyticsrepo.ListNodeTrafficHourly(ctx, analyticsrepo.NodeObservabilityFilter{
NodeID: nodeID,
Since: since,
})
if err != nil {
return nil, err
}
result := make([]*OpenFlareTrafficHourly, len(rows))
for index, row := range rows {
result[index] = &OpenFlareTrafficHourly{
NodeID: row.NodeID,
Hour: row.Hour,
RequestCount: row.RequestCount,
ErrorCount: row.ErrorCount,
UniqueVisitorCount: row.UniqueVisitorCount,
}
}
return result, nil
}
// ListOpenFlareActiveHealthEvents returns active health events across all nodes.
func ListOpenFlareActiveHealthEvents(ctx context.Context) ([]*OpenFlareHealthEvent, error) {
conn := db.DB(ctx)
+59 -42
View File
@@ -7,41 +7,51 @@ package analytics
import (
"context"
"fmt"
"time"
"github.com/Rain-kl/Wavelet/internal/db"
analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics"
"gorm.io/gorm"
)
func userAccessLogConn() error {
if db.ChConn == nil {
return fmt.Errorf("clickhouse native connection is not initialized")
}
return nil
}
// 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")
clause, args, ok := buildUserAccessLogFilterClause(filter)
if !ok {
return 0, nil
}
var count int64
query := applyFilter(ch.Model(&analyticsmodel.UserAccessLog{}), filter)
if err := query.Count(&count).Error; err != nil {
if err := userAccessLogConn(); err != nil {
return 0, err
}
tableName := analyticsmodel.UserAccessLog{}.TableName()
sql := fmt.Sprintf("SELECT count() FROM %s WHERE %s", tableName, clause)
var count uint64
if err := db.ChConn.QueryRow(ctx, sql, args...).Scan(&count); err != nil {
return 0, fmt.Errorf("count access logs: %w", err)
}
return safeUint64Count(count), nil
return 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 {
clause, args, ok := buildUserAccessLogFilterClause(filter)
if !ok {
return []analyticsmodel.UserAccessLog{}, 0, nil
}
if err := userAccessLogConn(); err != nil {
return nil, 0, err
}
var total int64
baseQuery := applyFilter(ch.Model(&analyticsmodel.UserAccessLog{}), filter)
if err := baseQuery.Count(&total).Error; err != nil {
tableName := analyticsmodel.UserAccessLog{}.TableName()
countSQL := fmt.Sprintf("SELECT count() FROM %s WHERE %s", tableName, clause)
var total uint64
if err := db.ChConn.QueryRow(ctx, countSQL, args...).Scan(&total); err != nil {
return nil, 0, fmt.Errorf("count access logs: %w", err)
}
if total == 0 {
@@ -56,34 +66,41 @@ func ListAccessLogs(ctx context.Context, filter AccessLogFilter, page, pageSize
}
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
listSQL := fmt.Sprintf(`
SELECT id, user_id, path, method, ip, user_agent, headers, status, latency, created_at
FROM %s
WHERE %s
ORDER BY created_at DESC, id DESC
LIMIT ? OFFSET ?`, tableName, clause)
listArgs := append(append([]any{}, args...), pageSize, offset)
rows, err := db.ChConn.Query(ctx, listSQL, listArgs...)
if err != nil {
return nil, 0, fmt.Errorf("list access logs: %w", err)
}
defer func() { _ = rows.Close() }()
return logs, safeUint64Count(total), nil
}
func applyFilter(query *gorm.DB, filter AccessLogFilter) *gorm.DB {
if filter.UserIDs != nil {
if len(filter.UserIDs) == 0 {
return query.Where("1 = 0")
logs := make([]analyticsmodel.UserAccessLog, 0, pageSize)
for rows.Next() {
var (
item analyticsmodel.UserAccessLog
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, 0, fmt.Errorf("scan access log row: %w", err)
}
query = query.Where("user_id IN ?", filter.UserIDs)
item.CreatedAt = createdAt
logs = append(logs, item)
}
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
return logs, total, nil
}
@@ -3,7 +3,13 @@
package analytics
import "time"
import (
"fmt"
"strings"
"time"
)
const userAccessLogFilterClauseCapacity = 4
// AccessLogFilter scopes ClickHouse user access log queries.
type AccessLogFilter struct {
@@ -14,4 +20,37 @@ type AccessLogFilter struct {
StartTime *time.Time
// EndTime filters created_at <= EndTime when non-nil.
EndTime *time.Time
}
func buildUserAccessLogFilterClause(filter AccessLogFilter) (string, []any, bool) {
if filter.UserIDs != nil && len(filter.UserIDs) == 0 {
return "", nil, false
}
parts := make([]string, 0, userAccessLogFilterClauseCapacity)
args := make([]any, 0, userAccessLogFilterClauseCapacity)
if filter.UserIDs != nil {
placeholders := make([]string, len(filter.UserIDs))
for index, userID := range filter.UserIDs {
placeholders[index] = "?"
args = append(args, userID)
}
parts = append(parts, fmt.Sprintf("user_id IN (%s)", strings.Join(placeholders, ", ")))
}
if trimmed := strings.TrimSpace(filter.Path); trimmed != "" {
parts = append(parts, "path LIKE ?")
args = append(args, "%"+trimmed+"%")
}
if filter.StartTime != nil {
parts = append(parts, "created_at >= ?")
args = append(args, *filter.StartTime)
}
if filter.EndTime != nil {
parts = append(parts, "created_at <= ?")
args = append(args, *filter.EndTime)
}
if len(parts) == 0 {
return "1", args, true
}
return strings.Join(parts, " AND "), args, true
}
@@ -38,15 +38,12 @@ 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")
if err := userAccessLogConn(); err != nil {
return nil, err
}
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
@@ -55,24 +52,26 @@ func GetDailyTrend(ctx context.Context, days int) ([]DailyTrend, error) {
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 {
rows, err := db.ChConn.Query(ctx, query, startTime)
if err != nil {
return nil, fmt.Errorf("get daily trend: %w", err)
}
defer func() { _ = rows.Close() }()
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
for rows.Next() {
var (
date time.Time
count uint64
)
if err := rows.Scan(&date, &count); err != nil {
return nil, fmt.Errorf("scan daily trend row: %w", err)
}
trendMap[date.Format("2006-01-02")] = count
}
result := make([]DailyTrend, 0, days)
@@ -88,9 +87,8 @@ func GetDailyTrend(ctx context.Context, days int) ([]DailyTrend, error) {
// 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")
if err := userAccessLogConn(); err != nil {
return nil, err
}
tableName := analyticsmodel.UserAccessLog{}.TableName()
@@ -103,20 +101,23 @@ func GetBrowserDistribution(ctx context.Context, startTime time.Time) ([]Browser
LIMIT 100
`, tableName)
type uaRow struct {
UserAgent string
Count uint64
}
var rows []uaRow
if err := ch.Raw(query, startTime).Scan(&rows).Error; err != nil {
rows, err := db.ChConn.Query(ctx, query, startTime)
if err != nil {
return nil, fmt.Errorf("get browser distribution: %w", err)
}
defer func() { _ = rows.Close() }()
browserCounts := make(map[string]uint64)
for _, row := range rows {
browser := ParseBrowserName(row.UserAgent)
browserCounts[browser] += row.Count
for rows.Next() {
var (
userAgent string
count uint64
)
if err := rows.Scan(&userAgent, &count); err != nil {
return nil, fmt.Errorf("scan browser distribution row: %w", err)
}
browser := ParseBrowserName(userAgent)
browserCounts[browser] += count
}
result := make([]BrowserShare, 0, len(browserCounts))
@@ -137,10 +138,8 @@ func GetTopActiveUsers(ctx context.Context, startTime time.Time, limit int) ([]T
if limit < 1 {
limit = 10
}
ch := db.ChDB(ctx)
if ch == nil {
return nil, fmt.Errorf("clickhouse gorm connection is not initialized")
if err := userAccessLogConn(); err != nil {
return nil, err
}
tableName := analyticsmodel.UserAccessLog{}.TableName()
@@ -153,9 +152,19 @@ func GetTopActiveUsers(ctx context.Context, startTime time.Time, limit int) ([]T
LIMIT ?
`, tableName)
var users []TopUser
if err := ch.Raw(query, startTime, limit).Scan(&users).Error; err != nil {
rows, err := db.ChConn.Query(ctx, query, startTime, limit)
if err != nil {
return nil, fmt.Errorf("get top active users: %w", err)
}
defer func() { _ = rows.Close() }()
var users []TopUser
for rows.Next() {
var item TopUser
if err := rows.Scan(&item.UserID, &item.Count); err != nil {
return nil, fmt.Errorf("scan top active user row: %w", err)
}
users = append(users, item)
}
return users, nil
}
@@ -12,24 +12,10 @@ import (
"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
@@ -52,56 +38,24 @@ func TestParseBrowserName(t *testing.T) {
}
}
func TestCountAccessLogs_EmptyUserIDs(t *testing.T) {
setupChGormDB(t)
t.Cleanup(func() { db.SetChDBForTest(nil) })
func TestBuildUserAccessLogFilterClause_EmptyUserIDs(t *testing.T) {
_, _, ok := buildUserAccessLogFilterClause(AccessLogFilter{UserIDs: []uint64{}})
assert.False(t, ok)
}
func TestCountAccessLogs_EmptyUserIDs(t *testing.T) {
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)
@@ -5,13 +5,6 @@ package analytics
import "math"
func safeUint64Count(count int64) uint64 {
if count < 0 {
return 0
}
return uint64(count)
}
func safeInt64Count(count uint64) int64 {
if count > math.MaxInt64 {
return math.MaxInt64
@@ -31,24 +31,3 @@ func TestSafeInt64Count(t *testing.T) {
}
}
func TestSafeUint64Count(t *testing.T) {
t.Parallel()
tests := []struct {
name string
count int64
want uint64
}{
{name: "zero", count: 0, want: 0},
{name: "positive", count: 42, want: 42},
{name: "negative clamps", count: -1, want: 0},
}
for _, tt := range tests {
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
if got := safeUint64Count(tt.count); got != tt.want {
t.Fatalf("safeUint64Count(%d) = %d, want %d", tt.count, got, tt.want)
}
})
}
}
@@ -0,0 +1,74 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package analytics
import (
"context"
"fmt"
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
)
const (
// CleanupModeTTLMaterialize expires rows via table TTL instead of ALTER DELETE mutations.
CleanupModeTTLMaterialize = "ttl_materialize"
// CleanupModeTruncate removes all rows via TRUNCATE TABLE.
CleanupModeTruncate = "truncate"
)
// CleanupOutcome describes a non-mutation ClickHouse cleanup operation.
type CleanupOutcome struct {
EligibleCount int64
Mode string
}
func countClickHouseRows(ctx context.Context, conn driver.Conn, countSQL string, countArgs []any) (int64, error) {
var count uint64
if err := conn.QueryRow(ctx, countSQL, countArgs...).Scan(&count); err != nil {
return 0, fmt.Errorf("count clickhouse rows: %w", err)
}
return safeInt64Count(count), nil
}
func materializeTableTTL(ctx context.Context, conn driver.Conn, tableName string) error {
sql := fmt.Sprintf("ALTER TABLE %s MATERIALIZE TTL", tableName)
if err := conn.Exec(ctx, sql); err != nil {
return fmt.Errorf("materialize ttl on %s: %w", tableName, err)
}
return nil
}
func expireRowsViaTTL(ctx context.Context, conn driver.Conn, tableName string, countSQL string, countArgs []any) (CleanupOutcome, error) {
count, err := countClickHouseRows(ctx, conn, countSQL, countArgs)
if err != nil {
return CleanupOutcome{}, err
}
if count == 0 {
return CleanupOutcome{Mode: CleanupModeTTLMaterialize}, nil
}
if err := materializeTableTTL(ctx, conn, tableName); err != nil {
return CleanupOutcome{}, err
}
return CleanupOutcome{
EligibleCount: count,
Mode: CleanupModeTTLMaterialize,
}, nil
}
func truncateClickHouseTable(ctx context.Context, conn driver.Conn, tableName string) (CleanupOutcome, error) {
count, err := countClickHouseRows(ctx, conn, "SELECT count() FROM "+tableName, nil)
if err != nil {
return CleanupOutcome{}, err
}
if count == 0 {
return CleanupOutcome{Mode: CleanupModeTruncate}, nil
}
if err := conn.Exec(ctx, "TRUNCATE TABLE "+tableName); err != nil {
return CleanupOutcome{}, fmt.Errorf("truncate %s: %w", tableName, err)
}
return CleanupOutcome{
EligibleCount: count,
Mode: CleanupModeTruncate,
}, nil
}
@@ -0,0 +1,70 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package analytics
import (
"context"
"fmt"
"github.com/Rain-kl/Wavelet/internal/config"
"github.com/Rain-kl/Wavelet/internal/db"
)
// ClickHouseOperationalStats summarizes ClickHouse merge/mutation pressure.
type ClickHouseOperationalStats struct {
Database string `json:"database"`
ActiveParts int64 `json:"active_parts"`
TotalRows int64 `json:"total_rows"`
PendingMutations int64 `json:"pending_mutations"`
AsyncInsertQueue int64 `json:"async_insert_queue"`
AsyncInsertBytes int64 `json:"async_insert_bytes"`
}
// GetClickHouseOperationalStats returns operational metrics for the configured database.
func GetClickHouseOperationalStats(ctx context.Context) (*ClickHouseOperationalStats, error) {
if db.ChConn == nil {
return nil, fmt.Errorf("clickhouse native connection is not initialized")
}
database := config.Config.ClickHouse.Database
stats := &ClickHouseOperationalStats{Database: database}
partsSQL := `
SELECT
count() AS active_parts,
ifNull(sum(rows), 0) AS total_rows
FROM system.parts
WHERE active AND database = ?`
var activeParts, totalRows uint64
if err := db.ChConn.QueryRow(ctx, partsSQL, database).Scan(&activeParts, &totalRows); err != nil {
return nil, fmt.Errorf("query system.parts: %w", err)
}
stats.ActiveParts = safeInt64Count(activeParts)
stats.TotalRows = safeInt64Count(totalRows)
mutationsSQL := `
SELECT count()
FROM system.mutations
WHERE is_done = 0 AND database = ?`
if err := db.ChConn.QueryRow(ctx, mutationsSQL, database).Scan(&stats.PendingMutations); err != nil {
return nil, fmt.Errorf("query system.mutations: %w", err)
}
asyncSQL := `
SELECT
count() AS queue_entries,
ifNull(sum(bytes), 0) AS queue_bytes
FROM system.asynchronous_inserts
WHERE database = ?`
var queueEntries, queueBytes uint64
if err := db.ChConn.QueryRow(ctx, asyncSQL, database).Scan(&queueEntries, &queueBytes); err != nil {
// Older ClickHouse versions may not expose asynchronous_inserts; treat as optional.
stats.AsyncInsertQueue = 0
stats.AsyncInsertBytes = 0
} else {
stats.AsyncInsertQueue = safeInt64Count(queueEntries)
stats.AsyncInsertBytes = safeInt64Count(queueBytes)
}
return stats, nil
}
@@ -90,7 +90,7 @@ func CountNodeAccessLogs(ctx context.Context, filter NodeAccessLogFilter) (int64
countSQL := fmt.Sprintf(`
SELECT
count() AS total_records,
uniqExactIf(trim(remote_addr), trim(remote_addr) != '') AS total_ips
uniqExactIf(remote_addr, remote_addr != '') AS total_ips
FROM %s
WHERE %s`, tableName, clause)
var totalRecords, totalIPs uint64
@@ -11,50 +11,55 @@ import (
// DeleteAllNodeAccessLogs deletes all node access logs.
func DeleteAllNodeAccessLogs(ctx context.Context) (int64, error) {
tableName := nodeAccessLogTableName()
return deleteNodeAccessLogsWithCount(ctx, "SELECT count() FROM "+tableName, nil, "ALTER TABLE "+tableName+" DELETE WHERE 1")
}
// DeleteNodeAccessLogsBefore deletes logs older than cutoff.
func DeleteNodeAccessLogsBefore(ctx context.Context, cutoff time.Time) (int64, error) {
tableName := nodeAccessLogTableName()
cutoff = cutoff.UTC()
return deleteNodeAccessLogsWithCount(
ctx,
fmt.Sprintf("SELECT count() FROM %s WHERE logged_at < ?", tableName),
[]any{cutoff},
fmt.Sprintf("ALTER TABLE %s DELETE WHERE logged_at < ?", tableName),
cutoff,
)
}
// DeleteNodeAccessLogsByNodeBefore deletes logs for a node older than cutoff.
func DeleteNodeAccessLogsByNodeBefore(ctx context.Context, nodeID string, before time.Time) (int64, error) {
tableName := nodeAccessLogTableName()
before = before.UTC()
return deleteNodeAccessLogsWithCount(
ctx,
fmt.Sprintf("SELECT count() FROM %s WHERE node_id = ? AND logged_at < ?", tableName),
[]any{nodeID, before},
fmt.Sprintf("ALTER TABLE %s DELETE WHERE node_id = ? AND logged_at < ?", tableName),
nodeID, before,
)
}
func deleteNodeAccessLogsWithCount(ctx context.Context, countSQL string, countArgs []any, deleteSQL string, deleteArgs ...any) (int64, error) {
conn, err := nodeAccessLogConn()
if err != nil {
return 0, err
}
var count uint64
if err := conn.QueryRow(ctx, countSQL, countArgs...).Scan(&count); err != nil {
return 0, fmt.Errorf("count node access logs for delete: %w", err)
outcome, err := truncateClickHouseTable(ctx, conn, nodeAccessLogTableName())
if err != nil {
return 0, err
}
if count == 0 {
return 0, nil
return outcome.EligibleCount, nil
}
// DeleteNodeAccessLogsBefore expires logs older than cutoff via table TTL.
func DeleteNodeAccessLogsBefore(ctx context.Context, cutoff time.Time) (int64, error) {
conn, err := nodeAccessLogConn()
if err != nil {
return 0, err
}
if err := conn.Exec(ctx, deleteSQL, deleteArgs...); err != nil {
return 0, fmt.Errorf("delete node access logs: %w", err)
tableName := nodeAccessLogTableName()
cutoff = cutoff.UTC()
outcome, err := expireRowsViaTTL(
ctx,
conn,
tableName,
fmt.Sprintf("SELECT count() FROM %s WHERE logged_at < ?", tableName),
[]any{cutoff},
)
if err != nil {
return 0, err
}
return safeInt64Count(count), nil
return outcome.EligibleCount, nil
}
// DeleteNodeAccessLogsByNodeBefore expires logs for a node older than cutoff via table TTL.
func DeleteNodeAccessLogsByNodeBefore(ctx context.Context, nodeID string, before time.Time) (int64, error) {
conn, err := nodeAccessLogConn()
if err != nil {
return 0, err
}
tableName := nodeAccessLogTableName()
before = before.UTC()
outcome, err := expireRowsViaTTL(
ctx,
conn,
tableName,
fmt.Sprintf("SELECT count() FROM %s WHERE node_id = ? AND logged_at < ?", tableName),
[]any{nodeID, before},
)
if err != nil {
return 0, err
}
return outcome.EligibleCount, nil
}
@@ -40,7 +40,7 @@ func buildNodeAccessLogFilterClause(filter NodeAccessLogFilter) (string, []any)
parts = append(parts, "node_id = ?")
args = append(args, trimmed)
}
if trimmed := strings.TrimSpace(filter.RemoteAddr); trimmed != "" {
if trimmed := normalizeNodeAccessLogRemoteAddr(filter.RemoteAddr); trimmed != "" {
parts = append(parts, "remote_addr LIKE ?")
args = append(args, trimmed+"%")
}
@@ -95,6 +95,10 @@ func nodeAccessLogOrderClause(sortBy string, sortOrder string) string {
return column + " " + direction + ", logged_at " + direction + ", id " + direction
}
func normalizeNodeAccessLogRemoteAddr(value string) string {
return strings.TrimSpace(value)
}
func normalizeNodeAccessLogSortOrder(sortOrder string) string {
if strings.EqualFold(strings.TrimSpace(sortOrder), "asc") {
return "asc"
@@ -142,9 +146,9 @@ func nodeAccessLogIPSummaryOrderClause(sortBy string, sortOrder string) string {
case "last_seen_at":
column = "last_seen_epoch"
case nodeAccessLogColumnRemoteAddr:
column = "trimmed_remote_addr"
column = nodeAccessLogColumnRemoteAddr
}
return column + " " + direction + ", last_seen_epoch DESC, trimmed_remote_addr ASC"
return column + " " + direction + ", last_seen_epoch DESC, remote_addr ASC"
}
func nodeAccessLogTableName() string {
@@ -79,8 +79,8 @@ SELECT
countIf(status_code < 400) AS success_count,
countIf(status_code >= 400 AND status_code < 500) AS client_error_count,
countIf(status_code >= 500) AS server_error_count,
uniqExactIf(trim(remote_addr), trim(remote_addr) != '') AS unique_ip_count,
uniqExactIf(trim(host), trim(host) != '') AS unique_host_count
uniqExactIf(remote_addr, remote_addr != '') AS unique_ip_count,
uniqExactIf(host, host != '') AS unique_host_count
FROM %s
WHERE %s
GROUP BY bucket_epoch
@@ -186,26 +186,26 @@ func IPAggregatesNodeAccessLogs(ctx context.Context, filter NodeAccessLogFilter,
queryClause := clause
queryArgs := append([]any{}, args...)
if exactRemoteAddr {
trimmed := strings.TrimSpace(filter.RemoteAddr)
trimmed := normalizeNodeAccessLogRemoteAddr(filter.RemoteAddr)
if trimmed == "" {
return []NodeAccessLogIPAggregate{}, nil
}
queryClause = combineNodeAccessLogSQLClauses(queryClause, "trim(remote_addr) = ?")
queryClause = combineNodeAccessLogSQLClauses(queryClause, "remote_addr = ?")
queryArgs = append(queryArgs, trimmed)
}
lastSeenExpr := nodeAccessLogEpochExpr()
tableName := nodeAccessLogTableName()
sql := fmt.Sprintf(`
SELECT
trim(remote_addr) AS trimmed_remote_addr,
remote_addr,
count() AS request_count,
countIf(status_code < 400) AS success_count,
countIf(status_code >= 400 AND status_code < 500) AS client_error_count,
countIf(status_code >= 500) AS server_error_count,
max(%s) AS last_seen_epoch
FROM %s
WHERE %s AND trim(remote_addr) != ''
GROUP BY trimmed_remote_addr`, lastSeenExpr, tableName, queryClause)
WHERE %s AND remote_addr != ''
GROUP BY remote_addr`, lastSeenExpr, tableName, queryClause)
rows, err := conn.Query(ctx, sql, queryArgs...)
if err != nil {
return nil, fmt.Errorf("ip aggregates node access logs: %w", err)
@@ -252,13 +252,13 @@ func IPSummariesNodeAccessLogs(ctx context.Context, filter NodeAccessLogFilter,
tableName := nodeAccessLogTableName()
sql := fmt.Sprintf(`
SELECT
trim(remote_addr) AS trimmed_remote_addr,
remote_addr,
count() AS total_requests,
sum(%s) AS recent_requests,
max(%s) AS last_seen_epoch
FROM %s
WHERE %s AND trim(remote_addr) != ''
GROUP BY trimmed_remote_addr
WHERE %s AND remote_addr != ''
GROUP BY remote_addr
ORDER BY %s`, recentClause, lastSeenExpr, tableName, clause, nodeAccessLogIPSummaryOrderClause(filter.SortBy, filter.SortOrder))
if filter.PageSize > 0 {
if filter.Page < 0 {
@@ -305,8 +305,8 @@ func CountIPSummaryNodeAccessLogs(ctx context.Context, filter NodeAccessLogFilte
SELECT count() FROM (
SELECT 1
FROM %s
WHERE %s AND trim(remote_addr) != ''
GROUP BY trim(remote_addr)
WHERE %s AND remote_addr != ''
GROUP BY remote_addr
)`, tableName, clause)
var totalIPs uint64
if err := conn.QueryRow(ctx, sql, args...).Scan(&totalIPs); err != nil {
@@ -327,7 +327,7 @@ func IPAggregatesForWAFNodeAccessLogs(ctx context.Context, filter NodeAccessLogF
tableName := nodeAccessLogTableName()
sql := fmt.Sprintf(`
SELECT
trim(remote_addr) AS trimmed_remote_addr,
remote_addr,
count() AS request_count,
countIf(status_code = 404) AS status_404_count,
countIf(status_code >= 400 AND status_code < 500) AS client_error_count,
@@ -335,8 +335,8 @@ SELECT
countIf(%s) AS ip_host_count,
max(%s) AS last_seen_epoch
FROM %s
WHERE %s AND trim(remote_addr) != ''
GROUP BY trimmed_remote_addr`, hostIsIPExpr, lastSeenExpr, tableName, clause)
WHERE %s AND remote_addr != ''
GROUP BY remote_addr`, hostIsIPExpr, lastSeenExpr, tableName, clause)
rows, err := conn.Query(ctx, sql, args...)
if err != nil {
return nil, fmt.Errorf("ip aggregates for waf node access logs: %w", err)
@@ -394,12 +394,12 @@ func mergeWAFIPStatusCodeCounts(ctx context.Context, filter NodeAccessLogFilter,
tableName := nodeAccessLogTableName()
sql := fmt.Sprintf(`
SELECT
trim(remote_addr) AS trimmed_remote_addr,
remote_addr,
status_code,
count() AS status_count
FROM %s
WHERE %s AND trim(remote_addr) != ''
GROUP BY trimmed_remote_addr, status_code`, tableName, clause)
WHERE %s AND remote_addr != ''
GROUP BY remote_addr, status_code`, tableName, clause)
rows, err := conn.Query(ctx, sql, args...)
if err != nil {
return fmt.Errorf("waf ip status code counts: %w", err)
@@ -6,6 +6,7 @@ package analytics
import (
"context"
"fmt"
"time"
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
"github.com/Rain-kl/Wavelet/internal/db"
@@ -244,6 +245,59 @@ func scanNodeObsFrpsRows(rows driver.Rows) ([]analyticsmodel.NodeObsFrps, error)
return result, nil
}
const nodeTrafficHourlyTableName = "of_node_traffic_hourly"
// NodeTrafficHourly is an hourly traffic rollup row.
type NodeTrafficHourly struct {
NodeID string
Hour time.Time
RequestCount int64
ErrorCount int64
UniqueVisitorCount int64
}
// ListNodeTrafficHourly returns hourly traffic rollup rows matching filter.
func ListNodeTrafficHourly(ctx context.Context, filter NodeObservabilityFilter) ([]NodeTrafficHourly, error) {
conn, err := observabilityConn()
if err != nil {
return nil, err
}
clause, args := buildNodeObservabilityFilterClause(filter, "hour")
sql := fmt.Sprintf(`
SELECT
node_id,
hour,
sum(request_count) AS request_count,
sum(error_count) AS error_count,
sum(unique_visitor_count) AS unique_visitor_count
FROM %s
WHERE %s
GROUP BY node_id, hour
ORDER BY hour ASC`, nodeTrafficHourlyTableName, clause)
rows, err := conn.Query(ctx, sql, args...)
if err != nil {
return nil, fmt.Errorf("list node traffic hourly: %w", err)
}
defer func() { _ = rows.Close() }()
result := make([]NodeTrafficHourly, 0)
for rows.Next() {
var (
item NodeTrafficHourly
requestCount, errorCount, uniqueVisitorCount uint64
)
if err := rows.Scan(&item.NodeID, &item.Hour, &requestCount, &errorCount, &uniqueVisitorCount); err != nil {
return nil, fmt.Errorf("scan node traffic hourly row: %w", err)
}
item.Hour = item.Hour.UTC()
item.RequestCount = safeInt64Count(requestCount)
item.ErrorCount = safeInt64Count(errorCount)
item.UniqueVisitorCount = safeInt64Count(uniqueVisitorCount)
result = append(result, item)
}
return result, nil
}
func scanNodeObsFrpcRows(rows driver.Rows) ([]analyticsmodel.NodeObsFrpc, error) {
var result []analyticsmodel.NodeObsFrpc
for rows.Next() {
@@ -11,113 +11,170 @@ import (
// DeleteAllNodeMetricSnapshots deletes all node metric snapshots.
func DeleteAllNodeMetricSnapshots(ctx context.Context) (int64, error) {
tableName := nodeMetricSnapshotTableName()
return deleteNodeObservabilityWithCount(ctx, "SELECT count() FROM "+tableName, nil, "ALTER TABLE "+tableName+" DELETE WHERE 1")
}
// DeleteNodeMetricSnapshotsBefore deletes metric snapshots captured before cutoff.
func DeleteNodeMetricSnapshotsBefore(ctx context.Context, cutoff time.Time) (int64, error) {
tableName := nodeMetricSnapshotTableName()
cutoff = cutoff.UTC()
return deleteNodeObservabilityWithCount(
ctx,
fmt.Sprintf("SELECT count() FROM %s WHERE captured_at < ?", tableName),
[]any{cutoff},
fmt.Sprintf("ALTER TABLE %s DELETE WHERE captured_at < ?", tableName),
cutoff,
)
}
// DeleteAllNodeRequestReports deletes all node request reports.
func DeleteAllNodeRequestReports(ctx context.Context) (int64, error) {
tableName := nodeRequestReportTableName()
return deleteNodeObservabilityWithCount(ctx, "SELECT count() FROM "+tableName, nil, "ALTER TABLE "+tableName+" DELETE WHERE 1")
}
// DeleteNodeRequestReportsBefore deletes request reports ending before cutoff.
func DeleteNodeRequestReportsBefore(ctx context.Context, cutoff time.Time) (int64, error) {
tableName := nodeRequestReportTableName()
cutoff = cutoff.UTC()
return deleteNodeObservabilityWithCount(
ctx,
fmt.Sprintf("SELECT count() FROM %s WHERE window_ended_at < ?", tableName),
[]any{cutoff},
fmt.Sprintf("ALTER TABLE %s DELETE WHERE window_ended_at < ?", tableName),
cutoff,
)
}
// DeleteAllNodeObsOpenresty deletes all OpenResty observations.
func DeleteAllNodeObsOpenresty(ctx context.Context) (int64, error) {
tableName := nodeObsOpenrestyTableName()
return deleteNodeObservabilityWithCount(ctx, "SELECT count() FROM "+tableName, nil, "ALTER TABLE "+tableName+" DELETE WHERE 1")
}
// DeleteNodeObsOpenrestyBefore deletes OpenResty observations captured before cutoff.
func DeleteNodeObsOpenrestyBefore(ctx context.Context, cutoff time.Time) (int64, error) {
tableName := nodeObsOpenrestyTableName()
cutoff = cutoff.UTC()
return deleteNodeObservabilityWithCount(
ctx,
fmt.Sprintf("SELECT count() FROM %s WHERE captured_at < ?", tableName),
[]any{cutoff},
fmt.Sprintf("ALTER TABLE %s DELETE WHERE captured_at < ?", tableName),
cutoff,
)
}
// DeleteAllNodeObsFrps deletes all FRPS observations.
func DeleteAllNodeObsFrps(ctx context.Context) (int64, error) {
tableName := nodeObsFrpsTableName()
return deleteNodeObservabilityWithCount(ctx, "SELECT count() FROM "+tableName, nil, "ALTER TABLE "+tableName+" DELETE WHERE 1")
}
// DeleteNodeObsFrpsBefore deletes FRPS observations captured before cutoff.
func DeleteNodeObsFrpsBefore(ctx context.Context, cutoff time.Time) (int64, error) {
tableName := nodeObsFrpsTableName()
cutoff = cutoff.UTC()
return deleteNodeObservabilityWithCount(
ctx,
fmt.Sprintf("SELECT count() FROM %s WHERE captured_at < ?", tableName),
[]any{cutoff},
fmt.Sprintf("ALTER TABLE %s DELETE WHERE captured_at < ?", tableName),
cutoff,
)
}
// DeleteAllNodeObsFrpc deletes all FRPC observations.
func DeleteAllNodeObsFrpc(ctx context.Context) (int64, error) {
tableName := nodeObsFrpcTableName()
return deleteNodeObservabilityWithCount(ctx, "SELECT count() FROM "+tableName, nil, "ALTER TABLE "+tableName+" DELETE WHERE 1")
}
// DeleteNodeObsFrpcBefore deletes FRPC observations captured before cutoff.
func DeleteNodeObsFrpcBefore(ctx context.Context, cutoff time.Time) (int64, error) {
tableName := nodeObsFrpcTableName()
cutoff = cutoff.UTC()
return deleteNodeObservabilityWithCount(
ctx,
fmt.Sprintf("SELECT count() FROM %s WHERE captured_at < ?", tableName),
[]any{cutoff},
fmt.Sprintf("ALTER TABLE %s DELETE WHERE captured_at < ?", tableName),
cutoff,
)
}
func deleteNodeObservabilityWithCount(ctx context.Context, countSQL string, countArgs []any, deleteSQL string, deleteArgs ...any) (int64, error) {
conn, err := observabilityConn()
if err != nil {
return 0, err
}
var count uint64
if err := conn.QueryRow(ctx, countSQL, countArgs...).Scan(&count); err != nil {
return 0, fmt.Errorf("count node observability rows for delete: %w", err)
outcome, err := truncateClickHouseTable(ctx, conn, nodeMetricSnapshotTableName())
if err != nil {
return 0, err
}
if count == 0 {
return 0, nil
}
if err := conn.Exec(ctx, deleteSQL, deleteArgs...); err != nil {
return 0, fmt.Errorf("delete node observability rows: %w", err)
}
return safeInt64Count(count), nil
return outcome.EligibleCount, nil
}
// DeleteNodeMetricSnapshotsBefore expires metric snapshots captured before cutoff via table TTL.
func DeleteNodeMetricSnapshotsBefore(ctx context.Context, cutoff time.Time) (int64, error) {
conn, err := observabilityConn()
if err != nil {
return 0, err
}
tableName := nodeMetricSnapshotTableName()
cutoff = cutoff.UTC()
outcome, err := expireRowsViaTTL(
ctx,
conn,
tableName,
fmt.Sprintf("SELECT count() FROM %s WHERE captured_at < ?", tableName),
[]any{cutoff},
)
if err != nil {
return 0, err
}
return outcome.EligibleCount, nil
}
// DeleteAllNodeRequestReports deletes all node request reports.
func DeleteAllNodeRequestReports(ctx context.Context) (int64, error) {
conn, err := observabilityConn()
if err != nil {
return 0, err
}
outcome, err := truncateClickHouseTable(ctx, conn, nodeRequestReportTableName())
if err != nil {
return 0, err
}
return outcome.EligibleCount, nil
}
// DeleteNodeRequestReportsBefore expires request reports ending before cutoff via table TTL.
func DeleteNodeRequestReportsBefore(ctx context.Context, cutoff time.Time) (int64, error) {
conn, err := observabilityConn()
if err != nil {
return 0, err
}
tableName := nodeRequestReportTableName()
cutoff = cutoff.UTC()
outcome, err := expireRowsViaTTL(
ctx,
conn,
tableName,
fmt.Sprintf("SELECT count() FROM %s WHERE window_ended_at < ?", tableName),
[]any{cutoff},
)
if err != nil {
return 0, err
}
return outcome.EligibleCount, nil
}
// DeleteAllNodeObsOpenresty deletes all OpenResty observations.
func DeleteAllNodeObsOpenresty(ctx context.Context) (int64, error) {
conn, err := observabilityConn()
if err != nil {
return 0, err
}
outcome, err := truncateClickHouseTable(ctx, conn, nodeObsOpenrestyTableName())
if err != nil {
return 0, err
}
return outcome.EligibleCount, nil
}
// DeleteNodeObsOpenrestyBefore expires OpenResty observations captured before cutoff via table TTL.
func DeleteNodeObsOpenrestyBefore(ctx context.Context, cutoff time.Time) (int64, error) {
conn, err := observabilityConn()
if err != nil {
return 0, err
}
tableName := nodeObsOpenrestyTableName()
cutoff = cutoff.UTC()
outcome, err := expireRowsViaTTL(
ctx,
conn,
tableName,
fmt.Sprintf("SELECT count() FROM %s WHERE captured_at < ?", tableName),
[]any{cutoff},
)
if err != nil {
return 0, err
}
return outcome.EligibleCount, nil
}
// DeleteAllNodeObsFrps deletes all FRPS observations.
func DeleteAllNodeObsFrps(ctx context.Context) (int64, error) {
conn, err := observabilityConn()
if err != nil {
return 0, err
}
outcome, err := truncateClickHouseTable(ctx, conn, nodeObsFrpsTableName())
if err != nil {
return 0, err
}
return outcome.EligibleCount, nil
}
// DeleteNodeObsFrpsBefore expires FRPS observations captured before cutoff via table TTL.
func DeleteNodeObsFrpsBefore(ctx context.Context, cutoff time.Time) (int64, error) {
conn, err := observabilityConn()
if err != nil {
return 0, err
}
tableName := nodeObsFrpsTableName()
cutoff = cutoff.UTC()
outcome, err := expireRowsViaTTL(
ctx,
conn,
tableName,
fmt.Sprintf("SELECT count() FROM %s WHERE captured_at < ?", tableName),
[]any{cutoff},
)
if err != nil {
return 0, err
}
return outcome.EligibleCount, nil
}
// DeleteAllNodeObsFrpc deletes all FRPC observations.
func DeleteAllNodeObsFrpc(ctx context.Context) (int64, error) {
conn, err := observabilityConn()
if err != nil {
return 0, err
}
outcome, err := truncateClickHouseTable(ctx, conn, nodeObsFrpcTableName())
if err != nil {
return 0, err
}
return outcome.EligibleCount, nil
}
// DeleteNodeObsFrpcBefore expires FRPC observations captured before cutoff via table TTL.
func DeleteNodeObsFrpcBefore(ctx context.Context, cutoff time.Time) (int64, error) {
conn, err := observabilityConn()
if err != nil {
return 0, err
}
tableName := nodeObsFrpcTableName()
cutoff = cutoff.UTC()
outcome, err := expireRowsViaTTL(
ctx,
conn,
tableName,
fmt.Sprintf("SELECT count() FROM %s WHERE captured_at < ?", tableName),
[]any{cutoff},
)
if err != nil {
return 0, err
}
return outcome.EligibleCount, nil
}
+1
View File
@@ -51,6 +51,7 @@ func RegisterAdminRoutes(apiV1Router *gin.RouterGroup) {
func registerAdminDiagnosticRoutes(adminRouter *gin.RouterGroup) {
// System status
adminRouter.GET("/status", admin_status.GetSystemStatus)
adminRouter.GET("/status/clickhouse", admin_status.GetClickHouseStatus)
// Database basic info & backup export
adminRouter.GET("/db-info", admin_status.GetDatabaseInfo)