mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-03 15:06:36 +08:00
fix(log): address CodeRabbit review findings for log database decoupling
This commit is contained in:
@@ -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)。
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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 {
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user