mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-08 16:46:37 +08:00
feat(log): decouple log storage from ClickHouse with switchable logstore
- New internal/repository/logstore abstraction: exported domain interfaces (AccessLogStore/ObservabilityStore/UserAccessLogStore/StatusStore), config-driven provider (Active/Build/Migrating/SetConfigReader), GORM implementation for PostgreSQL/SQLite (incl. hourly rollups computed in real time, migration listers, PG partition maintenance), and a ClickHouse wrapper preserving the native batch path; repository facade delegates to logstore; import-lint test enforces apps never import analyticsrepo. - ClickHouse is now optional: the log DB is either the main DB (postgres when database.enabled, else sqlite) or clickhouse; boot validation + first-run seed; log_database / log_db_migration are protected keys. - New user task 切换日志数据库 (of_log_db_switch): freeze log writes, drain batch writers, copy all 6 raw log tables by id (preserving IDs) with target-partition pre-creation for PG, flip log_database on success, clear the freeze flag on failure. - Per-store retention (log_retention_days_*) with expiry cleanup folded into the daily system_cleanup task; legacy database_auto_cleanup_* and of_database_auto_cleanup decommissioned. - goose migrations: 6 log tables in PG (2 monthly-partitioned) + SQLite, retention config seeds, schedule cleanup; GET /api/v1/admin/status/log-database endpoint; frontend retention settings, switch-task UI and status badge; changelog and docs updated. docs(plan): log database decoupling implementation plan docs(design): log database decoupling design (ClickHouse optional)
This commit is contained in:
@@ -11,13 +11,10 @@ import (
|
||||
"sync"
|
||||
"time"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/repository"
|
||||
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/config"
|
||||
"github.com/Rain-kl/Wavelet/internal/infra/persistence/batchwriter"
|
||||
analyticsmodel "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"
|
||||
)
|
||||
|
||||
@@ -56,12 +53,10 @@ var (
|
||||
frpcDedup *dedupSet
|
||||
)
|
||||
|
||||
// Init starts OpenFlare ClickHouse batch writers. Safe to call multiple times.
|
||||
// Init starts OpenFlare log batch writers. Safe to call multiple times.
|
||||
// Writers always initialize regardless of ClickHouse.enabled; the active log
|
||||
// store is resolved via logstore at flush time (PG/SQLite when CH is not active).
|
||||
func Init(ctx context.Context) {
|
||||
if !config.Config.ClickHouse.Enabled {
|
||||
return
|
||||
}
|
||||
|
||||
initOnce.Do(func() {
|
||||
metricSnapshotDedup = newDedupSet()
|
||||
edgeHealthDedup = newDedupSet()
|
||||
@@ -70,25 +65,25 @@ func Init(ctx context.Context) {
|
||||
|
||||
metricSnapshotWriter = mustNewObservabilityWriter(
|
||||
"metric_snapshots",
|
||||
withFlushRetries(analyticsrepo.BatchInsertNodeMetricSnapshots),
|
||||
withFlushRetries(flushNodeMetricSnapshots),
|
||||
metricSnapshotDedup,
|
||||
metricSnapshotKey,
|
||||
)
|
||||
edgeHealthWriter = mustNewObservabilityWriter(
|
||||
"edge_health",
|
||||
withFlushRetries(analyticsrepo.BatchInsertNodeEdgeHealth),
|
||||
withFlushRetries(flushNodeEdgeHealth),
|
||||
edgeHealthDedup,
|
||||
edgeHealthKey,
|
||||
)
|
||||
frpsWriter = mustNewObservabilityWriter(
|
||||
"frps_obs",
|
||||
withFlushRetries(analyticsrepo.BatchInsertNodeObsFrps),
|
||||
withFlushRetries(flushNodeObsFrps),
|
||||
frpsDedup,
|
||||
frpsKey,
|
||||
)
|
||||
frpcWriter = mustNewObservabilityWriter(
|
||||
"frpc_obs",
|
||||
withFlushRetries(analyticsrepo.BatchInsertNodeObsFrpc),
|
||||
withFlushRetries(flushNodeObsFrpc),
|
||||
frpcDedup,
|
||||
frpcKey,
|
||||
)
|
||||
@@ -129,6 +124,51 @@ func Stop(ctx context.Context) error {
|
||||
return firstErr
|
||||
}
|
||||
|
||||
// Drain 等待所有 OpenFlare 日志 writer 的在途批次落库:轮询队列 Depth 归零后
|
||||
// 再保持一个最大 flush 周期(observabilityFlushEvery)持续为空才返回;
|
||||
// 不停止 writer(迁移冻结后由 ensureWritable 拒绝新写入)。未初始化时直接返回 nil。
|
||||
func Drain(ctx context.Context) error {
|
||||
return drainWriters(ctx, WriterStats, observabilityFlushEvery)
|
||||
}
|
||||
|
||||
// drainWriters 轮询 stats 直至所有队列 Depth=0 并持续 quietPeriod 无新积压。
|
||||
func drainWriters(ctx context.Context, stats func() []batchwriter.Stats, quietPeriod time.Duration) error {
|
||||
if !running() {
|
||||
return nil
|
||||
}
|
||||
ticker := time.NewTicker(drainPollInterval)
|
||||
defer ticker.Stop()
|
||||
var quietSince time.Time
|
||||
for {
|
||||
if allDepthZero(stats()) {
|
||||
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:
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
// drainPollInterval 队列轮询间隔。
|
||||
const drainPollInterval = 50 * time.Millisecond
|
||||
|
||||
func allDepthZero(stats []batchwriter.Stats) bool {
|
||||
for _, s := range stats {
|
||||
if s.Depth > 0 {
|
||||
return false
|
||||
}
|
||||
}
|
||||
return true
|
||||
}
|
||||
|
||||
// WriterStats returns queue depth and failure counters for all OpenFlare writers.
|
||||
func WriterStats() []batchwriter.Stats {
|
||||
writers := []statsProvider{
|
||||
@@ -237,7 +277,7 @@ func mustNewNodeAccessLogWriter() *batchwriter.Writer[analyticsmodel.NodeAccessL
|
||||
}
|
||||
writer, err := batchwriter.New[analyticsmodel.NodeAccessLog](
|
||||
cfg,
|
||||
withFlushRetries(analyticsrepo.BatchInsertNodeAccessLogs),
|
||||
withFlushRetries(flushNodeAccessLogs),
|
||||
batchwriter.WithDropHandler[analyticsmodel.NodeAccessLog](func(item analyticsmodel.NodeAccessLog) {
|
||||
logger.WarnF(context.Background(), "[OpenFlare] node access log queue full, dropping log for node %s path %s", item.NodeID, item.Path)
|
||||
}),
|
||||
@@ -280,17 +320,59 @@ func withFlushRetries[T any](flush batchwriter.FlushFunc[T]) batchwriter.FlushFu
|
||||
}
|
||||
|
||||
func wireModelInsertHooks() {
|
||||
repository.SetObservabilityInsertHooks(repository.ObservabilityInsertHooks{
|
||||
QueueMetricSnapshot: QueueMetricSnapshot,
|
||||
QueueEdgeHealth: QueueEdgeHealth,
|
||||
QueueFrpsObservation: QueueFrpsObservation,
|
||||
QueueFrpcObservation: QueueFrpcObservation,
|
||||
logstore.SetObservabilityHooks(logstore.ObservabilityHooks{
|
||||
QueueMetricSnapshot: QueueMetricSnapshot,
|
||||
QueueEdgeHealth: QueueEdgeHealth,
|
||||
QueueNodeObsFrps: QueueFrpsObservation,
|
||||
QueueNodeObsFrpc: QueueFrpcObservation,
|
||||
})
|
||||
repository.SetAccessLogInsertHooks(repository.AccessLogInsertHooks{
|
||||
logstore.SetAccessLogHooks(logstore.AccessLogHooks{
|
||||
QueueNodeAccessLogs: QueueNodeAccessLogs,
|
||||
})
|
||||
}
|
||||
|
||||
// 以下 flush 函数作为 batchwriter 的落库目标:激活库由 logstore 在 flush 时决定。
|
||||
|
||||
func flushNodeMetricSnapshots(ctx context.Context, rows []analyticsmodel.NodeMetricSnapshot) error {
|
||||
s, err := logstore.Active(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return s.Observability.BatchInsertNodeMetricSnapshots(ctx, rows)
|
||||
}
|
||||
|
||||
func flushNodeEdgeHealth(ctx context.Context, rows []analyticsmodel.NodeEdgeHealth) error {
|
||||
s, err := logstore.Active(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return s.Observability.BatchInsertNodeEdgeHealth(ctx, rows)
|
||||
}
|
||||
|
||||
func flushNodeObsFrps(ctx context.Context, rows []analyticsmodel.NodeObsFrps) error {
|
||||
s, err := logstore.Active(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return s.Observability.BatchInsertNodeObsFrps(ctx, rows)
|
||||
}
|
||||
|
||||
func flushNodeObsFrpc(ctx context.Context, rows []analyticsmodel.NodeObsFrpc) error {
|
||||
s, err := logstore.Active(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return s.Observability.BatchInsertNodeObsFrpc(ctx, rows)
|
||||
}
|
||||
|
||||
func flushNodeAccessLogs(ctx context.Context, rows []analyticsmodel.NodeAccessLog) error {
|
||||
s, err := logstore.Active(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return s.AccessLogs.BatchInsertNodeAccessLogs(ctx, rows)
|
||||
}
|
||||
|
||||
func metricSnapshotKey(snapshot analyticsmodel.NodeMetricSnapshot) string {
|
||||
return fmt.Sprintf("%s|%d", snapshot.NodeID, snapshot.CapturedAt.UTC().UnixNano())
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user