From 7d93d3d2a180464a4241820a98ef1b240a60b3df Mon Sep 17 00:00:00 2001 From: ryan Date: Sat, 8 Aug 2026 20:36:04 +0800 Subject: [PATCH] fix(log): address remaining CodeRabbit suggestions for log database switch and migrations --- .../plans/2026-08-08-log-database-decoupling.md | 4 ++++ internal/apps/openflare/tasks/log_db_switch.go | 13 +++++++------ .../202608080003_drop_database_cleanup_schedule.sql | 3 ++- .../202608080003_drop_database_cleanup_schedule.sql | 3 ++- internal/repository/logstore/cleanup.go | 6 ++++++ 5 files changed, 21 insertions(+), 8 deletions(-) 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 aa511d75..2524f000 100644 --- a/docs/superpowers/plans/2026-08-08-log-database-decoupling.md +++ b/docs/superpowers/plans/2026-08-08-log-database-decoupling.md @@ -2250,6 +2250,10 @@ func (h *LogDBSwitchHandler) Execute(ctx context.Context, payload []byte) (*task } defer func() { _ = setMigrationFlag(ctx, "") }() // 失败也清除,保持源库可写 + if err := drainLogWriters(ctx); err != nil { + return nil, fmt.Errorf("排空日志写入队列失败: %w", err) + } + src, err := logstore.Active(ctx) if err != nil { return nil, err diff --git a/internal/apps/openflare/tasks/log_db_switch.go b/internal/apps/openflare/tasks/log_db_switch.go index 3585c642..12f2053b 100644 --- a/internal/apps/openflare/tasks/log_db_switch.go +++ b/internal/apps/openflare/tasks/log_db_switch.go @@ -88,12 +88,7 @@ func (h *LogDBSwitchHandler) Execute(ctx context.Context, payload []byte) (*task } task.AppendLog(ctx, "开始切换日志数据库:%s -> %s", source, p.Target) - // 冻结前先排空在途批次(chwriter + 用户访问日志 writer),避免迁移期间积压丢失; - // 只等 flush 完成,不停止 writer(置位后由 ensureWritable 拒绝新写入)。 - if err := drainLogWriters(ctx); err != nil { - return nil, fmt.Errorf("排空日志写入队列失败: %w", err) - } - + // 设置迁移冻结标记(置位后由 ensureWritable 拒绝新写入)。 if err := setMigrationFlag(ctx, "migrating"); err != nil { return nil, err } @@ -104,6 +99,12 @@ func (h *LogDBSwitchHandler) Execute(ctx context.Context, payload []byte) (*task } }() + // 冻结标记置位后再排空在途批次(chwriter + 用户访问日志 writer), + // 保证排空完成后不再有新批次进入源库。 + if err := drainLogWriters(ctx); err != nil { + return nil, fmt.Errorf("排空日志写入队列失败: %w", err) + } + src, err := logstore.Active(ctx) if err != nil { return nil, err 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 2a8dfed6..34663806 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 @@ -11,7 +11,8 @@ 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) +SELECT 102, 'OpenFlare 可观测数据自动清理', 'of_database_auto_cleanup', '0 3 * * *', '{}', TRUE, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP +WHERE NOT EXISTS (SELECT 1 FROM w_schedules WHERE task_type = 'of_database_auto_cleanup') 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/202608080003_drop_database_cleanup_schedule.sql b/internal/infra/persistence/migrator/goose/sqlite/202608080003_drop_database_cleanup_schedule.sql index 4ad60bbb..e8880951 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 @@ -10,7 +10,8 @@ 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); +SELECT 102, 'OpenFlare 可观测数据自动清理', 'of_database_auto_cleanup', '0 3 * * *', '{}', 1, CURRENT_TIMESTAMP, CURRENT_TIMESTAMP +WHERE NOT EXISTS (SELECT 1 FROM w_schedules WHERE task_type = 'of_database_auto_cleanup'); DROP TABLE IF EXISTS w_schedules_backup_of_database_auto_cleanup; diff --git a/internal/repository/logstore/cleanup.go b/internal/repository/logstore/cleanup.go index 3b18fb69..ea800178 100644 --- a/internal/repository/logstore/cleanup.go +++ b/internal/repository/logstore/cleanup.go @@ -5,11 +5,13 @@ package logstore import ( "context" + "errors" "fmt" "strconv" "time" "github.com/Rain-kl/Wavelet/internal/model" + "github.com/Rain-kl/Wavelet/pkg/logger" ) // CleanupSummary 汇总本次清理结果。 @@ -39,10 +41,14 @@ func retentionDaysForDatabase(ctx context.Context, dbName string) int { } 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