fix(log): address remaining CodeRabbit suggestions for log database switch and migrations

This commit is contained in:
ryan
2026-08-08 20:36:04 +08:00
parent a6fc2b7737
commit 7d93d3d2a1
5 changed files with 21 additions and 8 deletions
@@ -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
@@ -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
@@ -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;
@@ -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;
+6
View File
@@ -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