diff --git a/docs/docs.go b/docs/docs.go index 60b37f10..c32de5a9 100644 --- a/docs/docs.go +++ b/docs/docs.go @@ -1092,7 +1092,7 @@ const docTemplate = `{ } }, "400": { - "description": "ClickHouse 未启用或参数错误", + "description": "日志存储未启用或参数错误", "schema": { "$ref": "#/definitions/response.Any" } diff --git a/docs/superpowers/plans/2026-08-08-log-database-decoupling.md b/docs/superpowers/plans/2026-08-08-log-database-decoupling.md index 2d1e7117..aa511d75 100644 --- a/docs/superpowers/plans/2026-08-08-log-database-decoupling.md +++ b/docs/superpowers/plans/2026-08-08-log-database-decoupling.md @@ -1532,30 +1532,37 @@ var allowedInfraPersistence = []string{ } func TestAppsMustNotImportLogBackendDirectly(t *testing.T) { - out, err := exec.Command("go", "list", "-deps", "./internal/apps/...").Output() + t.Chdir("../../..") + out, err := exec.Command("go", "list", "-test", "-f", `{{.ImportPath}} {{join .Imports " "}}`, "./internal/apps/...").Output() if err != nil { t.Fatalf("go list: %v", err) } for _, line := range strings.Split(string(out), "\n") { - pkg := strings.TrimSpace(line) - if pkg == "" { + fields := strings.Fields(line) + if len(fields) == 0 { continue } - for _, forbidden := range forbiddenImports { - if pkg == forbidden { - t.Errorf("internal/apps must not import %s", forbidden) - } + pkg := fields[0] + if !strings.HasPrefix(pkg, "github.com/Rain-kl/Wavelet/internal/apps") { + continue } - if strings.HasPrefix(pkg, "github.com/Rain-kl/Wavelet/internal/infra/persistence/") { - allowed := false - for _, a := range allowedInfraPersistence { - if pkg == a || strings.HasPrefix(pkg, a+"/") { - allowed = true - break + for _, imp := range fields[1:] { + for _, forbidden := range forbiddenImports { + if imp == forbidden && !allowedAnalyticsDelegation[pkg] { + t.Errorf("%s must not import forbidden log backend %s", pkg, forbidden) } } - if !allowed { - t.Errorf("internal/apps must not import infra/persistence subpackage: %s", pkg) + if strings.HasPrefix(imp, "github.com/Rain-kl/Wavelet/internal/infra/persistence/") { + allowed := false + for _, a := range allowedInfraPersistence { + if imp == a || strings.HasPrefix(imp, a+"/") { + allowed = true + break + } + } + if !allowed { + t.Errorf("%s must not import infra/persistence subpackage directly: %s", pkg, imp) + } } } } @@ -2261,6 +2268,9 @@ func (h *LogDBSwitchHandler) Execute(ctx context.Context, payload []byte) (*task if err := copyAccessLogs(ctx, src, dst); err != nil { return nil, err } + if err := copyUserAccessLogs(ctx, src, dst); err != nil { + return nil, err + } if err := copyObservability(ctx, src, dst); err != nil { return nil, err } @@ -2329,10 +2339,13 @@ func buildTargetStore(ctx context.Context, database string) (*logstore.Store, er } func clearTargetLogTables(ctx context.Context, dst *logstore.Store, target string) error { - // 依次清空 6 张表:AccessLogs.DeleteAll、Observability.DeleteAll*(SQLite/PG 用 DeleteAll;CH 用 TRUNCATE 语义)。 + // 依次清空 6 张表:AccessLogs.DeleteAll、UserAccessLogs.DeleteAll、Observability.DeleteAll*(SQLite/PG 用 DeleteAll;CH 用 TRUNCATE 语义)。 if _, err := dst.AccessLogs.DeleteAll(ctx); err != nil { return fmt.Errorf("清空目标访问日志失败: %w", err) } + if _, err := dst.UserAccessLogs.DeleteAll(ctx); err != nil { + return fmt.Errorf("清空目标用户访问日志失败: %w", err) + } for _, fn := range []func(context.Context) (int64, error){ dst.Observability.DeleteAllMetricSnapshots, dst.Observability.DeleteAllEdgeHealth, diff --git a/docs/swagger.json b/docs/swagger.json index 6d9fa916..a6370e1d 100644 --- a/docs/swagger.json +++ b/docs/swagger.json @@ -1085,7 +1085,7 @@ } }, "400": { - "description": "ClickHouse 未启用或参数错误", + "description": "日志存储未启用或参数错误", "schema": { "$ref": "#/definitions/response.Any" } diff --git a/docs/swagger.yaml b/docs/swagger.yaml index d602f080..f025db8a 100644 --- a/docs/swagger.yaml +++ b/docs/swagger.yaml @@ -4970,7 +4970,7 @@ paths: $ref: '#/definitions/logs.accessLogsResponse' type: object "400": - description: ClickHouse 未启用或参数错误 + description: 日志存储未启用或参数错误 schema: $ref: '#/definitions/response.Any' "401": diff --git a/frontend/app/(main)/admin/tasks/components/task-manager.tsx b/frontend/app/(main)/admin/tasks/components/task-manager.tsx index 787578a5..82c97754 100644 --- a/frontend/app/(main)/admin/tasks/components/task-manager.tsx +++ b/frontend/app/(main)/admin/tasks/components/task-manager.tsx @@ -588,7 +588,11 @@ export function TaskManager() { } disabled={dispatching} > - + diff --git a/frontend/app/(main)/error-pages/page.tsx b/frontend/app/(main)/error-pages/page.tsx index b2f85e57..6612d45a 100644 --- a/frontend/app/(main)/error-pages/page.tsx +++ b/frontend/app/(main)/error-pages/page.tsx @@ -152,8 +152,7 @@ export default function ErrorPagesPage() {

- 开启后仅对 GET - 请求的匹配错误状态码返回自定义错误页;POST/PUT + 开启后仅对 GET 请求的匹配错误状态码返回自定义错误页;POST/PUT 等其它方法直接透传源站响应。

diff --git a/frontend/components/common/settings/operation-tab.tsx b/frontend/components/common/settings/operation-tab.tsx index c95d051d..d42ffe57 100644 --- a/frontend/components/common/settings/operation-tab.tsx +++ b/frontend/components/common/settings/operation-tab.tsx @@ -99,12 +99,17 @@ export function OperationTab({ const updateRetentionMutation = useMutation({ mutationFn: async (values: Record) => { for (const field of LOG_RETENTION_FIELDS) { + const raw = (values[field.key] ?? '').trim(); + const num = Number(raw); + if (!raw || !Number.isInteger(num) || num < 1) { + throw new Error(`${field.label}必须为大于等于 1 的整数`); + } const config = businessConfigs[field.key]; if (!config) { throw new Error(`缺少配置项: ${field.key}`); } await services.adminSystemConfig.updateSystemConfig(field.key, { - value: values[field.key] ?? '90', + value: String(num), description: config.description, }); } diff --git a/internal/apps/admin/logs/routers.go b/internal/apps/admin/logs/routers.go index b91ce94e..e8082bc3 100644 --- a/internal/apps/admin/logs/routers.go +++ b/internal/apps/admin/logs/routers.go @@ -237,7 +237,7 @@ 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 "无管理员权限" // @Router /api/v1/admin/logs/access [get] @@ -245,6 +245,7 @@ func GetAccessLogs(c *gin.Context) { ctx := c.Request.Context() store, err := logstore.Active(ctx) if err != nil { + logger.ErrorF(ctx, "获取日志存储实例失败: %v", err) response.AbortWithError(c, http.StatusBadRequest, "日志存储未启用,无法检索访问日志") return } @@ -346,6 +347,7 @@ func GetLogsAnalytics(c *gin.Context) { ctx := c.Request.Context() store, err := logstore.Active(ctx) if err != nil { + logger.ErrorF(ctx, "获取日志存储实例失败: %v", err) response.AbortWithError(c, http.StatusBadRequest, "日志存储未启用,无法获取分析数据") return } diff --git a/internal/apps/admin/status/clickhouse.go b/internal/apps/admin/status/clickhouse.go index 93c24305..0b0f59e4 100644 --- a/internal/apps/admin/status/clickhouse.go +++ b/internal/apps/admin/status/clickhouse.go @@ -5,6 +5,7 @@ package status import ( "context" + "errors" "net/http" "github.com/Rain-kl/Wavelet/internal/apps/openflare/chwriter" @@ -18,6 +19,7 @@ import ( "github.com/Rain-kl/Wavelet/internal/shared/response" "github.com/Rain-kl/Wavelet/pkg/logger" "github.com/gin-gonic/gin" + "gorm.io/gorm" ) // 日志库名取值(与 model 配置值、logstore provider 分支保持一致)。 @@ -54,6 +56,7 @@ 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 } @@ -61,6 +64,7 @@ func GetLogDatabaseStatus(c *gin.Context) { // 分支判定复用同一 store 实例的 ActiveDatabase,避免与 Active 解析之间出现 TOCTOU。 activeDB, err := store.Status.ActiveDatabase(ctx) if err != nil { + logger.ErrorF(ctx, "获取日志库状态失败: %v", err) response.AbortInternal(c, "获取日志库状态失败") return } @@ -97,6 +101,9 @@ func GetLogDatabaseStatus(c *gin.Context) { 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 } return v diff --git a/internal/apps/admin/status/clickhouse_test.go b/internal/apps/admin/status/clickhouse_test.go index 91404055..a9ea9598 100644 --- a/internal/apps/admin/status/clickhouse_test.go +++ b/internal/apps/admin/status/clickhouse_test.go @@ -72,6 +72,8 @@ func TestAvailableTargets(t *testing.T) { // TestGetLogDatabaseStatusSmoke 覆盖 handler 的 CH 激活分支(无需 DB/CH 连接)。 func TestGetLogDatabaseStatusSmoke(t *testing.T) { restoreConfig(t) + config.Config.Database.Enabled = false + config.Config.ClickHouse.Enabled = true logstore.ResetForTest() t.Cleanup(logstore.ResetForTest) diff --git a/internal/apps/upload/task/cleanup.go b/internal/apps/upload/task/cleanup.go index e874d086..378cf106 100644 --- a/internal/apps/upload/task/cleanup.go +++ b/internal/apps/upload/task/cleanup.go @@ -134,6 +134,7 @@ func (h *SystemCleanupHandler) Execute(ctx context.Context, _ []byte) (*task.Tas summary, err := logstore.CleanupExpired(ctx) switch { case err != nil: + logger.ErrorF(ctx, "清理过期日志失败: %v", err) task.AppendLog(ctx, "清理过期日志失败: %v", err) case summary.Deleted == 0: task.AppendLog(ctx, "没有需要清理的过期日志 (保留 %d 天)", summary.RetentionDays) diff --git a/internal/infra/persistence/migrator/goose/postgres/202608080002_log_retention_configs.sql b/internal/infra/persistence/migrator/goose/postgres/202608080002_log_retention_configs.sql index af5c14b7..5cd7a8c5 100644 --- a/internal/infra/persistence/migrator/goose/postgres/202608080002_log_retention_configs.sql +++ b/internal/infra/persistence/migrator/goose/postgres/202608080002_log_retention_configs.sql @@ -1,10 +1,12 @@ -- +goose Up -- 日志保留天数配置(business),替换旧的 database_auto_cleanup_* 键。 INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at) -VALUES - ('log_retention_days_postgres', '90', 'business', 0, 'PostgreSQL 日志保留天数(访问日志与可观测统一)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), - ('log_retention_days_sqlite', '90', 'business', 0, 'SQLite 日志保留天数', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), - ('log_retention_days_clickhouse','90', 'business', 0, 'ClickHouse 日志保留天数', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) +SELECT k, COALESCE((SELECT value FROM w_system_configs WHERE key = 'database_auto_cleanup_retention_days'), '90'), '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; DELETE FROM w_system_configs WHERE key IN ('database_auto_cleanup_enabled', 'database_auto_cleanup_retention_days'); @@ -13,7 +15,8 @@ DELETE FROM w_system_configs WHERE key IN ('database_auto_cleanup_enabled', 'dat INSERT INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at) VALUES ('database_auto_cleanup_enabled', 'true', 'business', 0, '数据库自动清理开关', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), - ('database_auto_cleanup_retention_days', '30', 'business', 0, '数据库保留天数', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) + ('database_auto_cleanup_retention_days', COALESCE((SELECT value FROM w_system_configs WHERE key = 'log_retention_days_postgres'), '30'), 'business', 0, '数据库保留天数', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) ON CONFLICT (key) DO NOTHING; DELETE FROM w_system_configs WHERE key IN ('log_retention_days_postgres', 'log_retention_days_sqlite', 'log_retention_days_clickhouse'); + diff --git a/internal/infra/persistence/migrator/goose/postgres/202608080003_drop_database_cleanup_schedule.sql b/internal/infra/persistence/migrator/goose/postgres/202608080003_drop_database_cleanup_schedule.sql index 651e604f..2a8dfed6 100644 --- a/internal/infra/persistence/migrator/goose/postgres/202608080003_drop_database_cleanup_schedule.sql +++ b/internal/infra/persistence/migrator/goose/postgres/202608080003_drop_database_cleanup_schedule.sql @@ -1,7 +1,18 @@ -- +goose Up +CREATE TABLE IF NOT EXISTS w_schedules_backup_of_database_auto_cleanup AS +SELECT * FROM w_schedules WHERE task_type = 'of_database_auto_cleanup'; + DELETE FROM w_schedules WHERE task_type = 'of_database_auto_cleanup'; -- +goose Down +INSERT INTO w_schedules (id, name, task_type, cron, payload, is_active, created_at, updated_at) +SELECT id, name, task_type, cron, payload, is_active, created_at, updated_at +FROM w_schedules_backup_of_database_auto_cleanup +ON CONFLICT (id) DO NOTHING; + INSERT INTO w_schedules (id, name, task_type, cron, payload, is_active, created_at, updated_at) VALUES (102, 'OpenFlare 可观测数据自动清理', 'of_database_auto_cleanup', '0 3 * * *', '{}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) ON CONFLICT (id) DO NOTHING; + +DROP TABLE IF EXISTS w_schedules_backup_of_database_auto_cleanup; + diff --git a/internal/infra/persistence/migrator/goose/sqlite/202608080002_log_retention_configs.sql b/internal/infra/persistence/migrator/goose/sqlite/202608080002_log_retention_configs.sql index afb38dfd..91346906 100644 --- a/internal/infra/persistence/migrator/goose/sqlite/202608080002_log_retention_configs.sql +++ b/internal/infra/persistence/migrator/goose/sqlite/202608080002_log_retention_configs.sql @@ -1,17 +1,19 @@ -- +goose Up -- 日志保留天数配置(business),替换旧的 database_auto_cleanup_* 键。 INSERT OR IGNORE INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at) -VALUES - ('log_retention_days_postgres', '90', 'business', 0, 'PostgreSQL 日志保留天数(访问日志与可观测统一)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), - ('log_retention_days_sqlite', '90', 'business', 0, 'SQLite 日志保留天数', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), - ('log_retention_days_clickhouse','90', 'business', 0, 'ClickHouse 日志保留天数', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP); +SELECT 'log_retention_days_postgres', COALESCE((SELECT value FROM w_system_configs WHERE key = 'database_auto_cleanup_retention_days'), '90'), 'business', 0, 'PostgreSQL 日志保留天数(访问日志与可观测统一)', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP +UNION ALL +SELECT 'log_retention_days_sqlite', COALESCE((SELECT value FROM w_system_configs WHERE key = 'database_auto_cleanup_retention_days'), '90'), 'business', 0, 'SQLite 日志保留天数', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP +UNION ALL +SELECT 'log_retention_days_clickhouse', COALESCE((SELECT value FROM w_system_configs WHERE key = 'database_auto_cleanup_retention_days'), '90'), 'business', 0, 'ClickHouse 日志保留天数', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP; DELETE FROM w_system_configs WHERE key IN ('database_auto_cleanup_enabled', 'database_auto_cleanup_retention_days'); -- +goose Down INSERT OR IGNORE INTO w_system_configs (key, value, type, visibility, description, created_at, updated_at) -VALUES - ('database_auto_cleanup_enabled', 'true', 'business', 0, '数据库自动清理开关', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP), - ('database_auto_cleanup_retention_days', '30', 'business', 0, '数据库保留天数', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP); +SELECT 'database_auto_cleanup_enabled', 'true', 'business', 0, '数据库自动清理开关', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP +UNION ALL +SELECT 'database_auto_cleanup_retention_days', COALESCE((SELECT value FROM w_system_configs WHERE key = 'log_retention_days_sqlite'), '30'), 'business', 0, '数据库保留天数', CURRENT_TIMESTAMP, CURRENT_TIMESTAMP; DELETE FROM w_system_configs WHERE key IN ('log_retention_days_postgres', 'log_retention_days_sqlite', 'log_retention_days_clickhouse'); + diff --git a/internal/infra/persistence/migrator/goose/sqlite/202608080003_drop_database_cleanup_schedule.sql b/internal/infra/persistence/migrator/goose/sqlite/202608080003_drop_database_cleanup_schedule.sql index 6f4c3a22..4ad60bbb 100644 --- a/internal/infra/persistence/migrator/goose/sqlite/202608080003_drop_database_cleanup_schedule.sql +++ b/internal/infra/persistence/migrator/goose/sqlite/202608080003_drop_database_cleanup_schedule.sql @@ -1,7 +1,16 @@ -- +goose Up +CREATE TABLE IF NOT EXISTS w_schedules_backup_of_database_auto_cleanup AS +SELECT * FROM w_schedules WHERE task_type = 'of_database_auto_cleanup'; + DELETE FROM w_schedules WHERE task_type = 'of_database_auto_cleanup'; -- +goose Down -INSERT INTO w_schedules (id, name, task_type, cron, payload, is_active, created_at, updated_at) -VALUES (102, 'OpenFlare 可观测数据自动清理', 'of_database_auto_cleanup', '0 3 * * *', '{}', 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP) -ON CONFLICT (id) DO NOTHING; +INSERT OR IGNORE INTO w_schedules (id, name, task_type, cron, payload, is_active, created_at, updated_at) +SELECT id, name, task_type, cron, payload, is_active, created_at, updated_at +FROM w_schedules_backup_of_database_auto_cleanup; + +INSERT OR IGNORE INTO w_schedules (id, name, task_type, cron, payload, is_active, created_at, updated_at) +VALUES (102, 'OpenFlare 可观测数据自动清理', 'of_database_auto_cleanup', '0 3 * * *', '{}', 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP); + +DROP TABLE IF EXISTS w_schedules_backup_of_database_auto_cleanup; + diff --git a/internal/platform/bootstrap/bootstrap_test.go b/internal/platform/bootstrap/bootstrap_test.go index 55bdc404..05c1eb62 100644 --- a/internal/platform/bootstrap/bootstrap_test.go +++ b/internal/platform/bootstrap/bootstrap_test.go @@ -109,6 +109,15 @@ func TestValidateAndSeedLogDatabaseUpdatesEmptyMarker(t *testing.T) { _, _, cleanup := testhelper.SetupTestEnvironment(t) defer cleanup() + dbPrev := config.Config.Database.Enabled + chPrev := config.Config.ClickHouse.Enabled + config.Config.Database.Enabled = true + config.Config.ClickHouse.Enabled = false + t.Cleanup(func() { + config.Config.Database.Enabled = dbPrev + config.Config.ClickHouse.Enabled = chPrev + }) + ctx := context.Background() // 标记行已存在但值为空,等同首次启动,应写入默认值(走更新路径)。 if err := repository.CreateSystemConfig(ctx, &model.SystemConfig{Key: model.ConfigKeyLogDatabase, Value: "", Type: "system"}); err != nil { @@ -125,13 +134,7 @@ func TestValidateAndSeedLogDatabaseUpdatesEmptyMarker(t *testing.T) { if err != nil { t.Fatalf("GetSystemConfigByKey(%s) error = %v", model.ConfigKeyLogDatabase, err) } - want := "sqlite" - if config.Config.Database.Enabled { - want = "postgres" - } - if config.Config.ClickHouse.Enabled { - want = "clickhouse" - } + want := "postgres" if cfg.Value != want { t.Fatalf("log_database = %q, want %q", cfg.Value, want) } diff --git a/internal/repository/logstore/cleanup.go b/internal/repository/logstore/cleanup.go index 71d43e23..3b18fb69 100644 --- a/internal/repository/logstore/cleanup.go +++ b/internal/repository/logstore/cleanup.go @@ -28,12 +28,13 @@ const defaultLogRetentionDays = 90 // partitionLeadMonths 清理时确保「当前月 + 未来 2 个月」分区持续存在。 const partitionLeadMonths = 2 -// retentionDaysForActive 按当前激活库读取保留天数(默认 90)。 -func retentionDaysForActive(ctx context.Context) int { +// retentionDaysForDatabase 按给定日志库读取保留天数(默认 90)。 +func retentionDaysForDatabase(ctx context.Context, dbName string) int { key := model.ConfigKeyLogRetentionDaysPostgres - if dbName, _ := resolveDatabase(ctx); dbName == dbNameSQLite { + switch dbName { + case dbNameSQLite: key = model.ConfigKeyLogRetentionDaysSQLite - } else if dbName == dbNameClickHouse { + case dbNameClickHouse: key = model.ConfigKeyLogRetentionDaysClickHouse } v, err := getConfig(ctx, key) @@ -49,14 +50,17 @@ func retentionDaysForActive(ctx context.Context) int { // CleanupExpired 按当前激活库保留天数清理过期日志(每日由 system_cleanup 调用)。 func CleanupExpired(ctx context.Context) (*CleanupSummary, error) { + dbName, err := resolveDatabase(ctx) + if err != nil { + return nil, fmt.Errorf("resolve active database: %w", err) + } s, err := Active(ctx) if err != nil { return nil, err } - days := retentionDaysForActive(ctx) + days := retentionDaysForDatabase(ctx, dbName) cutoff := time.Now().AddDate(0, 0, -days) - summary := &CleanupSummary{RetentionDays: days, Tables: []string{}} - summary.ActiveDatabase, _ = resolveDatabase(ctx) + summary := &CleanupSummary{ActiveDatabase: dbName, RetentionDays: days, Tables: []string{}} // PG 分区表仅在迁移时预建「当前+2 月」分区,此处确保分区持续存在, // 否则跨月后新写入会报 "no partition of relation found"(SQLite/CH 为 no-op)。 diff --git a/internal/repository/logstore/cleanup_test.go b/internal/repository/logstore/cleanup_test.go index cc2653db..0a2d1796 100644 --- a/internal/repository/logstore/cleanup_test.go +++ b/internal/repository/logstore/cleanup_test.go @@ -140,8 +140,8 @@ func TestCleanupExpiredSQLite(t *testing.T) { } } -// TestRetentionDaysForActive 覆盖保留天数读取:按激活库选 key、非法值回退默认 90。 -func TestRetentionDaysForActive(t *testing.T) { +// TestRetentionDaysForDatabase 覆盖保留天数读取:按激活库选 key、非法值回退默认 90。 +func TestRetentionDaysForDatabase(t *testing.T) { ResetForTest() SetConfigReader(func(_ context.Context, key string) (string, error) { switch key { @@ -152,8 +152,8 @@ func TestRetentionDaysForActive(t *testing.T) { } return "", nil }) - if got := retentionDaysForActive(context.Background()); got != 30 { - t.Fatalf("retentionDaysForActive = %d, want 30", got) + if got := retentionDaysForDatabase(context.Background(), "sqlite"); got != 30 { + t.Fatalf("retentionDaysForDatabase = %d, want 30", got) } // 非法值(非数字/<=0)回退默认 90。 @@ -166,16 +166,16 @@ func TestRetentionDaysForActive(t *testing.T) { } return "", nil }) - if got := retentionDaysForActive(context.Background()); got != 90 { - t.Fatalf("retentionDaysForActive invalid value = %d, want 90", got) + if got := retentionDaysForDatabase(context.Background(), "postgres"); got != 90 { + t.Fatalf("retentionDaysForDatabase invalid value = %d, want 90", got) } // reader 报错回退默认 90。 SetConfigReader(func(_ context.Context, _ string) (string, error) { return "", fmt.Errorf("boom") }) - if got := retentionDaysForActive(context.Background()); got != 90 { - t.Fatalf("retentionDaysForActive reader error = %d, want 90", got) + if got := retentionDaysForDatabase(context.Background(), "postgres"); got != 90 { + t.Fatalf("retentionDaysForDatabase reader error = %d, want 90", got) } } diff --git a/internal/repository/logstore/imports_test.go b/internal/repository/logstore/imports_test.go index 303db298..a0dd3d51 100644 --- a/internal/repository/logstore/imports_test.go +++ b/internal/repository/logstore/imports_test.go @@ -33,8 +33,7 @@ var allowedInfraPersistence = []string{ func TestAppsMustNotImportLogBackendDirectly(t *testing.T) { t.Chdir("../../..") // module root,保证 ./internal/apps/... 可解析 - // -test 同时列出测试二进制依赖,覆盖仅测试文件引入的底层日志实现。 - out, err := exec.Command("go", "list", "-test", "-deps", "-f", `{{.ImportPath}} {{join .Imports " "}}`, "./internal/apps/...").Output() + out, err := exec.Command("go", "list", "-test", "-f", `{{.ImportPath}} {{join .Imports " "}}`, "./internal/apps/...").Output() if err != nil { t.Fatalf("go list: %v", err) } @@ -44,23 +43,26 @@ func TestAppsMustNotImportLogBackendDirectly(t *testing.T) { continue } pkg := fields[0] - for _, forbidden := range forbiddenImports { - for _, imp := range fields[1:] { - if imp == forbidden && !allowedAnalyticsDelegation[pkg] { - t.Errorf("internal/apps must not import %s (via %s)", forbidden, pkg) - } - } + if !strings.HasPrefix(pkg, "github.com/Rain-kl/Wavelet/internal/apps") { + continue } - if strings.HasPrefix(pkg, "github.com/Rain-kl/Wavelet/internal/infra/persistence/") { - allowed := false - for _, a := range allowedInfraPersistence { - if pkg == a || strings.HasPrefix(pkg, a+"/") { - allowed = true - break + for _, imp := range fields[1:] { + for _, forbidden := range forbiddenImports { + if imp == forbidden && !allowedAnalyticsDelegation[pkg] { + t.Errorf("%s must not import forbidden log backend %s", pkg, forbidden) } } - if !allowed { - t.Errorf("internal/apps must not import infra/persistence subpackage: %s", pkg) + if strings.HasPrefix(imp, "github.com/Rain-kl/Wavelet/internal/infra/persistence/") { + allowed := false + for _, a := range allowedInfraPersistence { + if imp == a || strings.HasPrefix(imp, a+"/") { + allowed = true + break + } + } + if !allowed { + t.Errorf("%s must not import infra/persistence subpackage directly: %s", pkg, imp) + } } } } diff --git a/internal/repository/logstore/postgres_store.go b/internal/repository/logstore/postgres_store.go index 86149da8..ca048b77 100644 --- a/internal/repository/logstore/postgres_store.go +++ b/internal/repository/logstore/postgres_store.go @@ -1076,7 +1076,7 @@ FROM ( ` + counterDelta("disk_write_bytes") + ` AS write_delta FROM of_node_metric_snapshots WHERE ` + where + ` -) +) AS deltas GROUP BY hour_epoch ORDER BY hour_epoch ASC` type row struct { diff --git a/internal/repository/logstore/provider.go b/internal/repository/logstore/provider.go index 8bcdd07c..14df5cb7 100644 --- a/internal/repository/logstore/provider.go +++ b/internal/repository/logstore/provider.go @@ -8,10 +8,12 @@ import ( "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" ) // logDatabaseKey / logMigrationKey 对应 model.ConfigKeyLogDatabase / ConfigKeyLogDBMigration。 @@ -33,12 +35,16 @@ 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 + storeMu sync.RWMutex + active *Store + activeDB string + lastResolveDB string + lastResolveTime time.Time ) // SetConfigReader 注入系统配置读取函数(bootstrap 调用,测试可注入内存实现)。 @@ -124,6 +130,9 @@ func buildStore(ctx context.Context, database string, skipFreeze bool) (*Store, 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" @@ -134,11 +143,21 @@ func Init(ctx context.Context) { _, _ = Active(ctx) } +// InvalidateCache 清空日志库解析缓存(在修改 log_database 配置后显式调用)。 +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 } @@ -151,20 +170,35 @@ func ActiveDatabase(ctx context.Context) (string, error) { // resolveDatabase 读取 log_database:值缺失或 reader 未装配(首启)时按启动规则 seed; // 已装配 reader 的真实读取错误直接透出,避免把读失败当首次启动。 func resolveDatabase(ctx context.Context) (string, error) { + storeMu.RLock() + if active != nil && time.Since(lastResolveTime) < resolveCacheTTL { + db := lastResolveDB + storeMu.RUnlock() + return db, nil + } + storeMu.RUnlock() + v, err := getConfig(ctx, logDatabaseKey) if err != nil && !errors.Is(err, errConfigReaderNotWired) { return "", err } - if v != "" { - return v, nil + + resolved := v + if resolved == "" { + // 首次启动 seed:CH 启用 → clickhouse;否则随主库。 + resolved = dbNameSQLite + if config.Config.Database.Enabled { + resolved = dbNamePostgres + } + if config.Config.ClickHouse.Enabled { + resolved = dbNameClickHouse + } } - // 首次启动 seed:CH 启用 → clickhouse;否则随主库。 - defaultDB := dbNameSQLite - if config.Config.Database.Enabled { - defaultDB = dbNamePostgres - } - if config.Config.ClickHouse.Enabled { - defaultDB = dbNameClickHouse - } - return defaultDB, nil + + storeMu.Lock() + lastResolveDB = resolved + lastResolveTime = time.Now() + storeMu.Unlock() + + return resolved, nil } diff --git a/internal/repository/openflare_observability.go b/internal/repository/openflare_observability.go index 8917b57c..c2d4606c 100644 --- a/internal/repository/openflare_observability.go +++ b/internal/repository/openflare_observability.go @@ -18,6 +18,7 @@ import ( analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics" analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics" "github.com/Rain-kl/Wavelet/internal/repository/logstore" + "github.com/Rain-kl/Wavelet/pkg/logger" ) const ( @@ -105,16 +106,20 @@ func ListOpenFlareMetricSnapshotsSince(ctx context.Context, nodeID string, since // avoiding stale reads of the previous CH store after a migration. func ListOpenFlareLatestMetricSnapshotsSince(ctx context.Context, nodeID string, since time.Time) ([]*model.OpenFlareMetricSnapshot, error) { active, err := logstore.ActiveDatabase(ctx) - if err == nil && active == logStoreNameClickHouse { - rows, err := analyticsrepo.ListLatestNodeMetricSnapshots(ctx, analyticsrepo.NodeObservabilityFilter{ + if err != nil { + logger.ErrorF(ctx, "failed to resolve active log database for latest metric snapshots: %v", err) + } else if active == logStoreNameClickHouse { + rows, chErr := analyticsrepo.ListLatestNodeMetricSnapshots(ctx, analyticsrepo.NodeObservabilityFilter{ NodeID: nodeID, Since: since, }) - if err == nil { + if chErr == nil { return fromAnalyticsNodeMetricSnapshots(rows), nil } + logger.ErrorF(ctx, "clickhouse fast-path ListLatestNodeMetricSnapshots failed: %v", chErr) + return nil, chErr } - // Routes through the active log store (CH unavailable or PG/SQLite active). + // Routes through the active log store (PG/SQLite active). all, listErr := ListOpenFlareMetricSnapshotsSince(ctx, nodeID, since, 0) if listErr != nil { return nil, listErr diff --git a/internal/testhelper/test_helper.go b/internal/testhelper/test_helper.go index 9c152aa4..47fb3c12 100644 --- a/internal/testhelper/test_helper.go +++ b/internal/testhelper/test_helper.go @@ -401,21 +401,31 @@ func SetupLogStoresForTest(t *testing.T) { logstore.SetAccessLogHooks(logstore.AccessLogHooks{ QueueNodeAccessLogs: func(logs []analyticsmodel.NodeAccessLog) { - _ = store.AccessLogs.BatchInsertNodeAccessLogs(context.Background(), logs) + if err := store.AccessLogs.BatchInsertNodeAccessLogs(context.Background(), logs); err != nil { + t.Errorf("batch insert node access logs failed in test hook: %v", err) + } }, }) logstore.SetObservabilityHooks(logstore.ObservabilityHooks{ QueueMetricSnapshot: func(record analyticsmodel.NodeMetricSnapshot) { - _ = store.Observability.BatchInsertNodeMetricSnapshots(context.Background(), []analyticsmodel.NodeMetricSnapshot{record}) + if err := store.Observability.BatchInsertNodeMetricSnapshots(context.Background(), []analyticsmodel.NodeMetricSnapshot{record}); err != nil { + t.Errorf("batch insert node metric snapshots failed in test hook: %v", err) + } }, QueueEdgeHealth: func(record analyticsmodel.NodeEdgeHealth) { - _ = store.Observability.BatchInsertNodeEdgeHealth(context.Background(), []analyticsmodel.NodeEdgeHealth{record}) + if err := store.Observability.BatchInsertNodeEdgeHealth(context.Background(), []analyticsmodel.NodeEdgeHealth{record}); err != nil { + t.Errorf("batch insert node edge health failed in test hook: %v", err) + } }, QueueNodeObsFrps: func(record analyticsmodel.NodeObsFrps) { - _ = store.Observability.BatchInsertNodeObsFrps(context.Background(), []analyticsmodel.NodeObsFrps{record}) + if err := store.Observability.BatchInsertNodeObsFrps(context.Background(), []analyticsmodel.NodeObsFrps{record}); err != nil { + t.Errorf("batch insert node obs frps failed in test hook: %v", err) + } }, QueueNodeObsFrpc: func(record analyticsmodel.NodeObsFrpc) { - _ = store.Observability.BatchInsertNodeObsFrpc(context.Background(), []analyticsmodel.NodeObsFrpc{record}) + if err := store.Observability.BatchInsertNodeObsFrpc(context.Background(), []analyticsmodel.NodeObsFrpc{record}); err != nil { + t.Errorf("batch insert node obs frpc failed in test hook: %v", err) + } }, })