From 58624db39732ad9d00781e77eb8cd68a8a3e22aa Mon Sep 17 00:00:00 2001 From: ryan Date: Thu, 2 Jul 2026 15:47:38 +0800 Subject: [PATCH] =?UTF-8?q?perf(clickhouse):=20Phase=202=20legacy=20govern?= =?UTF-8?q?ance=20=E2=80=94=20TTL=20cleanup,=20unified=20pool,=20MV,=20ops?= =?UTF-8?q?=20API?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - 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 --- docs/changelog/index.md | 1 + docs/docs.go | 224 ++++++++++++++- docs/plan/clickhouse-cpu-optimization.md | 2 +- docs/swagger.json | 224 ++++++++++++++- docs/swagger.yaml | 142 +++++++++- internal/apps/admin/logs/routers.go | 4 +- internal/apps/admin/status/clickhouse.go | 40 +++ internal/apps/openflare/dashboard/logics.go | 6 +- .../apps/openflare/observability/analytics.go | 22 ++ .../openflare/observability/analytics_test.go | 17 ++ .../openflare/observability/node_logics.go | 6 +- .../apps/openflare/tasks/database_cleanup.go | 56 ++-- internal/db/clickhouse.go | 85 +----- ...07020002_create_node_traffic_hourly_mv.sql | 28 ++ internal/model/openflare_observability.go | 32 +++ internal/repository/analytics/access_log.go | 101 ++++--- .../repository/analytics/access_log_filter.go | 41 ++- .../repository/analytics/access_log_stats.go | 77 ++--- .../repository/analytics/access_log_test.go | 56 +--- .../repository/analytics/clickhouse_count.go | 7 - .../analytics/clickhouse_count_test.go | 21 -- .../analytics/clickhouse_maintenance.go | 74 +++++ .../repository/analytics/clickhouse_stats.go | 70 +++++ .../repository/analytics/node_access_log.go | 2 +- .../analytics/node_access_log_delete.go | 83 +++--- .../analytics/node_access_log_filter.go | 10 +- .../analytics/node_access_log_stats.go | 36 +-- .../analytics/node_observability.go | 54 ++++ .../analytics/node_observability_delete.go | 265 +++++++++++------- internal/router/v1/admin.go | 1 + 30 files changed, 1351 insertions(+), 436 deletions(-) create mode 100644 internal/apps/admin/status/clickhouse.go create mode 100644 internal/db/migrator/goose/clickhouse/202607020002_create_node_traffic_hourly_mv.sql create mode 100644 internal/repository/analytics/clickhouse_maintenance.go create mode 100644 internal/repository/analytics/clickhouse_stats.go diff --git a/docs/changelog/index.md b/docs/changelog/index.md index 3eabfecf..405fc536 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -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。 diff --git a/docs/docs.go b/docs/docs.go index 3b17ec0c..8bca2bb7 100644 --- a/docs/docs.go +++ b/docs/docs.go @@ -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": { diff --git a/docs/plan/clickhouse-cpu-optimization.md b/docs/plan/clickhouse-cpu-optimization.md index dd5e238d..16df5483 100644 --- a/docs/plan/clickhouse-cpu-optimization.md +++ b/docs/plan/clickhouse-cpu-optimization.md @@ -1,7 +1,7 @@ # ClickHouse CPU 性能优化计划 > PLAN_ID: `63ba981b` -> 状态: 已完成 +> 状态: 已完成(含 Phase 2 遗留治理) > 目标: 完成 P0–P2 优化,降低 ClickHouse CPU 占用 ## 背景 diff --git a/docs/swagger.json b/docs/swagger.json index b079dd27..677daec8 100644 --- a/docs/swagger.json +++ b/docs/swagger.json @@ -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": { diff --git a/docs/swagger.yaml b/docs/swagger.yaml index 2178e11a..4672242c 100644 --- a/docs/swagger.yaml +++ b/docs/swagger.yaml @@ -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 diff --git a/internal/apps/admin/logs/routers.go b/internal/apps/admin/logs/routers.go index a29be1eb..d85e004b 100644 --- a/internal/apps/admin/logs/routers.go +++ b/internal/apps/admin/logs/routers.go @@ -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 } diff --git a/internal/apps/admin/status/clickhouse.go b/internal/apps/admin/status/clickhouse.go new file mode 100644 index 00000000..a2f8ea05 --- /dev/null +++ b/internal/apps/admin/status/clickhouse.go @@ -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)) +} \ No newline at end of file diff --git a/internal/apps/openflare/dashboard/logics.go b/internal/apps/openflare/dashboard/logics.go index ce266fc5..0a11b8ef 100644 --- a/internal/apps/openflare/dashboard/logics.go +++ b/internal/apps/openflare/dashboard/logics.go @@ -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), diff --git a/internal/apps/openflare/observability/analytics.go b/internal/apps/openflare/observability/analytics.go index 8f648b30..dce95a8d 100644 --- a/internal/apps/openflare/observability/analytics.go +++ b/internal/apps/openflare/observability/analytics.go @@ -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) diff --git a/internal/apps/openflare/observability/analytics_test.go b/internal/apps/openflare/observability/analytics_test.go index c2bf53e1..967030b2 100644 --- a/internal/apps/openflare/observability/analytics_test.go +++ b/internal/apps/openflare/observability/analytics_test.go @@ -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() diff --git a/internal/apps/openflare/observability/node_logics.go b/internal/apps/openflare/observability/node_logics.go index a2f4d1ff..9efaaa56 100644 --- a/internal/apps/openflare/observability/node_logics.go +++ b/internal/apps/openflare/observability/node_logics.go @@ -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), diff --git a/internal/apps/openflare/tasks/database_cleanup.go b/internal/apps/openflare/tasks/database_cleanup.go index 4d70e181..6bbae7c0 100644 --- a/internal/apps/openflare/tasks/database_cleanup.go +++ b/internal/apps/openflare/tasks/database_cleanup.go @@ -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 } diff --git a/internal/db/clickhouse.go b/internal/db/clickhouse.go index 092203f1..cc087eb6 100644 --- a/internal/db/clickhouse.go +++ b/internal/db/clickhouse.go @@ -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. diff --git a/internal/db/migrator/goose/clickhouse/202607020002_create_node_traffic_hourly_mv.sql b/internal/db/migrator/goose/clickhouse/202607020002_create_node_traffic_hourly_mv.sql new file mode 100644 index 00000000..4f4e6667 --- /dev/null +++ b/internal/db/migrator/goose/clickhouse/202607020002_create_node_traffic_hourly_mv.sql @@ -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; \ No newline at end of file diff --git a/internal/model/openflare_observability.go b/internal/model/openflare_observability.go index 4872bd21..5a580ef2 100644 --- a/internal/model/openflare_observability.go +++ b/internal/model/openflare_observability.go @@ -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) diff --git a/internal/repository/analytics/access_log.go b/internal/repository/analytics/access_log.go index 317000cf..0bf52c81 100644 --- a/internal/repository/analytics/access_log.go +++ b/internal/repository/analytics/access_log.go @@ -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 } \ No newline at end of file diff --git a/internal/repository/analytics/access_log_filter.go b/internal/repository/analytics/access_log_filter.go index 984840da..2f81a91f 100644 --- a/internal/repository/analytics/access_log_filter.go +++ b/internal/repository/analytics/access_log_filter.go @@ -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 } \ No newline at end of file diff --git a/internal/repository/analytics/access_log_stats.go b/internal/repository/analytics/access_log_stats.go index 7187717e..2d0ed92e 100644 --- a/internal/repository/analytics/access_log_stats.go +++ b/internal/repository/analytics/access_log_stats.go @@ -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 } \ No newline at end of file diff --git a/internal/repository/analytics/access_log_test.go b/internal/repository/analytics/access_log_test.go index e0abce72..4ba181bc 100644 --- a/internal/repository/analytics/access_log_test.go +++ b/internal/repository/analytics/access_log_test.go @@ -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) diff --git a/internal/repository/analytics/clickhouse_count.go b/internal/repository/analytics/clickhouse_count.go index 1b19b475..168006c1 100644 --- a/internal/repository/analytics/clickhouse_count.go +++ b/internal/repository/analytics/clickhouse_count.go @@ -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 diff --git a/internal/repository/analytics/clickhouse_count_test.go b/internal/repository/analytics/clickhouse_count_test.go index c440cd73..c9cb4d72 100644 --- a/internal/repository/analytics/clickhouse_count_test.go +++ b/internal/repository/analytics/clickhouse_count_test.go @@ -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) - } - }) - } -} \ No newline at end of file diff --git a/internal/repository/analytics/clickhouse_maintenance.go b/internal/repository/analytics/clickhouse_maintenance.go new file mode 100644 index 00000000..57224d3c --- /dev/null +++ b/internal/repository/analytics/clickhouse_maintenance.go @@ -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 +} \ No newline at end of file diff --git a/internal/repository/analytics/clickhouse_stats.go b/internal/repository/analytics/clickhouse_stats.go new file mode 100644 index 00000000..ed4cd8f3 --- /dev/null +++ b/internal/repository/analytics/clickhouse_stats.go @@ -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 +} \ No newline at end of file diff --git a/internal/repository/analytics/node_access_log.go b/internal/repository/analytics/node_access_log.go index 005f94c6..333160f4 100644 --- a/internal/repository/analytics/node_access_log.go +++ b/internal/repository/analytics/node_access_log.go @@ -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 diff --git a/internal/repository/analytics/node_access_log_delete.go b/internal/repository/analytics/node_access_log_delete.go index a0593cc8..f910dd16 100644 --- a/internal/repository/analytics/node_access_log_delete.go +++ b/internal/repository/analytics/node_access_log_delete.go @@ -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 } \ No newline at end of file diff --git a/internal/repository/analytics/node_access_log_filter.go b/internal/repository/analytics/node_access_log_filter.go index d4962bb1..d5309f33 100644 --- a/internal/repository/analytics/node_access_log_filter.go +++ b/internal/repository/analytics/node_access_log_filter.go @@ -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 { diff --git a/internal/repository/analytics/node_access_log_stats.go b/internal/repository/analytics/node_access_log_stats.go index d7444318..adb8e489 100644 --- a/internal/repository/analytics/node_access_log_stats.go +++ b/internal/repository/analytics/node_access_log_stats.go @@ -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) diff --git a/internal/repository/analytics/node_observability.go b/internal/repository/analytics/node_observability.go index 7c76ce4b..db1232a8 100644 --- a/internal/repository/analytics/node_observability.go +++ b/internal/repository/analytics/node_observability.go @@ -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() { diff --git a/internal/repository/analytics/node_observability_delete.go b/internal/repository/analytics/node_observability_delete.go index 1f7c0fa3..978d9427 100644 --- a/internal/repository/analytics/node_observability_delete.go +++ b/internal/repository/analytics/node_observability_delete.go @@ -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 +} \ No newline at end of file diff --git a/internal/router/v1/admin.go b/internal/router/v1/admin.go index 819d42d1..f3900195 100644 --- a/internal/router/v1/admin.go +++ b/internal/router/v1/admin.go @@ -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)