diff --git a/docs/docs.go b/docs/docs.go index c1870837..67371613 100644 --- a/docs/docs.go +++ b/docs/docs.go @@ -49,13 +49,25 @@ const docTemplate = `{ "200": { "description": "成功返回 PoW 难题", "schema": { - "$ref": "#/definitions/cap.ChallengeResponse" + "allOf": [ + { + "$ref": "#/definitions/response.Any" + }, + { + "type": "object", + "properties": { + "data": { + "$ref": "#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.ChallengeResponse" + } + } + } + ] } }, "500": { "description": "内部服务错误", "schema": { - "$ref": "#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse" + "$ref": "#/definitions/response.Any" } } } @@ -89,19 +101,31 @@ const docTemplate = `{ "200": { "description": "核销成功,返回 X-Cap-Token", "schema": { - "$ref": "#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse" + "allOf": [ + { + "$ref": "#/definitions/response.Any" + }, + { + "type": "object", + "properties": { + "data": { + "$ref": "#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse" + } + } + } + ] } }, "400": { "description": "参数错误或核销失败", "schema": { - "$ref": "#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse" + "$ref": "#/definitions/response.Any" } }, "500": { "description": "内部服务错误", "schema": { - "$ref": "#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse" + "$ref": "#/definitions/response.Any" } } } @@ -1000,7 +1024,7 @@ const docTemplate = `{ "SessionCookie": [] } ], - "description": "分页并按照用户、接口路径、时间范围等维度检索 ClickHouse 用户访问日志列表(需要管理员权限,ClickHouse 未启用时报错)", + "description": "分页并按照用户、接口路径、时间范围等维度检索用户访问日志列表(需要管理员权限)", "produces": [ "application/json" ], @@ -1068,7 +1092,7 @@ const docTemplate = `{ } }, "400": { - "description": "ClickHouse 未启用或参数错误", + "description": "参数错误", "schema": { "$ref": "#/definitions/response.Any" } @@ -1084,6 +1108,12 @@ const docTemplate = `{ "schema": { "$ref": "#/definitions/response.Any" } + }, + "500": { + "description": "内部错误", + "schema": { + "$ref": "#/definitions/response.Any" + } } } } @@ -1095,7 +1125,7 @@ const docTemplate = `{ "SessionCookie": [] } ], - "description": "聚合统计最近 7 天的每日访问趋势、浏览器分布以及前 10 名最活跃用户排行(需要管理员权限,ClickHouse 未启用时报错)", + "description": "聚合统计最近 7 天的每日访问趋势、浏览器分布以及前 10 名最活跃用户排行(需要管理员权限)", "produces": [ "application/json" ], @@ -1122,12 +1152,6 @@ const docTemplate = `{ ] } }, - "400": { - "description": "ClickHouse 未启用", - "schema": { - "$ref": "#/definitions/response.Any" - } - }, "401": { "description": "未登录", "schema": { @@ -1139,6 +1163,12 @@ const docTemplate = `{ "schema": { "$ref": "#/definitions/response.Any" } + }, + "500": { + "description": "内部错误", + "schema": { + "$ref": "#/definitions/response.Any" + } } } } @@ -1853,6 +1883,61 @@ const docTemplate = `{ } } }, + "/api/v1/admin/status/log-database": { + "get": { + "security": [ + { + "SessionCookie": [] + } + ], + "description": "返回当前日志主库、迁移状态、各库保留天数与合法迁移目标,需要管理员权限", + "produces": [ + "application/json" + ], + "tags": [ + "admin" + ], + "summary": "获取日志数据库状态", + "responses": { + "200": { + "description": "获取成功", + "schema": { + "allOf": [ + { + "$ref": "#/definitions/response.Any" + }, + { + "type": "object", + "properties": { + "data": { + "$ref": "#/definitions/status.LogDatabaseStatus" + } + } + } + ] + } + }, + "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": [ @@ -3673,6 +3758,11 @@ const docTemplate = `{ ], "summary": "获取用户列表", "parameters": [ + { + "type": "string", + "name": "email", + "in": "query" + }, { "minimum": 1, "type": "integer", @@ -3891,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": [ { @@ -4079,7 +4255,7 @@ const docTemplate = `{ }, "/api/v1/custom/hello": { "get": { - "description": "A sample business API for customization", + "description": "Scaffold demo API; product APIs use semantic paths under apps/\u003cdomain\u003e", "produces": [ "application/json" ], @@ -5531,32 +5707,6 @@ const docTemplate = `{ } } }, - "cap.ChallengeResponse": { - "type": "object", - "properties": { - "challenge": { - "type": "object", - "properties": { - "c": { - "type": "integer" - }, - "d": { - "type": "integer" - }, - "s": { - "type": "integer" - } - } - }, - "expires": { - "description": "ms timestamp", - "type": "integer" - }, - "token": { - "type": "string" - } - } - }, "cap.challengeRequest": { "type": "object", "properties": { @@ -5671,6 +5821,32 @@ const docTemplate = `{ } } }, + "github_com_Rain-kl_Wavelet_internal_apps_cap.ChallengeResponse": { + "type": "object", + "properties": { + "challenge": { + "type": "object", + "properties": { + "c": { + "type": "integer" + }, + "d": { + "type": "integer" + }, + "s": { + "type": "integer" + } + } + }, + "expires": { + "description": "ms timestamp", + "type": "integer" + }, + "token": { + "type": "string" + } + } + }, "github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse": { "type": "object", "properties": { @@ -6884,6 +7060,29 @@ const docTemplate = `{ } } }, + "status.LogDatabaseStatus": { + "type": "object", + "properties": { + "active_database": { + "type": "string" + }, + "available_targets": { + "type": "array", + "items": { + "type": "string" + } + }, + "migration": { + "type": "string" + }, + "retention_days": { + "type": "object", + "additionalProperties": { + "type": "integer" + } + } + } + }, "status.SystemStatusResponse": { "type": "object", "properties": { @@ -7483,6 +7682,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/swagger.json b/docs/swagger.json index 07f5eb39..1c4df5ed 100644 --- a/docs/swagger.json +++ b/docs/swagger.json @@ -42,13 +42,25 @@ "200": { "description": "成功返回 PoW 难题", "schema": { - "$ref": "#/definitions/cap.ChallengeResponse" + "allOf": [ + { + "$ref": "#/definitions/response.Any" + }, + { + "type": "object", + "properties": { + "data": { + "$ref": "#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.ChallengeResponse" + } + } + } + ] } }, "500": { "description": "内部服务错误", "schema": { - "$ref": "#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse" + "$ref": "#/definitions/response.Any" } } } @@ -82,19 +94,31 @@ "200": { "description": "核销成功,返回 X-Cap-Token", "schema": { - "$ref": "#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse" + "allOf": [ + { + "$ref": "#/definitions/response.Any" + }, + { + "type": "object", + "properties": { + "data": { + "$ref": "#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse" + } + } + } + ] } }, "400": { "description": "参数错误或核销失败", "schema": { - "$ref": "#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse" + "$ref": "#/definitions/response.Any" } }, "500": { "description": "内部服务错误", "schema": { - "$ref": "#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse" + "$ref": "#/definitions/response.Any" } } } @@ -993,7 +1017,7 @@ "SessionCookie": [] } ], - "description": "分页并按照用户、接口路径、时间范围等维度检索 ClickHouse 用户访问日志列表(需要管理员权限,ClickHouse 未启用时报错)", + "description": "分页并按照用户、接口路径、时间范围等维度检索用户访问日志列表(需要管理员权限)", "produces": [ "application/json" ], @@ -1061,7 +1085,7 @@ } }, "400": { - "description": "ClickHouse 未启用或参数错误", + "description": "参数错误", "schema": { "$ref": "#/definitions/response.Any" } @@ -1077,6 +1101,12 @@ "schema": { "$ref": "#/definitions/response.Any" } + }, + "500": { + "description": "内部错误", + "schema": { + "$ref": "#/definitions/response.Any" + } } } } @@ -1088,7 +1118,7 @@ "SessionCookie": [] } ], - "description": "聚合统计最近 7 天的每日访问趋势、浏览器分布以及前 10 名最活跃用户排行(需要管理员权限,ClickHouse 未启用时报错)", + "description": "聚合统计最近 7 天的每日访问趋势、浏览器分布以及前 10 名最活跃用户排行(需要管理员权限)", "produces": [ "application/json" ], @@ -1115,12 +1145,6 @@ ] } }, - "400": { - "description": "ClickHouse 未启用", - "schema": { - "$ref": "#/definitions/response.Any" - } - }, "401": { "description": "未登录", "schema": { @@ -1132,6 +1156,12 @@ "schema": { "$ref": "#/definitions/response.Any" } + }, + "500": { + "description": "内部错误", + "schema": { + "$ref": "#/definitions/response.Any" + } } } } @@ -1846,6 +1876,61 @@ } } }, + "/api/v1/admin/status/log-database": { + "get": { + "security": [ + { + "SessionCookie": [] + } + ], + "description": "返回当前日志主库、迁移状态、各库保留天数与合法迁移目标,需要管理员权限", + "produces": [ + "application/json" + ], + "tags": [ + "admin" + ], + "summary": "获取日志数据库状态", + "responses": { + "200": { + "description": "获取成功", + "schema": { + "allOf": [ + { + "$ref": "#/definitions/response.Any" + }, + { + "type": "object", + "properties": { + "data": { + "$ref": "#/definitions/status.LogDatabaseStatus" + } + } + } + ] + } + }, + "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": [ @@ -3666,6 +3751,11 @@ ], "summary": "获取用户列表", "parameters": [ + { + "type": "string", + "name": "email", + "in": "query" + }, { "minimum": 1, "type": "integer", @@ -3884,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": [ { @@ -4072,7 +4248,7 @@ }, "/api/v1/custom/hello": { "get": { - "description": "A sample business API for customization", + "description": "Scaffold demo API; product APIs use semantic paths under apps/\u003cdomain\u003e", "produces": [ "application/json" ], @@ -5524,32 +5700,6 @@ } } }, - "cap.ChallengeResponse": { - "type": "object", - "properties": { - "challenge": { - "type": "object", - "properties": { - "c": { - "type": "integer" - }, - "d": { - "type": "integer" - }, - "s": { - "type": "integer" - } - } - }, - "expires": { - "description": "ms timestamp", - "type": "integer" - }, - "token": { - "type": "string" - } - } - }, "cap.challengeRequest": { "type": "object", "properties": { @@ -5664,6 +5814,32 @@ } } }, + "github_com_Rain-kl_Wavelet_internal_apps_cap.ChallengeResponse": { + "type": "object", + "properties": { + "challenge": { + "type": "object", + "properties": { + "c": { + "type": "integer" + }, + "d": { + "type": "integer" + }, + "s": { + "type": "integer" + } + } + }, + "expires": { + "description": "ms timestamp", + "type": "integer" + }, + "token": { + "type": "string" + } + } + }, "github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse": { "type": "object", "properties": { @@ -6877,6 +7053,29 @@ } } }, + "status.LogDatabaseStatus": { + "type": "object", + "properties": { + "active_database": { + "type": "string" + }, + "available_targets": { + "type": "array", + "items": { + "type": "string" + } + }, + "migration": { + "type": "string" + }, + "retention_days": { + "type": "object", + "additionalProperties": { + "type": "integer" + } + } + } + }, "status.SystemStatusResponse": { "type": "object", "properties": { @@ -7476,6 +7675,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 0554159e..fc1b970c 100644 --- a/docs/swagger.yaml +++ b/docs/swagger.yaml @@ -40,23 +40,6 @@ definitions: - max_size_mb - ttl_minutes type: object - cap.ChallengeResponse: - properties: - challenge: - properties: - c: - type: integer - d: - type: integer - s: - type: integer - type: object - expires: - description: ms timestamp - type: integer - token: - type: string - type: object cap.challengeRequest: properties: scope: @@ -132,6 +115,23 @@ definitions: ttl_minutes: type: integer type: object + github_com_Rain-kl_Wavelet_internal_apps_cap.ChallengeResponse: + properties: + challenge: + properties: + c: + type: integer + d: + type: integer + s: + type: integer + type: object + expires: + description: ms timestamp + type: integer + token: + type: string + type: object github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse: properties: error: @@ -949,6 +949,21 @@ definitions: version: type: string type: object + status.LogDatabaseStatus: + properties: + active_database: + type: string + available_targets: + items: + type: string + type: array + migration: + type: string + retention_days: + additionalProperties: + type: integer + type: object + type: object status.SystemStatusResponse: properties: alloc: @@ -1358,6 +1373,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: @@ -1425,11 +1457,16 @@ paths: "200": description: 成功返回 PoW 难题 schema: - $ref: '#/definitions/cap.ChallengeResponse' + allOf: + - $ref: '#/definitions/response.Any' + - properties: + data: + $ref: '#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.ChallengeResponse' + type: object "500": description: 内部服务错误 schema: - $ref: '#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse' + $ref: '#/definitions/response.Any' summary: 生成人机验证难题 tags: - cap @@ -1451,15 +1488,20 @@ paths: "200": description: 核销成功,返回 X-Cap-Token schema: - $ref: '#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse' + allOf: + - $ref: '#/definitions/response.Any' + - properties: + data: + $ref: '#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse' + type: object "400": description: 参数错误或核销失败 schema: - $ref: '#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse' + $ref: '#/definitions/response.Any' "500": description: 内部服务错误 schema: - $ref: '#/definitions/github_com_Rain-kl_Wavelet_internal_apps_cap.RedeemResponse' + $ref: '#/definitions/response.Any' summary: 校验人机验证解答 tags: - cap @@ -1997,7 +2039,7 @@ paths: - admin /api/v1/admin/logs/access: get: - description: 分页并按照用户、接口路径、时间范围等维度检索 ClickHouse 用户访问日志列表(需要管理员权限,ClickHouse 未启用时报错) + description: 分页并按照用户、接口路径、时间范围等维度检索用户访问日志列表(需要管理员权限) parameters: - default: 1 description: 页码 @@ -2038,7 +2080,7 @@ paths: $ref: '#/definitions/logs.accessLogsResponse' type: object "400": - description: ClickHouse 未启用或参数错误 + description: 参数错误 schema: $ref: '#/definitions/response.Any' "401": @@ -2049,6 +2091,10 @@ paths: description: 无管理员权限 schema: $ref: '#/definitions/response.Any' + "500": + description: 内部错误 + schema: + $ref: '#/definitions/response.Any' security: - SessionCookie: [] summary: 获取用户访问日志 @@ -2056,7 +2102,7 @@ paths: - admin /api/v1/admin/logs/analytics: get: - description: 聚合统计最近 7 天的每日访问趋势、浏览器分布以及前 10 名最活跃用户排行(需要管理员权限,ClickHouse 未启用时报错) + description: 聚合统计最近 7 天的每日访问趋势、浏览器分布以及前 10 名最活跃用户排行(需要管理员权限) produces: - application/json responses: @@ -2069,10 +2115,6 @@ paths: data: $ref: '#/definitions/logs.logsAnalyticsResponse' type: object - "400": - description: ClickHouse 未启用 - schema: - $ref: '#/definitions/response.Any' "401": description: 未登录 schema: @@ -2081,6 +2123,10 @@ paths: description: 无管理员权限 schema: $ref: '#/definitions/response.Any' + "500": + description: 内部错误 + schema: + $ref: '#/definitions/response.Any' security: - SessionCookie: [] summary: 获取访问日志分析数据 @@ -2496,6 +2542,38 @@ paths: summary: 获取系统状态信息 tags: - admin + /api/v1/admin/status/log-database: + get: + description: 返回当前日志主库、迁移状态、各库保留天数与合法迁移目标,需要管理员权限 + produces: + - application/json + responses: + "200": + description: 获取成功 + schema: + allOf: + - $ref: '#/definitions/response.Any' + - properties: + data: + $ref: '#/definitions/status.LogDatabaseStatus' + type: object + "401": + description: 未登录 + schema: + $ref: '#/definitions/response.Any' + "403": + description: 无管理员权限 + schema: + $ref: '#/definitions/response.Any' + "500": + description: 内部错误 + schema: + $ref: '#/definitions/response.Any' + security: + - SessionCookie: [] + summary: 获取日志数据库状态 + tags: + - admin /api/v1/admin/system-configs: get: description: 返回所有系统配置列表,支持按配置类型(system/business)过滤,需要管理员权限 @@ -3589,6 +3667,9 @@ paths: get: description: 分页返回用户列表,支持按用户 ID 和用户名筛选,需要管理员权限 parameters: + - in: query + name: email + type: string - in: query minimum: 1 name: page @@ -3772,6 +3853,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: @@ -3843,7 +3977,7 @@ paths: - config /api/v1/custom/hello: get: - description: A sample business API for customization + description: Scaffold demo API; product APIs use semantic paths under apps/ produces: - application/json responses: diff --git a/frontend/app/(main)/admin/tasks/components/task-manager.tsx b/frontend/app/(main)/admin/tasks/components/task-manager.tsx index 788e7f99..7fe0587b 100644 --- a/frontend/app/(main)/admin/tasks/components/task-manager.tsx +++ b/frontend/app/(main)/admin/tasks/components/task-manager.tsx @@ -1,6 +1,6 @@ 'use client'; -import { useCallback, useEffect, useState } from 'react'; +import { useCallback, useEffect, useMemo, useState } from 'react'; import { useTranslations } from 'next-intl'; import { toast } from 'sonner'; import { Button } from '@/components/ui/button'; @@ -20,12 +20,24 @@ import { import { Calendar as CalendarIcon, Clock, + Database, Info, Layers, Play, } from 'lucide-react'; +import { + Select, + SelectContent, + SelectItem, + SelectTrigger, + SelectValue, +} from '@/components/ui/select'; -import type { DispatchTaskRequest, TaskMeta } from '@/lib/services/admin'; +import type { + DispatchTaskRequest, + LogDatabaseStatus, + TaskMeta, +} from '@/lib/services/admin'; import services from '@/lib/services'; import { buildTaskPayload } from '@/lib/task-param-utils'; import { ErrorInline } from '@/components/layout/error'; @@ -69,6 +81,18 @@ const TASK_CONFIGS: Record< gradient: 'from-rose-500/10 via-rose-500/5 to-transparent border-rose-200/50 dark:border-rose-800/50 hover:border-rose-400 dark:hover:border-rose-500', }, + logs_db_switch: { + icon: Database, + color: 'text-teal-600 dark:text-teal-400', + gradient: + 'from-teal-500/10 via-teal-500/5 to-transparent border-teal-200/50 dark:border-teal-800/50 hover:border-teal-400 dark:hover:border-teal-500', + }, +}; + +const LOG_DATABASE_LABELS: Record = { + postgres: 'PostgreSQL(主库)', + sqlite: 'SQLite(主库)', + clickhouse: 'ClickHouse', }; const DEFAULT_TASK_CONFIG = { @@ -196,9 +220,37 @@ export function TaskManager() { } }, [t]); + const [logDbStatus, setLogDbStatus] = useState( + null, + ); + + const fetchLogDbStatus = useCallback(async () => { + try { + const data = await services.adminStatus.getLogDatabaseStatus(); + setLogDbStatus(data); + } catch { + setLogDbStatus(null); + } + }, []); + useEffect(() => { fetchTaskTypes(); - }, [fetchTaskTypes]); + fetchLogDbStatus(); + }, [fetchTaskTypes, fetchLogDbStatus]); + + const availableLogDbTargets = useMemo( + () => logDbStatus?.available_targets ?? [], + [logDbStatus], + ); + + const retentionSummary = useMemo(() => { + const days = logDbStatus?.retention_days ?? {}; + const parts: string[] = []; + if (days.postgres != null) parts.push(`PG ${days.postgres}`); + if (days.sqlite != null) parts.push(`SQLite ${days.sqlite}`); + if (days.clickhouse != null) parts.push(`CH ${days.clickhouse}`); + return parts.join(' / '); + }, [logDbStatus]); useEffect(() => { if (selectedTaskType) { @@ -353,6 +405,45 @@ export function TaskManager() { + {task.type === 'logs_db_switch' && logDbStatus && ( +
+
+ + 日志主库 + + + {LOG_DATABASE_LABELS[logDbStatus.active_database] || + logDbStatus.active_database} + +
+
+ + 保留天数 + + + {retentionSummary || '-'} + +
+
+ + 迁移状态 + + + {logDbStatus.migration === 'migrating' + ? '迁移中' + : '空闲'} + +
+
+ )} +
+ ) : param.name === 'target' && + getSelectedTaskMeta()?.type === 'logs_db_switch' && + availableLogDbTargets.length > 0 ? ( + ) : ( ('/status'); } + static async getLogDatabaseStatus(): Promise { + return this.get('/status/log-database'); + } + static async getUpdateStatus(): Promise { return this.get('/update'); } diff --git a/frontend/lib/services/admin/types.ts b/frontend/lib/services/admin/types.ts index ca9a46d0..51e52f7c 100644 --- a/frontend/lib/services/admin/types.ts +++ b/frontend/lib/services/admin/types.ts @@ -401,6 +401,20 @@ export interface ToggleAuthSourceRequest { is_active: boolean; } +/** + * 日志数据库状态 + */ +export interface LogDatabaseStatus { + /** 当前日志主库:postgres | sqlite | clickhouse */ + active_database: string; + /** 迁移状态:idle | migrating */ + migration: string; + /** 各日志库保留天数 */ + retention_days: Record; + /** 当前主库的合法迁移目标 */ + available_targets: string[]; +} + /** * 系统状态信息 */ diff --git a/frontend/lib/services/index.ts b/frontend/lib/services/index.ts index c101ee5f..fd2f0240 100644 --- a/frontend/lib/services/index.ts +++ b/frontend/lib/services/index.ts @@ -124,6 +124,7 @@ export type { CreateUserRequest, UpdateUserRequest, SystemStatus, + LogDatabaseStatus, AppUpdateStatus, Schedule, CreateScheduleRequest, diff --git a/internal/apps/admin/logs/routers.go b/internal/apps/admin/logs/routers.go index ab333651..73a10181 100644 --- a/internal/apps/admin/logs/routers.go +++ b/internal/apps/admin/logs/routers.go @@ -14,10 +14,9 @@ import ( "time" "github.com/Rain-kl/Wavelet/internal/apps/admin" - "github.com/Rain-kl/Wavelet/internal/infra/config" - "github.com/Rain-kl/Wavelet/internal/infra/persistence" "github.com/Rain-kl/Wavelet/internal/repository" analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics" + "github.com/Rain-kl/Wavelet/internal/repository/logstore" "github.com/Rain-kl/Wavelet/pkg/logger" "github.com/gin-gonic/gin" @@ -227,7 +226,7 @@ func enrichAccessLogsWithUsers(ctx context.Context, list []accessLogItem) { // GetAccessLogs 获取 ClickHouse 异步采集的访问日志 // @Summary 获取用户访问日志 -// @Description 分页并按照用户、接口路径、时间范围等维度检索 ClickHouse 用户访问日志列表(需要管理员权限,ClickHouse 未启用时报错) +// @Description 分页并按照用户、接口路径、时间范围等维度检索用户访问日志列表(需要管理员权限) // @Tags admin // @Produce json // @Security SessionCookie @@ -238,14 +237,16 @@ func enrichAccessLogsWithUsers(ctx context.Context, list []accessLogItem) { // @Param start_time query string false "起始时间(RFC3339 或 YYYY-MM-DD HH:MM:SS)" // @Param end_time query string false "结束时间(RFC3339 或 YYYY-MM-DD HH:MM:SS)" // @Success 200 {object} response.Any{data=logs.accessLogsResponse} "访问日志列表" -// @Failure 400 {object} response.Any "ClickHouse 未启用或参数错误" +// @Failure 400 {object} response.Any "参数错误" // @Failure 401 {object} response.Any "未登录" // @Failure 403 {object} response.Any "无管理员权限" +// @Failure 500 {object} response.Any "内部错误" // @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 { - response.AbortWithError(c, http.StatusBadRequest, "ClickHouse 存储服务未启用,无法检索访问日志") + store, err := logstore.Active(ctx) + if err != nil { + response.AbortInternal(c, "日志存储初始化失败") return } @@ -271,7 +272,7 @@ func GetAccessLogs(c *gin.Context) { return } - logs, total, err := analyticsrepo.ListAccessLogs(ctx, filter, page, pageSize) + logs, total, err := store.UserAccessLogs.List(ctx, filter, page, pageSize) if err != nil { response.AbortWithError(c, http.StatusInternalServerError, err.Error()) return @@ -333,25 +334,26 @@ type logsAnalyticsResponse struct { // GetLogsAnalytics 获取 ClickHouse 访问日志图表聚合指标 // @Summary 获取访问日志分析数据 -// @Description 聚合统计最近 7 天的每日访问趋势、浏览器分布以及前 10 名最活跃用户排行(需要管理员权限,ClickHouse 未启用时报错) +// @Description 聚合统计最近 7 天的每日访问趋势、浏览器分布以及前 10 名最活跃用户排行(需要管理员权限) // @Tags admin // @Produce json // @Security SessionCookie // @Success 200 {object} response.Any{data=logs.logsAnalyticsResponse} "分析统计数据" -// @Failure 400 {object} response.Any "ClickHouse 未启用" +// @Failure 500 {object} response.Any "内部错误" // @Failure 401 {object} response.Any "未登录" // @Failure 403 {object} response.Any "无管理员权限" // @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 { - response.AbortWithError(c, http.StatusBadRequest, "ClickHouse 存储服务未启用,无法获取分析数据") + store, err := logstore.Active(ctx) + if err != nil { + response.AbortInternal(c, "日志存储初始化失败") return } startTime := time.Now().AddDate(0, 0, -(analyticsDays - 1)).Truncate(hoursInDay * time.Hour) - trendPoints, err := analyticsrepo.GetDailyTrend(ctx, analyticsDays) + trendPoints, err := store.UserAccessLogs.GetDailyTrend(ctx, analyticsDays) if err != nil { response.AbortWithError(c, http.StatusInternalServerError, "查询访问趋势失败: "+err.Error()) return @@ -364,7 +366,7 @@ func GetLogsAnalytics(c *gin.Context) { } } - browserPoints, err := analyticsrepo.GetBrowserDistribution(ctx, startTime) + browserPoints, err := store.UserAccessLogs.GetBrowserDistribution(ctx, startTime) if err != nil { response.AbortWithError(c, http.StatusInternalServerError, "查询浏览器分布失败: "+err.Error()) return @@ -377,7 +379,7 @@ func GetLogsAnalytics(c *gin.Context) { } } - topUserPoints, err := analyticsrepo.GetTopActiveUsers(ctx, startTime, topActiveLimit) + topUserPoints, err := store.UserAccessLogs.GetTopActiveUsers(ctx, startTime, topActiveLimit) if err != nil { response.AbortWithError(c, http.StatusInternalServerError, "查询活跃用户失败: "+err.Error()) return diff --git a/internal/apps/admin/logs/tasks.go b/internal/apps/admin/logs/tasks.go new file mode 100644 index 00000000..e03369dc --- /dev/null +++ b/internal/apps/admin/logs/tasks.go @@ -0,0 +1,219 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package logs + +import ( + "context" + "encoding/json" + "errors" + "fmt" + + "github.com/Rain-kl/Wavelet/internal/apps/risk_control" + "github.com/Rain-kl/Wavelet/internal/infra/config" + "github.com/Rain-kl/Wavelet/internal/infra/task" + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/internal/repository" + "github.com/Rain-kl/Wavelet/internal/repository/logstore" + "github.com/Rain-kl/Wavelet/pkg/logger" +) + +const ( + // LogDBSwitchTask 切换日志数据库任务标识。 + LogDBSwitchTask = "logs:db_switch" + // TaskTypeLogDBSwitch 管理端任务类型。 + TaskTypeLogDBSwitch = "logs_db_switch" + + copyBatchSize = 1000 + targetPostgres = "postgres" + targetSQLite = "sqlite" + targetClickHouse = "clickhouse" +) + +// LogDBSwitchMeta 描述切换日志数据库任务。 +var LogDBSwitchMeta = task.TaskMeta{ + Type: TaskTypeLogDBSwitch, + AsynqTask: LogDBSwitchTask, + Name: "切换日志数据库", + Description: "复制迁移用户访问日志并在成功后切换日志主库(期间禁止日志写入)", + SupportsTime: false, + MaxRetry: task.DefaultMaxRetry, + Queue: task.QueueDefault, + Retryable: true, + Params: []task.TaskParam{ + {Name: "target", Label: "目标日志库", Type: "string", Required: true, + Placeholder: "postgres|sqlite|clickhouse", Description: "迁移目标:postgres(主库为 PG 时)、sqlite(主库为 SQLite 时)或 clickhouse"}, + }, +} + +type logDBSwitchPayload struct { + Target string `json:"target"` +} + +// LogDBSwitchHandler 切换日志数据库任务处理器。 +type LogDBSwitchHandler struct{} + +// ValidatePayload 校验并规范化参数。 +func (h *LogDBSwitchHandler) ValidatePayload(payload []byte) ([]byte, error) { + var p logDBSwitchPayload + if err := json.Unmarshal(payload, &p); err != nil { + return nil, fmt.Errorf("参数解析失败: %w", err) + } + p.Target = normalizeTarget(p.Target) + if !validTarget(p.Target) { + return nil, fmt.Errorf("目标日志库不合法: %s", p.Target) + } + out, err := json.Marshal(p) + if err != nil { + return nil, err + } + return out, nil +} + +func normalizeTarget(v string) string { + switch v { + case targetPostgres, "postgresql": + return targetPostgres + case targetSQLite, "sqlite3": + return targetSQLite + case targetClickHouse, "ch": + return targetClickHouse + } + return v +} + +func validTarget(v string) bool { + return v == targetPostgres || v == targetSQLite || v == targetClickHouse +} + +// Execute 执行迁移。 +func (h *LogDBSwitchHandler) Execute(ctx context.Context, payload []byte) (*task.TaskResult, error) { + var p logDBSwitchPayload + if err := json.Unmarshal(payload, &p); err != nil { + return nil, fmt.Errorf("参数解析失败: %w", err) + } + p.Target = normalizeTarget(p.Target) + if err := validateSwitch(ctx, p.Target); err != nil { + return nil, err + } + + source, err := currentLogDatabase(ctx) + if err != nil { + task.AppendLog(ctx, "读取日志主库失败: %v", err) + return nil, err + } + task.AppendLog(ctx, "开始切换日志数据库:%s -> %s", source, p.Target) + + if err := setMigrationFlag(ctx, "migrating"); err != nil { + return nil, err + } + defer func() { + if err := setMigrationFlag(ctx, ""); err != nil { + logger.ErrorF(ctx, "清除日志迁移冻结标记失败: %v", err) + } + }() + + if err := risk_control.Drain(ctx); err != nil { + return nil, fmt.Errorf("排空日志写入队列失败: %w", err) + } + + src, err := logstore.Active(ctx) + if err != nil { + return nil, err + } + dst, err := logstore.BuildForMigration(ctx, p.Target) + if err != nil { + return nil, err + } + + if _, err := dst.UserAccessLogs.DeleteAll(ctx); err != nil { + return nil, fmt.Errorf("清空目标用户访问日志失败: %w", err) + } + from, to, err := src.UserAccessLogs.MigrationRange(ctx) + if err != nil { + return nil, fmt.Errorf("读取源库时间范围失败: %w", err) + } + if !from.IsZero() && !to.IsZero() { + if err := dst.UserAccessLogs.EnsurePartitions(ctx, from, to); err != nil { + return nil, fmt.Errorf("预建目标分区失败: %w", err) + } + } + + if err := copyUserAccessLogs(ctx, src, dst); err != nil { + return nil, err + } + if err := flipLogDatabase(ctx, p.Target); err != nil { + return nil, err + } + logstore.InvalidateCache() + task.AppendLog(ctx, "日志数据库已切换为 %s,写入恢复", p.Target) + return &task.TaskResult{Message: fmt.Sprintf("日志数据库已从 %s 切换为 %s", source, p.Target)}, nil +} + +func validateSwitch(ctx context.Context, target string) error { + source, err := currentLogDatabase(ctx) + if err != nil { + return err + } + if source == target { + return errors.New("目标日志库与当前日志库相同,无需迁移") + } + switch target { + case targetClickHouse: + if !config.Config.ClickHouse.Enabled { + return errors.New("ClickHouse 未启用,无法迁移到 ClickHouse") + } + case targetPostgres: + if !config.Config.Database.Enabled { + return errors.New("PostgreSQL 未启用(当前主库为 SQLite),无法迁移到 PostgreSQL") + } + case targetSQLite: + if config.Config.Database.Enabled { + return errors.New("当前主库为 PostgreSQL,日志库不能设置为 SQLite") + } + } + return nil +} + +func currentLogDatabase(ctx context.Context) (string, error) { + cfg, err := repository.GetSystemConfigByKey(ctx, model.ConfigKeyLogDatabase) + if err != nil { + return "", fmt.Errorf("读取日志主库失败: %w", err) + } + if cfg.Value == "" { + return "", errors.New("日志主库配置为空") + } + return cfg.Value, nil +} + +func setMigrationFlag(ctx context.Context, v string) error { + return repository.SaveOrUpdateSystemConfig(ctx, model.ConfigKeyLogDBMigration, v) +} + +func flipLogDatabase(ctx context.Context, target string) error { + return repository.SaveOrUpdateSystemConfig(ctx, model.ConfigKeyLogDatabase, target) +} + +func copyUserAccessLogs(ctx context.Context, src, dst *logstore.Store) error { + var afterID uint64 + var copied int + for { + rows, err := src.UserAccessLogs.ListForMigration(ctx, afterID, copyBatchSize) + if err != nil { + return fmt.Errorf("读取源用户访问日志失败: %w", err) + } + if len(rows) == 0 { + break + } + if err := dst.UserAccessLogs.BatchInsert(ctx, rows); err != nil { + return fmt.Errorf("写入目标用户访问日志失败: %w", err) + } + afterID = rows[len(rows)-1].ID + copied += len(rows) + task.AppendLog(ctx, "已复制用户访问日志 %d 条", copied) + if len(rows) < copyBatchSize { + break + } + } + return nil +} diff --git a/internal/apps/admin/status/log_database.go b/internal/apps/admin/status/log_database.go new file mode 100644 index 00000000..230ca476 --- /dev/null +++ b/internal/apps/admin/status/log_database.go @@ -0,0 +1,102 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package status + +import ( + "context" + "errors" + "net/http" + + "github.com/Rain-kl/Wavelet/internal/infra/config" + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/internal/repository" + "github.com/Rain-kl/Wavelet/internal/repository/logstore" + "github.com/Rain-kl/Wavelet/internal/shared/response" + "github.com/Rain-kl/Wavelet/pkg/logger" + "github.com/gin-gonic/gin" + "gorm.io/gorm" +) + +const ( + logDBNamePostgres = "postgres" + logDBNameSQLite = "sqlite" + logDBNameClickHouse = "clickhouse" + defaultLogRetentionDays = 30 +) + +// LogDatabaseStatus 日志库状态。 +type LogDatabaseStatus struct { + ActiveDatabase string `json:"active_database"` + Migration string `json:"migration"` + RetentionDays map[string]int `json:"retention_days"` + AvailableTargets []string `json:"available_targets"` +} + +// GetLogDatabaseStatus 返回当前日志库状态。 +// @Summary 获取日志数据库状态 +// @Description 返回当前日志主库、迁移状态、各库保留天数与合法迁移目标,需要管理员权限 +// @Tags admin +// @Produce json +// @Security SessionCookie +// @Success 200 {object} response.Any{data=status.LogDatabaseStatus} "获取成功" +// @Failure 401 {object} response.Any "未登录" +// @Failure 403 {object} response.Any "无管理员权限" +// @Failure 500 {object} response.Any "内部错误" +// @Router /api/v1/admin/status/log-database [get] +func GetLogDatabaseStatus(c *gin.Context) { + ctx := c.Request.Context() + store, err := logstore.Active(ctx) + if err != nil { + logger.ErrorF(ctx, "获取日志存储实例失败: %v", err) + response.AbortInternal(c, "日志存储初始化失败") + return + } + activeDB, err := store.Status.ActiveDatabase(ctx) + if err != nil { + logger.ErrorF(ctx, "获取日志库状态失败: %v", err) + response.AbortInternal(c, "获取日志库状态失败") + return + } + migration := "idle" + if logstore.Migrating(ctx) { + migration = "migrating" + } + c.JSON(http.StatusOK, response.OK(LogDatabaseStatus{ + ActiveDatabase: activeDB, + Migration: migration, + RetentionDays: map[string]int{ + logDBNamePostgres: retentionOr(ctx, model.ConfigKeyLogRetentionDaysPostgres), + logDBNameSQLite: retentionOr(ctx, model.ConfigKeyLogRetentionDaysSQLite), + logDBNameClickHouse: retentionOr(ctx, model.ConfigKeyLogRetentionDaysClickHouse), + }, + AvailableTargets: availableLogTargets(activeDB), + })) +} + +func retentionOr(ctx context.Context, key string) int { + v, err := repository.GetIntByKey(ctx, key) + if err != nil { + if !errors.Is(err, gorm.ErrRecordNotFound) { + logger.ErrorF(ctx, "读取日志保留天数配置失败 key=%s: %v", key, err) + } + return defaultLogRetentionDays + } + if v < 1 { + return defaultLogRetentionDays + } + return v +} + +func availableLogTargets(active string) []string { + if active == logDBNameClickHouse { + if config.Config.Database.Enabled { + return []string{logDBNamePostgres} + } + return []string{logDBNameSQLite} + } + if config.Config.ClickHouse.Enabled { + return []string{logDBNameClickHouse} + } + return []string{} +} diff --git a/internal/apps/admin/system_config/errs.go b/internal/apps/admin/system_config/errs.go index 0ec1df84..312e4372 100644 --- a/internal/apps/admin/system_config/errs.go +++ b/internal/apps/admin/system_config/errs.go @@ -11,5 +11,6 @@ const ( ConfigKeyRequired = "配置键不能为空" ConfigValueRequired = "配置值不能为空" ConfigKeyExists = "配置键已存在" + protectedConfigKeyMessage = "该配置项由系统任务管理,禁止手动修改" StorageDriverSwitchRequiresMigration = "存在存量文件,请通过存储迁移任务切换存储引擎" ) diff --git a/internal/apps/admin/system_config/logics.go b/internal/apps/admin/system_config/logics.go index 6e03ae74..bab02077 100644 --- a/internal/apps/admin/system_config/logics.go +++ b/internal/apps/admin/system_config/logics.go @@ -17,7 +17,14 @@ import ( "gorm.io/gorm" ) +func isProtectedConfigKey(key string) bool { + return key == model.ConfigKeyLogDatabase || key == model.ConfigKeyLogDBMigration +} + func createSystemConfig(ctx context.Context, req CreateSystemConfigRequest) error { + if isProtectedConfigKey(req.Key) { + return errors.New(protectedConfigKeyMessage) + } exists, err := repository.SystemConfigExists(ctx, req.Key) if err != nil { return err @@ -53,6 +60,9 @@ func getSystemConfig(ctx context.Context, key string) (model.SystemConfig, error } func updateSystemConfig(ctx context.Context, key string, req UpdateSystemConfigRequest) error { + if isProtectedConfigKey(key) { + return errors.New(protectedConfigKeyMessage) + } config, err := repository.GetAdminSystemConfigByKey(ctx, key) if err != nil { return err diff --git a/internal/apps/admin/system_config/routers.go b/internal/apps/admin/system_config/routers.go index d1a73218..2be447d1 100644 --- a/internal/apps/admin/system_config/routers.go +++ b/internal/apps/admin/system_config/routers.go @@ -64,6 +64,10 @@ func CreateSystemConfig(c *gin.Context) { response.AbortBadRequest(c, err.Error()) return } + if isProtectedConfigKey(req.Key) { + response.AbortBadRequest(c, protectedConfigKeyMessage) + return + } if err := createSystemConfig(c.Request.Context(), req); err != nil { if err.Error() == ConfigKeyExists { @@ -156,6 +160,10 @@ func UpdateSystemConfig(c *gin.Context) { } key := c.Param("key") + if isProtectedConfigKey(key) { + response.AbortBadRequest(c, protectedConfigKeyMessage) + return + } if err := updateSystemConfig(c.Request.Context(), key, req); err != nil { if errors.Is(err, gorm.ErrRecordNotFound) { response.AbortNotFound(c, SystemConfigNotFound) diff --git a/internal/apps/admin/system_config/routers_test.go b/internal/apps/admin/system_config/routers_test.go index a32be2de..befef2b5 100644 --- a/internal/apps/admin/system_config/routers_test.go +++ b/internal/apps/admin/system_config/routers_test.go @@ -27,7 +27,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/shared/response" ) -const expectedDefaultConfigsCount = 30 +const expectedDefaultConfigsCount = 35 func setupTestRouter(authUser *model.User) *gin.Engine { r := testhelper.NewTestGinEngine() @@ -168,8 +168,15 @@ func TestListSystemConfigs(t *testing.T) { var configs []model.SystemConfig _ = json.Unmarshal(dataBytes, &configs) - if len(configs) != 1 || configs[0].Key != model.ConfigKeyMaxAPIKeysPerUser { - t.Errorf("expected 1 business config (max_api_keys_per_user), got %d: %v", len(configs), configs) + if len(configs) != 4 { + t.Errorf("expected 4 business configs, got %d: %v", len(configs), configs) + } + keys := make(map[string]struct{}, len(configs)) + for _, cfg := range configs { + keys[cfg.Key] = struct{}{} + } + if _, ok := keys[model.ConfigKeyMaxAPIKeysPerUser]; !ok { + t.Errorf("missing business config %s", model.ConfigKeyMaxAPIKeysPerUser) } }) } diff --git a/internal/apps/risk_control/logics.go b/internal/apps/risk_control/logics.go index 04ca2d30..e69aac6e 100644 --- a/internal/apps/risk_control/logics.go +++ b/internal/apps/risk_control/logics.go @@ -7,11 +7,12 @@ import ( "context" "sync" - "github.com/Rain-kl/Wavelet/internal/infra/config" + "time" + "github.com/Rain-kl/Wavelet/internal/infra/persistence/batchwriter" "github.com/Rain-kl/Wavelet/internal/model/analytics" "github.com/Rain-kl/Wavelet/internal/platform/lifecycle" - analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics" + "github.com/Rain-kl/Wavelet/internal/repository/logstore" "github.com/Rain-kl/Wavelet/pkg/logger" ) @@ -20,12 +21,8 @@ var ( logWriter *batchwriter.Writer[*analytics.UserAccessLog] ) -// InitLogWriter initializes the ClickHouse access-log batch writer. +// InitLogWriter initializes the access-log batch writer for the active log database. func InitLogWriter(ctx context.Context) { - if !config.Config.ClickHouse.Enabled { - return - } - logWriterMu.Lock() defer logWriterMu.Unlock() if logWriter != nil { @@ -41,7 +38,11 @@ func InitLogWriter(ctx context.Context) { } rows = append(rows, *item) } - return analyticsrepo.BatchInsert(ctx, rows) + store, err := logstore.Active(ctx) + if err != nil { + return err + } + return store.UserAccessLogs.BatchInsert(ctx, rows) }, batchwriter.WithDropHandler[*analytics.UserAccessLog](func(item *analytics.UserAccessLog) { path := "" @@ -51,7 +52,7 @@ func InitLogWriter(ctx context.Context) { logger.WarnF(context.Background(), "[RiskControl] Log queue full, dropping log item for path: %s", path) }), batchwriter.WithFlushErrorHandler[*analytics.UserAccessLog](func(ctx context.Context, items []*analytics.UserAccessLog, err error) { - logger.ErrorF(ctx, "[RiskControl] Send ClickHouse batch failed (batch=%d): %v", len(items), err) + logger.ErrorF(ctx, "[RiskControl] flush access-log batch failed (batch=%d): %v", len(items), err) }), ) if err != nil { @@ -109,3 +110,36 @@ func currentLogWriter() *batchwriter.Writer[*analytics.UserAccessLog] { defer logWriterMu.RUnlock() return logWriter } + +const drainPollInterval = 50 * time.Millisecond + +// Drain waits until the in-memory access-log queue has been empty for one flush interval. +func Drain(ctx context.Context) error { + writer := currentLogWriter() + if writer == nil { + return nil + } + quietPeriod := batchwriter.DefaultConfig().FlushInterval + if quietPeriod <= 0 { + quietPeriod = time.Second + } + ticker := time.NewTicker(drainPollInterval) + defer ticker.Stop() + var quietSince time.Time + for { + if writer.Len() == 0 { + if quietSince.IsZero() { + quietSince = time.Now() + } else if time.Since(quietSince) >= quietPeriod { + return nil + } + } else { + quietSince = time.Time{} + } + select { + case <-ctx.Done(): + return ctx.Err() + case <-ticker.C: + } + } +} diff --git a/internal/apps/upload/task/cleanup.go b/internal/apps/upload/task/cleanup.go index 3dc4d9a1..3e96351c 100644 --- a/internal/apps/upload/task/cleanup.go +++ b/internal/apps/upload/task/cleanup.go @@ -19,6 +19,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/infra/task" "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/repository" + "github.com/Rain-kl/Wavelet/internal/repository/logstore" "github.com/Rain-kl/Wavelet/pkg/logger" "gorm.io/gorm" ) @@ -136,11 +137,23 @@ func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*task.Tas ) } - msg := fmt.Sprintf("系统清理完成。成功清理未使用的上传文件 %d/%d 个;清理历史推送审计日志 %d 条;清理任务执行日志 %d 条。", + var logDeleted int64 + logSummary, logErr := logstore.CleanupExpired(ctx) + if logErr != nil { + task.AppendLog(ctx, "清理过期用户访问日志失败: %v", logErr) + logger.ErrorF(ctx, "清理过期用户访问日志失败: %v", logErr) + } else { + logDeleted = logSummary.Deleted + task.AppendLog(ctx, "成功清理过期用户访问日志 %d 条(%s 保留 %d 天)", + logSummary.Deleted, logSummary.ActiveDatabase, logSummary.RetentionDays) + } + + msg := fmt.Sprintf("系统清理完成。成功清理未使用的上传文件 %d/%d 个;清理历史推送审计日志 %d 条;清理任务执行日志 %d 条;清理过期访问日志 %d 条。", totalDeleted, totalProcessed, pushHistoryCount, taskLogStats.HighFrequencyDeleted+taskLogStats.LowFrequencyDeleted, + logDeleted, ) task.AppendLog(ctx, "%s", msg) return &task.TaskResult{Message: msg}, nil diff --git a/internal/apps/upload/task/tasks_test.go b/internal/apps/upload/task/tasks_test.go index 6cae6b51..2999d261 100644 --- a/internal/apps/upload/task/tasks_test.go +++ b/internal/apps/upload/task/tasks_test.go @@ -134,7 +134,7 @@ func TestSystemCleanupHandler_Execute(t *testing.T) { // 验证结果 require.NoError(t, err) require.NotNil(t, result) - assert.Contains(t, result.Message, "系统清理完成。成功清理未使用的上传文件 2/2 个;清理历史推送审计日志 1 条;清理任务执行日志 1 条。") + assert.Contains(t, result.Message, "系统清理完成。成功清理未使用的上传文件 2/2 个;清理历史推送审计日志 1 条;清理任务执行日志 1 条;清理过期访问日志 0 条。") // 验证数据库状态:pending 且超过1小时的应被标记为 deleted var pendingCount int64 @@ -189,7 +189,7 @@ func TestSystemCleanupHandler_ExecuteNoFiles(t *testing.T) { require.NoError(t, err) require.NotNil(t, result) - assert.Contains(t, result.Message, "系统清理完成。成功清理未使用的上传文件 0/0 个;清理历史推送审计日志 0 条;清理任务执行日志 0 条。") + assert.Contains(t, result.Message, "系统清理完成。成功清理未使用的上传文件 0/0 个;清理历史推送审计日志 0 条;清理任务执行日志 0 条;清理过期访问日志 0 条。") } func TestSystemCleanupHandler_ImplementsTaskHandler(t *testing.T) { diff --git a/internal/infra/persistence/migrator/goose/postgres/202608160001_create_user_access_logs.sql b/internal/infra/persistence/migrator/goose/postgres/202608160001_create_user_access_logs.sql new file mode 100644 index 00000000..f3ff2532 --- /dev/null +++ b/internal/infra/persistence/migrator/goose/postgres/202608160001_create_user_access_logs.sql @@ -0,0 +1,20 @@ +-- +goose Up +CREATE TABLE IF NOT EXISTS w_user_access_logs ( + id BIGINT NOT NULL, + user_id BIGINT NOT NULL DEFAULT 0, + path VARCHAR(2048) NOT NULL DEFAULT '', + method VARCHAR(16) NOT NULL DEFAULT '', + ip VARCHAR(128) NOT NULL DEFAULT '', + user_agent TEXT NOT NULL DEFAULT '', + headers TEXT NOT NULL DEFAULT '', + status INTEGER NOT NULL DEFAULT 0, + latency BIGINT NOT NULL DEFAULT 0, + created_at TIMESTAMPTZ NOT NULL DEFAULT CURRENT_TIMESTAMP, + PRIMARY KEY (id, created_at) +) PARTITION BY RANGE (created_at); + +CREATE INDEX IF NOT EXISTS idx_w_user_access_logs_user_id ON w_user_access_logs (user_id, created_at DESC); +CREATE INDEX IF NOT EXISTS idx_w_user_access_logs_created_at ON w_user_access_logs (created_at DESC); + +-- +goose Down +DROP TABLE IF EXISTS w_user_access_logs; diff --git a/internal/infra/persistence/migrator/goose/postgres/202608160002_log_database_configs.sql b/internal/infra/persistence/migrator/goose/postgres/202608160002_log_database_configs.sql new file mode 100644 index 00000000..81643145 --- /dev/null +++ b/internal/infra/persistence/migrator/goose/postgres/202608160002_log_database_configs.sql @@ -0,0 +1,24 @@ +-- +goose Up +INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at) +VALUES + ('log_database', '', 'system', 0, '当前日志主库(postgres/sqlite/clickhouse),由切换任务写入', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('log_db_migration', '', 'system', 0, '日志库迁移冻结标记(空或 migrating)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) +ON CONFLICT (key) DO NOTHING; + +INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at) +SELECT k, '30', 'business', 0, descr, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP +FROM (VALUES + ('log_retention_days_postgres', 'PostgreSQL 用户访问日志保留天数'), + ('log_retention_days_sqlite', 'SQLite 用户访问日志保留天数'), + ('log_retention_days_clickhouse', 'ClickHouse 用户访问日志保留天数') +) AS v(k, descr) +ON CONFLICT (key) DO NOTHING; + +-- +goose Down +DELETE FROM w_system_configs WHERE key IN ( + 'log_database', + 'log_db_migration', + 'log_retention_days_postgres', + 'log_retention_days_sqlite', + 'log_retention_days_clickhouse' +); diff --git a/internal/infra/persistence/migrator/goose/sqlite/202608160001_create_user_access_logs.sql b/internal/infra/persistence/migrator/goose/sqlite/202608160001_create_user_access_logs.sql new file mode 100644 index 00000000..e8abea95 --- /dev/null +++ b/internal/infra/persistence/migrator/goose/sqlite/202608160001_create_user_access_logs.sql @@ -0,0 +1,22 @@ +-- +goose Up +CREATE TABLE IF NOT EXISTS w_user_access_logs ( + id INTEGER NOT NULL, + user_id INTEGER NOT NULL DEFAULT 0, + path TEXT NOT NULL DEFAULT '', + method TEXT NOT NULL DEFAULT '', + ip TEXT NOT NULL DEFAULT '', + user_agent TEXT NOT NULL DEFAULT '', + headers TEXT NOT NULL DEFAULT '', + status INTEGER NOT NULL DEFAULT 0, + latency INTEGER NOT NULL DEFAULT 0, + created_at DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, + PRIMARY KEY (id, created_at) +); + +CREATE INDEX IF NOT EXISTS idx_w_user_access_logs_user_id ON w_user_access_logs (user_id, created_at DESC); +CREATE INDEX IF NOT EXISTS idx_w_user_access_logs_created_at ON w_user_access_logs (created_at DESC); + +-- +goose Down +DROP TABLE IF EXISTS w_user_access_logs; +DROP INDEX IF EXISTS idx_w_user_access_logs_user_id; +DROP INDEX IF EXISTS idx_w_user_access_logs_created_at; diff --git a/internal/infra/persistence/migrator/goose/sqlite/202608160002_log_database_configs.sql b/internal/infra/persistence/migrator/goose/sqlite/202608160002_log_database_configs.sql new file mode 100644 index 00000000..0a6918c6 --- /dev/null +++ b/internal/infra/persistence/migrator/goose/sqlite/202608160002_log_database_configs.sql @@ -0,0 +1,17 @@ +-- +goose Up +INSERT OR IGNORE INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at) +VALUES + ('log_database', '', 'system', 0, '当前日志主库(postgres/sqlite/clickhouse),由切换任务写入', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('log_db_migration', '', 'system', 0, '日志库迁移冻结标记(空或 migrating)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('log_retention_days_postgres', '30', 'business', 0, 'PostgreSQL 用户访问日志保留天数', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('log_retention_days_sqlite', '30', 'business', 0, 'SQLite 用户访问日志保留天数', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), + ('log_retention_days_clickhouse', '30', 'business', 0, 'ClickHouse 用户访问日志保留天数', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP); + +-- +goose Down +DELETE FROM w_system_configs WHERE key IN ( + 'log_database', + 'log_db_migration', + 'log_retention_days_postgres', + 'log_retention_days_sqlite', + 'log_retention_days_clickhouse' +); diff --git a/internal/infra/task/handlers/register.go b/internal/infra/task/handlers/register.go index 907ad17c..ddb8f73f 100644 --- a/internal/infra/task/handlers/register.go +++ b/internal/infra/task/handlers/register.go @@ -6,6 +6,7 @@ package handlers import ( + "github.com/Rain-kl/Wavelet/internal/apps/admin/logs" "github.com/Rain-kl/Wavelet/internal/apps/admin/push" "github.com/Rain-kl/Wavelet/internal/apps/upload" "github.com/Rain-kl/Wavelet/internal/apps/user" @@ -35,4 +36,8 @@ func Register() { // push task.RegisterHandler(push.SendNotificationTask, &push.PushHandler{}) task.RegisterTaskMeta(push.SendNotificationMeta) + + // logs + task.RegisterHandler(logs.LogDBSwitchTask, &logs.LogDBSwitchHandler{}) + task.RegisterTaskMeta(logs.LogDBSwitchMeta) } diff --git a/internal/model/system_configs.go b/internal/model/system_configs.go index 1ba0e5b9..86493926 100644 --- a/internal/model/system_configs.go +++ b/internal/model/system_configs.go @@ -37,6 +37,11 @@ const ( ConfigKeyLoginSessionTTLHours = "login_session_ttl_hours" // 登录会话过期时间 (小时,0表示浏览器关闭后自动退出登录,-1表示永不过期) ConfigKeyUpdateUpstreamRepository = "update_upstream_repository" // GitHub Actions Release 上游仓库 ConfigKeyStorageConfig = "storage_config" // 文件存储配置 (JSON) + ConfigKeyLogDatabase = "log_database" // 当前日志主库(postgres/sqlite/clickhouse),受保护 + ConfigKeyLogDBMigration = "log_db_migration" // 日志库迁移冻结标记(空/migrating),受保护 + ConfigKeyLogRetentionDaysPostgres = "log_retention_days_postgres" // PostgreSQL 用户访问日志保留天数 + ConfigKeyLogRetentionDaysSQLite = "log_retention_days_sqlite" // SQLite 用户访问日志保留天数 + ConfigKeyLogRetentionDaysClickHouse = "log_retention_days_clickhouse" // ClickHouse 用户访问日志保留天数 ) const ( diff --git a/internal/platform/bootstrap/bootstrap.go b/internal/platform/bootstrap/bootstrap.go index 537b64c9..d005cb26 100644 --- a/internal/platform/bootstrap/bootstrap.go +++ b/internal/platform/bootstrap/bootstrap.go @@ -7,16 +7,23 @@ package bootstrap import ( "context" + "errors" + "fmt" + "log" "sync" admin_push "github.com/Rain-kl/Wavelet/internal/apps/admin/push" "github.com/Rain-kl/Wavelet/internal/apps/admin/push/custom_events" "github.com/Rain-kl/Wavelet/internal/apps/risk_control" + "github.com/Rain-kl/Wavelet/internal/infra/config" taskhandlers "github.com/Rain-kl/Wavelet/internal/infra/task/handlers" + "github.com/Rain-kl/Wavelet/internal/model" "github.com/Rain-kl/Wavelet/internal/platform/lifecycle" "github.com/Rain-kl/Wavelet/internal/repository" + "github.com/Rain-kl/Wavelet/internal/repository/logstore" "github.com/Rain-kl/Wavelet/pkg/cache/ram" "github.com/Rain-kl/Wavelet/pkg/logger" + "gorm.io/gorm" ) // Options selects role-specific runtime bootstrap steps for the current process. @@ -127,6 +134,20 @@ func RegisterAll() { // Call from cmd entry points after wiring registration and database migration, not from router. func Init(ctx context.Context, opts Options) { initRuntimeOnce.Do(func() { + if err := validateAndSeedLogDatabase(ctx); err != nil { + logger.ErrorF(ctx, "[Bootstrap] 日志主库配置校验失败: %v", err) + log.Fatalf("[Bootstrap] 日志主库配置校验失败: %v", err) + } + + logstore.SetConfigReader(func(ctx context.Context, key string) (string, error) { + cfg, err := repository.GetSystemConfigByKey(ctx, key) + if err != nil { + return "", err + } + return cfg.Value, nil + }) + logstore.Init(ctx) + // Register config cache loader RegisterCache(repository.ConfigCacheType, CacheRegistry{ Loader: repository.ConfigLoader{}, @@ -146,6 +167,44 @@ func Init(ctx context.Context, opts Options) { }) } +func validateAndSeedLogDatabase(ctx context.Context) error { + cfg, err := repository.GetSystemConfigByKey(ctx, model.ConfigKeyLogDatabase) + if err != nil && !errors.Is(err, gorm.ErrRecordNotFound) { + return fmt.Errorf("读取日志主库配置失败: %w", err) + } + current := cfg.Value + if current == "" { + current = "sqlite" + if config.Config.Database.Enabled { + current = "postgres" + } + if config.Config.ClickHouse.Enabled { + current = "clickhouse" + } + if err := repository.SaveOrUpdateSystemConfig(ctx, model.ConfigKeyLogDatabase, current); err != nil { + return fmt.Errorf("初始化日志主库配置失败: %w", err) + } + return nil + } + switch current { + case "clickhouse": + if !config.Config.ClickHouse.Enabled { + return errors.New("当前日志主库为 ClickHouse 但 ClickHouse 未启用。请先重新启用 ClickHouse 配置并启动,在任务管理运行「切换日志数据库」迁移到 PostgreSQL/SQLite 后再禁用 ClickHouse") + } + case "postgres": + if !config.Config.Database.Enabled { + return errors.New("当前日志主库为 PostgreSQL 但 PostgreSQL 未启用(当前为 SQLite 主库)。请运行「切换日志数据库」迁回 SQLite 或启用 PostgreSQL") + } + case "sqlite": + if config.Config.Database.Enabled { + return errors.New("当前日志主库为 SQLite 但当前主库为 PostgreSQL。请运行「切换日志数据库」迁移到 PostgreSQL") + } + default: + return fmt.Errorf("未知的日志主库配置: %s", current) + } + return nil +} + // Stop stops all batch writers and background resources. func Stop(ctx context.Context) { lifecycle.Stop(ctx) diff --git a/internal/repository/analytics/access_log.go b/internal/repository/analytics/access_log.go index 3a5c30f5..e7b1cb53 100644 --- a/internal/repository/analytics/access_log.go +++ b/internal/repository/analytics/access_log.go @@ -8,6 +8,8 @@ import ( "context" "fmt" + "time" + "github.com/Rain-kl/Wavelet/internal/infra/persistence" analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" "gorm.io/gorm" @@ -69,6 +71,28 @@ func ListAccessLogs(ctx context.Context, filter AccessLogFilter, page, pageSize return logs, safeUint64Count(total), nil } +// DeleteAllUserAccessLogs hard-deletes all user access logs via TRUNCATE. +func DeleteAllUserAccessLogs(ctx context.Context) (int64, error) { + if db.ChConn == nil { + return 0, fmt.Errorf("clickhouse connection is not initialized") + } + if err := db.ChConn.Exec(ctx, "TRUNCATE TABLE "+analyticsmodel.UserAccessLog{}.TableName()); err != nil { + return 0, fmt.Errorf("truncate user access logs: %w", err) + } + return 0, nil +} + +// DeleteUserAccessLogsBefore deletes user access logs older than cutoff. +func DeleteUserAccessLogsBefore(ctx context.Context, cutoff time.Time) (int64, error) { + if db.ChConn == nil { + return 0, fmt.Errorf("clickhouse connection is not initialized") + } + if err := db.ChConn.Exec(ctx, "ALTER TABLE "+analyticsmodel.UserAccessLog{}.TableName()+" DELETE WHERE created_at < ?", cutoff); err != nil { + return 0, fmt.Errorf("delete expired user access logs: %w", err) + } + return 0, nil +} + func safeUint64Count(count int64) uint64 { if count < 0 { return 0 diff --git a/internal/repository/logstore/cleanup.go b/internal/repository/logstore/cleanup.go new file mode 100644 index 00000000..78439f69 --- /dev/null +++ b/internal/repository/logstore/cleanup.go @@ -0,0 +1,93 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package logstore + +import ( + "context" + "errors" + "fmt" + "strconv" + "time" + + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/pkg/logger" +) + +const ( + defaultLogRetentionDays = 30 + partitionLeadMonths = 2 + userAccessLogTable = "w_user_access_logs" +) + +// CleanupSummary 汇总本次清理结果。 +type CleanupSummary struct { + ActiveDatabase string `json:"active_database"` + RetentionDays int `json:"retention_days"` + Deleted int64 `json:"deleted"` +} + +// CleanupExpired 按当前日志库保留天数删除过期用户访问日志,并预建 PG 分区。 +func CleanupExpired(ctx context.Context) (CleanupSummary, error) { + active, err := ActiveDatabase(ctx) + if err != nil { + return CleanupSummary{}, err + } + days := retentionDaysForDatabase(ctx, active) + summary := CleanupSummary{ActiveDatabase: active, RetentionDays: days} + + store, err := Active(ctx) + if err != nil { + return summary, err + } + now := time.Now().UTC() + if err := store.UserAccessLogs.EnsurePartitions(ctx, now, now.AddDate(0, partitionLeadMonths, 0)); err != nil { + logger.WarnF(ctx, "logstore: ensure partitions during cleanup failed: %v", err) + } + cutoff := now.AddDate(0, 0, -days) + deleted, err := store.UserAccessLogs.DeleteBefore(ctx, cutoff) + if err != nil { + return summary, fmt.Errorf("delete expired user access logs: %w", err) + } + summary.Deleted = deleted + return summary, nil +} + +func retentionDaysForDatabase(ctx context.Context, dbName string) int { + key := model.ConfigKeyLogRetentionDaysPostgres + switch dbName { + case dbNameSQLite: + key = model.ConfigKeyLogRetentionDaysSQLite + case dbNameClickHouse: + key = model.ConfigKeyLogRetentionDaysClickHouse + } + v, err := getConfig(ctx, key) + if err != nil { + if !errors.Is(err, errConfigReaderNotWired) { + logger.ErrorF(ctx, "读取日志保留天数配置失败(key=%s),回退默认 %d 天: %v", key, defaultLogRetentionDays, err) + } + return defaultLogRetentionDays + } + days, perr := strconv.Atoi(v) + if perr != nil || days <= 0 { + logger.ErrorF(ctx, "日志保留天数配置非法(key=%s, value=%q),回退默认 %d 天", key, v, defaultLogRetentionDays) + return defaultLogRetentionDays + } + return days +} + +func partitionStatementsRange(from, to time.Time) []string { + var out []string + start := time.Date(from.Year(), from.Month(), 1, 0, 0, 0, 0, time.UTC) + end := time.Date(to.Year(), to.Month(), 1, 0, 0, 0, 0, time.UTC).AddDate(0, 1, 0) + for ; start.Before(end); start = start.AddDate(0, 1, 0) { + monthEnd := start.AddDate(0, 1, 0) + suffix := start.Format("200601") + fromDay := start.Format("2006-01-02") + toDay := monthEnd.Format("2006-01-02") + out = append(out, fmt.Sprintf( + "CREATE TABLE IF NOT EXISTS %s_%s PARTITION OF %s FOR VALUES FROM ('%s') TO ('%s')", + userAccessLogTable, suffix, userAccessLogTable, fromDay, toDay)) + } + return out +} diff --git a/internal/repository/logstore/clickhouse.go b/internal/repository/logstore/clickhouse.go new file mode 100644 index 00000000..1e151cc5 --- /dev/null +++ b/internal/repository/logstore/clickhouse.go @@ -0,0 +1,146 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package logstore + +import ( + "context" + "fmt" + "time" + + "github.com/ClickHouse/clickhouse-go/v2/lib/driver" + db "github.com/Rain-kl/Wavelet/internal/infra/persistence" + analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" + analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics" +) + +type clickhouseUserAccessLogStore struct { + skipFreeze bool +} + +func newClickHouseUserAccessLogStore() *clickhouseUserAccessLogStore { + return &clickhouseUserAccessLogStore{} +} + +var ( + _ UserAccessLogStore = (*clickhouseUserAccessLogStore)(nil) + _ StatusStore = (*clickhouseUserAccessLogStore)(nil) +) + +func (s *clickhouseUserAccessLogStore) ActiveDatabase(_ context.Context) (string, error) { + return dbNameClickHouse, nil +} + +func (s *clickhouseUserAccessLogStore) ensureWritable(ctx context.Context) error { + if !s.skipFreeze && Migrating(ctx) { + return ErrMigrating + } + return nil +} + +func (s *clickhouseUserAccessLogStore) BatchInsert(ctx context.Context, logs []analyticsmodel.UserAccessLog) error { + if len(logs) == 0 { + return nil + } + if err := s.ensureWritable(ctx); err != nil { + return err + } + return analyticsrepo.BatchInsert(ctx, logs) +} + +func (s *clickhouseUserAccessLogStore) DeleteAll(ctx context.Context) (int64, error) { + if err := s.ensureWritable(ctx); err != nil { + return 0, err + } + return analyticsrepo.DeleteAllUserAccessLogs(ctx) +} + +func (s *clickhouseUserAccessLogStore) DeleteBefore(ctx context.Context, cutoff time.Time) (int64, error) { + if err := s.ensureWritable(ctx); err != nil { + return 0, err + } + return analyticsrepo.DeleteUserAccessLogsBefore(ctx, cutoff) +} + +func (s *clickhouseUserAccessLogStore) Count(ctx context.Context, filter analyticsrepo.AccessLogFilter) (uint64, error) { + return analyticsrepo.CountAccessLogs(ctx, filter) +} + +func (s *clickhouseUserAccessLogStore) List(ctx context.Context, filter analyticsrepo.AccessLogFilter, page, pageSize int) ([]analyticsmodel.UserAccessLog, uint64, error) { + return analyticsrepo.ListAccessLogs(ctx, filter, page, pageSize) +} + +func (s *clickhouseUserAccessLogStore) GetDailyTrend(ctx context.Context, days int) ([]analyticsrepo.DailyTrend, error) { + return analyticsrepo.GetDailyTrend(ctx, days) +} + +func (s *clickhouseUserAccessLogStore) GetBrowserDistribution(ctx context.Context, startTime time.Time) ([]analyticsrepo.BrowserShare, error) { + return analyticsrepo.GetBrowserDistribution(ctx, startTime) +} + +func (s *clickhouseUserAccessLogStore) GetTopActiveUsers(ctx context.Context, startTime time.Time, limit int) ([]analyticsrepo.TopUser, error) { + return analyticsrepo.GetTopActiveUsers(ctx, startTime, limit) +} + +func (s *clickhouseUserAccessLogStore) EnsurePartitions(_ context.Context, _, _ time.Time) error { + return nil +} + +func (s *clickhouseUserAccessLogStore) MigrationRange(ctx context.Context) (time.Time, time.Time, error) { + if db.ChConn == nil { + return time.Time{}, time.Time{}, fmt.Errorf("clickhouse connection is not initialized") + } + table := analyticsmodel.UserAccessLog{}.TableName() + var minTime, maxTime *time.Time + if err := db.ChConn.QueryRow(ctx, "SELECT min(created_at), max(created_at) FROM "+table).Scan(&minTime, &maxTime); err != nil { + return time.Time{}, time.Time{}, fmt.Errorf("query migration range %s: %w", table, err) + } + if minTime == nil || maxTime == nil { + return time.Time{}, time.Time{}, nil + } + return minTime.UTC(), maxTime.UTC(), nil +} + +func (s *clickhouseUserAccessLogStore) ListForMigration(ctx context.Context, afterID uint64, limit int) ([]analyticsmodel.UserAccessLog, error) { + if db.ChConn == nil { + return nil, fmt.Errorf("clickhouse connection is not initialized") + } + if limit <= 0 { + limit = migrationPageSize + } + table := analyticsmodel.UserAccessLog{}.TableName() + columns := analyticsmodel.UserAccessLog{}.InsertColumns() + rows, err := db.ChConn.Query(ctx, fmt.Sprintf( + "SELECT %s FROM %s WHERE id > ? ORDER BY id ASC LIMIT ?", + columns, table, + ), afterID, limit) + if err != nil { + return nil, fmt.Errorf("list user access logs for migration: %w", err) + } + defer func() { _ = rows.Close() }() + return scanUserAccessLogs(rows) +} + +func scanUserAccessLogs(rows driver.Rows) ([]analyticsmodel.UserAccessLog, error) { + var result []analyticsmodel.UserAccessLog + for rows.Next() { + var item analyticsmodel.UserAccessLog + if err := rows.Scan( + &item.ID, + &item.UserID, + &item.Path, + &item.Method, + &item.IP, + &item.UserAgent, + &item.Headers, + &item.Status, + &item.Latency, + &item.CreatedAt, + ); err != nil { + return nil, fmt.Errorf("scan user access log row: %w", err) + } + item.CreatedAt = item.CreatedAt.UTC() + result = append(result, item) + } + return result, nil +} diff --git a/internal/repository/logstore/gorm.go b/internal/repository/logstore/gorm.go new file mode 100644 index 00000000..e00dbe18 --- /dev/null +++ b/internal/repository/logstore/gorm.go @@ -0,0 +1,327 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package logstore + +import ( + "context" + "errors" + "fmt" + "sort" + "strings" + "time" + + "github.com/Rain-kl/Wavelet/internal/infra/persistence/idgen" + analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" + analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics" + "gorm.io/gorm" +) + +const ( + insertBatchSize = 500 + migrationPageSize = 100 + defaultPageSize = 20 + defaultTopN = 10 + topUserAgents = 100 + dayDuration = 24 * time.Hour +) + +type gormLogStore struct { + db *gorm.DB + skipFreeze bool +} + +func newGormStore(db *gorm.DB) *gormLogStore { return &gormLogStore{db: db} } + +type userAccessLogGormStore struct { + *gormLogStore +} + +func newUserAccessLogGormStore(db *gorm.DB) *userAccessLogGormStore { + return &userAccessLogGormStore{gormLogStore: newGormStore(db)} +} + +var ( + _ UserAccessLogStore = (*userAccessLogGormStore)(nil) + _ StatusStore = (*userAccessLogGormStore)(nil) +) + +func (s *gormLogStore) ActiveDatabase(_ context.Context) (string, error) { + if isPostgresDialect(s.db) { + return dbNamePostgres, nil + } + return dbNameSQLite, nil +} + +func (s *gormLogStore) ensureWritable(ctx context.Context) error { + if !s.skipFreeze && Migrating(ctx) { + return ErrMigrating + } + return nil +} + +func (s *userAccessLogGormStore) BatchInsert(ctx context.Context, logs []analyticsmodel.UserAccessLog) error { + if len(logs) == 0 { + return nil + } + if err := s.ensureWritable(ctx); err != nil { + return err + } + for i := range logs { + if logs[i].ID == 0 { + logs[i].ID = idgen.NextUint64ID() + } + } + return s.db.WithContext(ctx).CreateInBatches(logs, insertBatchSize).Error +} + +func (s *userAccessLogGormStore) DeleteAll(ctx context.Context) (int64, error) { + if err := s.ensureWritable(ctx); err != nil { + return 0, err + } + res := s.db.WithContext(ctx).Where("1 = 1").Delete(&analyticsmodel.UserAccessLog{}) + return res.RowsAffected, res.Error +} + +func (s *userAccessLogGormStore) DeleteBefore(ctx context.Context, cutoff time.Time) (int64, error) { + if err := s.ensureWritable(ctx); err != nil { + return 0, err + } + res := s.db.WithContext(ctx).Where("created_at < ?", cutoff).Delete(&analyticsmodel.UserAccessLog{}) + if res.Error != nil && isMissingRelation(res.Error) { + return 0, nil + } + return res.RowsAffected, res.Error +} + +func (s *userAccessLogGormStore) ListForMigration(ctx context.Context, afterID uint64, limit int) ([]analyticsmodel.UserAccessLog, error) { + var rows []analyticsmodel.UserAccessLog + q := s.db.WithContext(ctx).Model(&analyticsmodel.UserAccessLog{}). + Where("id > ?", afterID). + Order("id ASC"). + Limit(limitOr(limit, migrationPageSize)) + if err := q.Find(&rows).Error; err != nil { + return nil, err + } + return rows, nil +} + +func (s *userAccessLogGormStore) MigrationRange(ctx context.Context) (time.Time, time.Time, error) { + return gormMigrationRange(ctx, s.db, "created_at", analyticsmodel.UserAccessLog{}, func(v *analyticsmodel.UserAccessLog) time.Time { + return v.CreatedAt + }) +} + +func (s *userAccessLogGormStore) Count(ctx context.Context, filter analyticsrepo.AccessLogFilter) (uint64, error) { + where, args, ok := buildUserAccessLogWhere(filter) + if !ok { + return 0, nil + } + var total int64 + if err := s.db.WithContext(ctx).Model(&analyticsmodel.UserAccessLog{}).Where(where, args...).Count(&total).Error; err != nil { + return 0, err + } + return countToUint64(total), nil +} + +func (s *userAccessLogGormStore) List(ctx context.Context, filter analyticsrepo.AccessLogFilter, page, pageSize int) ([]analyticsmodel.UserAccessLog, uint64, error) { + where, args, ok := buildUserAccessLogWhere(filter) + if !ok { + return []analyticsmodel.UserAccessLog{}, 0, nil + } + var total int64 + if err := s.db.WithContext(ctx).Model(&analyticsmodel.UserAccessLog{}).Where(where, args...).Count(&total).Error; err != nil { + return nil, 0, err + } + if total == 0 { + return []analyticsmodel.UserAccessLog{}, 0, nil + } + var rows []analyticsmodel.UserAccessLog + q := s.db.WithContext(ctx).Where(where, args...).Order("created_at DESC, id DESC") + if err := q.Limit(limitOr(pageSize, defaultPageSize)).Offset(offsetOf(page, pageSize)).Find(&rows).Error; err != nil { + return nil, 0, err + } + return rows, countToUint64(total), nil +} + +func buildUserAccessLogWhere(filter analyticsrepo.AccessLogFilter) (string, []any, bool) { + if filter.UserIDs != nil && len(filter.UserIDs) == 0 { + return "", nil, false + } + var parts []string + var args []any + if filter.UserIDs != nil { + parts = append(parts, "user_id IN ?") + args = append(args, filter.UserIDs) + } + 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 = 1", args, true + } + return strings.Join(parts, " AND "), args, true +} + +func (s *userAccessLogGormStore) GetDailyTrend(ctx context.Context, days int) ([]analyticsrepo.DailyTrend, error) { + if days <= 0 { + days = 7 + } + start := time.Now().AddDate(0, 0, -(days - 1)).Truncate(dayDuration) + type row struct { + Date string + Cnt uint64 + } + var rows []row + err := s.db.WithContext(ctx).Model(&analyticsmodel.UserAccessLog{}). + Select(dailyTrendDateSQL(s.db)+" AS date, COUNT(*) AS cnt"). + Where("created_at >= ?", start). + Group("date").Order("date ASC").Scan(&rows).Error + if err != nil { + return nil, err + } + counts := make(map[string]uint64, len(rows)) + for _, r := range rows { + counts[r.Date] = r.Cnt + } + out := make([]analyticsrepo.DailyTrend, 0, days) + for i := 0; i < days; i++ { + d := start.AddDate(0, 0, i).Format("2006-01-02") + out = append(out, analyticsrepo.DailyTrend{Date: d, Count: counts[d]}) + } + return out, nil +} + +func (s *userAccessLogGormStore) GetBrowserDistribution(ctx context.Context, startTime time.Time) ([]analyticsrepo.BrowserShare, error) { + type row struct { + UserAgent string + Cnt uint64 + } + var rows []row + err := s.db.WithContext(ctx).Model(&analyticsmodel.UserAccessLog{}). + Select("user_agent, COUNT(*) AS cnt"). + Where("created_at >= ?", startTime). + Group("user_agent").Order("cnt DESC").Limit(topUserAgents).Scan(&rows).Error + if err != nil { + return nil, err + } + counts := make(map[string]uint64) + for _, r := range rows { + counts[analyticsrepo.ParseBrowserName(r.UserAgent)] += r.Cnt + } + out := make([]analyticsrepo.BrowserShare, 0, len(counts)) + for label, count := range counts { + out = append(out, analyticsrepo.BrowserShare{Browser: label, Count: count}) + } + sort.Slice(out, func(i, j int) bool { return out[i].Count > out[j].Count }) + return out, nil +} + +func (s *userAccessLogGormStore) GetTopActiveUsers(ctx context.Context, startTime time.Time, limit int) ([]analyticsrepo.TopUser, error) { + type row struct { + UserID uint64 + Cnt uint64 + } + var rows []row + err := s.db.WithContext(ctx).Model(&analyticsmodel.UserAccessLog{}). + Select("user_id, COUNT(*) AS cnt"). + Where("user_id <> 0 AND created_at >= ?", startTime). + Group("user_id").Order("cnt DESC").Limit(limitOr(limit, defaultTopN)).Scan(&rows).Error + if err != nil { + return nil, err + } + out := make([]analyticsrepo.TopUser, len(rows)) + for i, r := range rows { + out[i] = analyticsrepo.TopUser{UserID: r.UserID, Count: r.Cnt} + } + return out, nil +} + +func (s *userAccessLogGormStore) EnsurePartitions(ctx context.Context, from, to time.Time) error { + if !isPostgresDialect(s.db) { + return nil + } + for _, sql := range partitionStatementsRange(from, to) { + if err := s.db.WithContext(ctx).Exec(sql).Error; err != nil { + return fmt.Errorf("ensure partition: %w", err) + } + } + return nil +} + +func gormMigrationRange[T any]( + ctx context.Context, + gdb *gorm.DB, + column string, + model T, + timeOf func(*T) time.Time, +) (time.Time, time.Time, error) { + var first, last T + found := false + for _, order := range []string{"ASC", "DESC"} { + out := &first + if order == "DESC" { + out = &last + } + res := gdb.WithContext(ctx).Model(model).Order(column + " " + order).Limit(1).Take(out) + if res.Error != nil && !errors.Is(res.Error, gorm.ErrRecordNotFound) { + return time.Time{}, time.Time{}, fmt.Errorf("query migration range %s: %w", column, res.Error) + } + if res.Error == nil { + found = true + } + } + if !found { + return time.Time{}, time.Time{}, nil + } + return timeOf(&first).UTC(), timeOf(&last).UTC(), nil +} + +func limitOr(v, def int) int { + if v <= 0 { + return def + } + return v +} + +func offsetOf(page, pageSize int) int { + if page < 1 { + page = 1 + } + return (page - 1) * limitOr(pageSize, defaultPageSize) +} + +func countToUint64(v int64) uint64 { + if v < 0 { + return 0 + } + return uint64(v) +} + +func isPostgresDialect(db *gorm.DB) bool { + return db != nil && db.Dialector != nil && db.Name() == "postgres" +} + +func dailyTrendDateSQL(db *gorm.DB) string { + if isPostgresDialect(db) { + return "to_char(created_at, 'YYYY-MM-DD')" + } + return "strftime('%Y-%m-%d', created_at)" +} + +func isMissingRelation(err error) bool { + if err == nil { + return false + } + msg := strings.ToLower(err.Error()) + return strings.Contains(msg, "no such table") || strings.Contains(msg, "does not exist") +} diff --git a/internal/repository/logstore/gorm_test.go b/internal/repository/logstore/gorm_test.go new file mode 100644 index 00000000..38807ddc --- /dev/null +++ b/internal/repository/logstore/gorm_test.go @@ -0,0 +1,60 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package logstore + +import ( + "context" + "testing" + "time" + + analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" + analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics" + "github.com/glebarez/sqlite" + "github.com/stretchr/testify/require" + "gorm.io/gorm" +) + +func newTestUserAccessStore(t *testing.T) *userAccessLogGormStore { + t.Helper() + gdb, err := gorm.Open(sqlite.Open("file:logstore-"+t.Name()+"?mode=memory&cache=shared"), &gorm.Config{}) + require.NoError(t, err) + require.NoError(t, gdb.AutoMigrate(&analyticsmodel.UserAccessLog{})) + return newUserAccessLogGormStore(gdb) +} + +func TestGormUserAccessLogCountList(t *testing.T) { + ua := newTestUserAccessStore(t) + ctx := context.Background() + now := time.Now().UTC().Truncate(time.Second) + require.NoError(t, ua.BatchInsert(ctx, []analyticsmodel.UserAccessLog{ + {UserID: 10, Path: "/api/v1/users", Method: "GET", Status: 200, CreatedAt: now}, + {UserID: 20, Path: "/api/v1/admin", Method: "GET", Status: 200, CreatedAt: now}, + {UserID: 10, Path: "/api/v1/other", Method: "POST", Status: 201, CreatedAt: now}, + })) + + count, err := ua.Count(ctx, analyticsrepo.AccessLogFilter{UserIDs: []uint64{10}, Path: "users"}) + require.NoError(t, err) + require.Equal(t, uint64(1), count) + + rows, total, err := ua.List(ctx, analyticsrepo.AccessLogFilter{UserIDs: []uint64{10}, Path: "users"}, 1, 10) + require.NoError(t, err) + require.Equal(t, uint64(1), total) + require.Len(t, rows, 1) + require.Equal(t, "/api/v1/users", rows[0].Path) + require.NotZero(t, rows[0].ID) +} + +func TestGormUserAccessLogFreeze(t *testing.T) { + ua := newTestUserAccessStore(t) + SetConfigReader(func(_ context.Context, key string) (string, error) { + if key == logMigrationKey { + return "migrating", nil + } + return "", nil + }) + t.Cleanup(ResetForTest) + + err := ua.BatchInsert(context.Background(), []analyticsmodel.UserAccessLog{{UserID: 1, CreatedAt: time.Now()}}) + require.ErrorIs(t, err, ErrMigrating) +} diff --git a/internal/repository/logstore/logstore.go b/internal/repository/logstore/logstore.go new file mode 100644 index 00000000..ef583010 --- /dev/null +++ b/internal/repository/logstore/logstore.go @@ -0,0 +1,43 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +// Package logstore abstracts user access-log storage across ClickHouse, PostgreSQL and SQLite. +package logstore + +import ( + "context" + "errors" + "time" + + analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" + analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics" +) + +// ErrMigrating 表示日志数据库正在迁移,当前禁止写入。 +var ErrMigrating = errors.New("log database is migrating, writes are disabled") + +// UserAccessLogStore 用户访问日志(w_user_access_logs)。 +type UserAccessLogStore interface { + BatchInsert(ctx context.Context, logs []analyticsmodel.UserAccessLog) error + DeleteAll(ctx context.Context) (int64, error) + DeleteBefore(ctx context.Context, cutoff time.Time) (int64, error) + Count(ctx context.Context, filter analyticsrepo.AccessLogFilter) (uint64, error) + List(ctx context.Context, filter analyticsrepo.AccessLogFilter, page, pageSize int) ([]analyticsmodel.UserAccessLog, uint64, error) + GetDailyTrend(ctx context.Context, days int) ([]analyticsrepo.DailyTrend, error) + GetBrowserDistribution(ctx context.Context, startTime time.Time) ([]analyticsrepo.BrowserShare, error) + GetTopActiveUsers(ctx context.Context, startTime time.Time, limit int) ([]analyticsrepo.TopUser, error) + ListForMigration(ctx context.Context, afterID uint64, limit int) ([]analyticsmodel.UserAccessLog, error) + MigrationRange(ctx context.Context) (from, to time.Time, err error) + EnsurePartitions(ctx context.Context, from, to time.Time) error +} + +// StatusStore 日志库状态。 +type StatusStore interface { + ActiveDatabase(ctx context.Context) (string, error) +} + +// Store 当前生效日志库。 +type Store struct { + UserAccessLogs UserAccessLogStore + Status StatusStore +} diff --git a/internal/repository/logstore/provider.go b/internal/repository/logstore/provider.go new file mode 100644 index 00000000..2fa5b808 --- /dev/null +++ b/internal/repository/logstore/provider.go @@ -0,0 +1,189 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package logstore + +import ( + "context" + "errors" + "fmt" + "sync" + "time" + + "github.com/Rain-kl/Wavelet/internal/infra/config" + db "github.com/Rain-kl/Wavelet/internal/infra/persistence" + "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/pkg/logger" +) + +const ( + logDatabaseKey = model.ConfigKeyLogDatabase + logMigrationKey = model.ConfigKeyLogDBMigration +) + +const ( + dbNamePostgres = "postgres" + dbNameSQLite = "sqlite" + dbNameClickHouse = "clickhouse" +) + +var errConfigReaderNotWired = errors.New("logstore: config reader not wired") + +// ConfigReader 读取系统配置字符串值,由 bootstrap 注入(避免 logstore ↔ repository 循环依赖)。 +type ConfigReader func(ctx context.Context, key string) (string, error) + +const resolveCacheTTL = 1 * time.Second + +var ( + configReader ConfigReader + + storeMu sync.RWMutex + active *Store + activeDB string + lastResolveDB string + lastResolveTime time.Time +) + +// SetConfigReader 注入系统配置读取函数(bootstrap 调用,测试可注入内存实现)。 +func SetConfigReader(fn ConfigReader) { configReader = fn } + +func getConfig(ctx context.Context, key string) (string, error) { + if configReader == nil { + return "", errConfigReaderNotWired + } + return configReader(ctx, key) +} + +// Active 返回当前生效的日志库 Store。 +func Active(ctx context.Context) (*Store, error) { + current, err := resolveDatabase(ctx) + if err != nil { + return nil, err + } + storeMu.RLock() + if active != nil && activeDB == current { + s := active + storeMu.RUnlock() + return s, nil + } + storeMu.RUnlock() + + storeMu.Lock() + defer storeMu.Unlock() + if active != nil && activeDB == current { + return active, nil + } + s, err := buildStore(ctx, current, false) + if err != nil { + return nil, err + } + active = s + activeDB = current + return s, nil +} + +// Build 直接按目标构造 store(不经 Active 缓存)。 +func Build(ctx context.Context, database string) (*Store, error) { + return buildStore(ctx, database, false) +} + +// BuildForMigration 构造迁移目标 store,跳过冻结检查。 +func BuildForMigration(ctx context.Context, database string) (*Store, error) { + return buildStore(ctx, database, true) +} + +func buildStore(ctx context.Context, database string, skipFreeze bool) (*Store, error) { + switch database { + case dbNameClickHouse: + ual := newClickHouseUserAccessLogStore() + ual.skipFreeze = skipFreeze + return &Store{UserAccessLogs: ual, Status: ual}, nil + case dbNamePostgres, dbNameSQLite: + gdb := db.DB(ctx) + ual := newUserAccessLogGormStore(gdb) + ual.skipFreeze = skipFreeze + return &Store{UserAccessLogs: ual, Status: ual}, nil + default: + return nil, fmt.Errorf("unsupported log database: %s", database) + } +} + +// Migrating 返回日志库是否处于迁移冻结状态。 +func Migrating(ctx context.Context) bool { + v, err := getConfig(ctx, logMigrationKey) + if err != nil { + if !errors.Is(err, errConfigReaderNotWired) { + logger.ErrorF(ctx, "read log migration config failed: %v", err) + } + return false + } + return v == "migrating" +} + +// Init 预热激活 store,并兜底预建当前月及未来分区。 +func Init(ctx context.Context) { + s, err := Active(ctx) + if err != nil { + return + } + now := time.Now().UTC() + if err := s.UserAccessLogs.EnsurePartitions(ctx, now, now.AddDate(0, partitionLeadMonths, 0)); err != nil { + logger.WarnF(ctx, "logstore: ensure startup partitions failed: %v", err) + } +} + +// InvalidateCache 清空日志库解析缓存。 +func InvalidateCache() { + storeMu.Lock() + defer storeMu.Unlock() + lastResolveTime = time.Time{} + lastResolveDB = "" +} + +// ResetForTest 清空缓存的激活 store 与 config reader。 +func ResetForTest() { + storeMu.Lock() + active = nil + activeDB = "" + lastResolveDB = "" + lastResolveTime = time.Time{} + storeMu.Unlock() + configReader = nil +} + +// ActiveDatabase 返回当前日志主库名。 +func ActiveDatabase(ctx context.Context) (string, error) { + return resolveDatabase(ctx) +} + +func resolveDatabase(ctx context.Context) (string, error) { + storeMu.RLock() + if active != nil && time.Since(lastResolveTime) < resolveCacheTTL { + name := lastResolveDB + storeMu.RUnlock() + return name, nil + } + storeMu.RUnlock() + + v, err := getConfig(ctx, logDatabaseKey) + if err != nil && !errors.Is(err, errConfigReaderNotWired) { + return "", err + } + + resolved := v + if resolved == "" { + resolved = dbNameSQLite + if config.Config.Database.Enabled { + resolved = dbNamePostgres + } + if config.Config.ClickHouse.Enabled { + resolved = dbNameClickHouse + } + } + + storeMu.Lock() + lastResolveDB = resolved + lastResolveTime = time.Now() + storeMu.Unlock() + return resolved, nil +} diff --git a/internal/router/v1/admin.go b/internal/router/v1/admin.go index 819d42d1..043cb162 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/log-database", admin_status.GetLogDatabaseStatus) // Database basic info & backup export adminRouter.GET("/db-info", admin_status.GetDatabaseInfo) diff --git a/internal/testhelper/test_helper.go b/internal/testhelper/test_helper.go index f5b29ac9..77204610 100644 --- a/internal/testhelper/test_helper.go +++ b/internal/testhelper/test_helper.go @@ -282,6 +282,36 @@ func getSeedConfigsPart2() []model.SystemConfig { Type: configTypeSystem, Description: "文件存储驱动及连接配置(JSON)", }, + { + Key: model.ConfigKeyLogDatabase, + Value: "sqlite", + Type: configTypeSystem, + Description: "当前日志主库", + }, + { + Key: model.ConfigKeyLogDBMigration, + Value: "", + Type: configTypeSystem, + Description: "日志库迁移冻结标记", + }, + { + Key: model.ConfigKeyLogRetentionDaysPostgres, + Value: "30", + Type: "business", + Description: "PostgreSQL 用户访问日志保留天数", + }, + { + Key: model.ConfigKeyLogRetentionDaysSQLite, + Value: "30", + Type: "business", + Description: "SQLite 用户访问日志保留天数", + }, + { + Key: model.ConfigKeyLogRetentionDaysClickHouse, + Value: "30", + Type: "business", + Description: "ClickHouse 用户访问日志保留天数", + }, } }