feat(log): 解耦用户访问日志存储,支持切换日志主库

用户访问日志可在 ClickHouse、PostgreSQL、SQLite 之间切换。
关闭 ClickHouse 时由主库承接写入与查询;切换任务会冻结写入、复制数据后翻转主库。
启动时校验日志主库与运行配置一致,定期清理按各库保留天数删除过期记录。
This commit is contained in:
ryan
2026-08-16 11:17:55 +08:00
parent 6a53619dd2
commit e52592b16d
34 changed files with 2325 additions and 145 deletions
@@ -8,6 +8,8 @@ import (
"context"
"fmt"
"time"
"github.com/Rain-kl/Wavelet/internal/infra/persistence"
analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics"
"gorm.io/gorm"
@@ -69,6 +71,28 @@ func ListAccessLogs(ctx context.Context, filter AccessLogFilter, page, pageSize
return logs, safeUint64Count(total), nil
}
// DeleteAllUserAccessLogs hard-deletes all user access logs via TRUNCATE.
func DeleteAllUserAccessLogs(ctx context.Context) (int64, error) {
if db.ChConn == nil {
return 0, fmt.Errorf("clickhouse connection is not initialized")
}
if err := db.ChConn.Exec(ctx, "TRUNCATE TABLE "+analyticsmodel.UserAccessLog{}.TableName()); err != nil {
return 0, fmt.Errorf("truncate user access logs: %w", err)
}
return 0, nil
}
// DeleteUserAccessLogsBefore deletes user access logs older than cutoff.
func DeleteUserAccessLogsBefore(ctx context.Context, cutoff time.Time) (int64, error) {
if db.ChConn == nil {
return 0, fmt.Errorf("clickhouse connection is not initialized")
}
if err := db.ChConn.Exec(ctx, "ALTER TABLE "+analyticsmodel.UserAccessLog{}.TableName()+" DELETE WHERE created_at < ?", cutoff); err != nil {
return 0, fmt.Errorf("delete expired user access logs: %w", err)
}
return 0, nil
}
func safeUint64Count(count int64) uint64 {
if count < 0 {
return 0
+93
View File
@@ -0,0 +1,93 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package logstore
import (
"context"
"errors"
"fmt"
"strconv"
"time"
"github.com/Rain-kl/Wavelet/internal/model"
"github.com/Rain-kl/Wavelet/pkg/logger"
)
const (
defaultLogRetentionDays = 30
partitionLeadMonths = 2
userAccessLogTable = "w_user_access_logs"
)
// CleanupSummary 汇总本次清理结果。
type CleanupSummary struct {
ActiveDatabase string `json:"active_database"`
RetentionDays int `json:"retention_days"`
Deleted int64 `json:"deleted"`
}
// CleanupExpired 按当前日志库保留天数删除过期用户访问日志,并预建 PG 分区。
func CleanupExpired(ctx context.Context) (CleanupSummary, error) {
active, err := ActiveDatabase(ctx)
if err != nil {
return CleanupSummary{}, err
}
days := retentionDaysForDatabase(ctx, active)
summary := CleanupSummary{ActiveDatabase: active, RetentionDays: days}
store, err := Active(ctx)
if err != nil {
return summary, err
}
now := time.Now().UTC()
if err := store.UserAccessLogs.EnsurePartitions(ctx, now, now.AddDate(0, partitionLeadMonths, 0)); err != nil {
logger.WarnF(ctx, "logstore: ensure partitions during cleanup failed: %v", err)
}
cutoff := now.AddDate(0, 0, -days)
deleted, err := store.UserAccessLogs.DeleteBefore(ctx, cutoff)
if err != nil {
return summary, fmt.Errorf("delete expired user access logs: %w", err)
}
summary.Deleted = deleted
return summary, nil
}
func retentionDaysForDatabase(ctx context.Context, dbName string) int {
key := model.ConfigKeyLogRetentionDaysPostgres
switch dbName {
case dbNameSQLite:
key = model.ConfigKeyLogRetentionDaysSQLite
case dbNameClickHouse:
key = model.ConfigKeyLogRetentionDaysClickHouse
}
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
}
func partitionStatementsRange(from, to time.Time) []string {
var out []string
start := time.Date(from.Year(), from.Month(), 1, 0, 0, 0, 0, time.UTC)
end := time.Date(to.Year(), to.Month(), 1, 0, 0, 0, 0, time.UTC).AddDate(0, 1, 0)
for ; start.Before(end); start = start.AddDate(0, 1, 0) {
monthEnd := start.AddDate(0, 1, 0)
suffix := start.Format("200601")
fromDay := start.Format("2006-01-02")
toDay := monthEnd.Format("2006-01-02")
out = append(out, fmt.Sprintf(
"CREATE TABLE IF NOT EXISTS %s_%s PARTITION OF %s FOR VALUES FROM ('%s') TO ('%s')",
userAccessLogTable, suffix, userAccessLogTable, fromDay, toDay))
}
return out
}
+146
View File
@@ -0,0 +1,146 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package logstore
import (
"context"
"fmt"
"time"
"github.com/ClickHouse/clickhouse-go/v2/lib/driver"
db "github.com/Rain-kl/Wavelet/internal/infra/persistence"
analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics"
analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics"
)
type clickhouseUserAccessLogStore struct {
skipFreeze bool
}
func newClickHouseUserAccessLogStore() *clickhouseUserAccessLogStore {
return &clickhouseUserAccessLogStore{}
}
var (
_ UserAccessLogStore = (*clickhouseUserAccessLogStore)(nil)
_ StatusStore = (*clickhouseUserAccessLogStore)(nil)
)
func (s *clickhouseUserAccessLogStore) ActiveDatabase(_ context.Context) (string, error) {
return dbNameClickHouse, nil
}
func (s *clickhouseUserAccessLogStore) ensureWritable(ctx context.Context) error {
if !s.skipFreeze && Migrating(ctx) {
return ErrMigrating
}
return nil
}
func (s *clickhouseUserAccessLogStore) BatchInsert(ctx context.Context, logs []analyticsmodel.UserAccessLog) error {
if len(logs) == 0 {
return nil
}
if err := s.ensureWritable(ctx); err != nil {
return err
}
return analyticsrepo.BatchInsert(ctx, logs)
}
func (s *clickhouseUserAccessLogStore) DeleteAll(ctx context.Context) (int64, error) {
if err := s.ensureWritable(ctx); err != nil {
return 0, err
}
return analyticsrepo.DeleteAllUserAccessLogs(ctx)
}
func (s *clickhouseUserAccessLogStore) DeleteBefore(ctx context.Context, cutoff time.Time) (int64, error) {
if err := s.ensureWritable(ctx); err != nil {
return 0, err
}
return analyticsrepo.DeleteUserAccessLogsBefore(ctx, cutoff)
}
func (s *clickhouseUserAccessLogStore) Count(ctx context.Context, filter analyticsrepo.AccessLogFilter) (uint64, error) {
return analyticsrepo.CountAccessLogs(ctx, filter)
}
func (s *clickhouseUserAccessLogStore) List(ctx context.Context, filter analyticsrepo.AccessLogFilter, page, pageSize int) ([]analyticsmodel.UserAccessLog, uint64, error) {
return analyticsrepo.ListAccessLogs(ctx, filter, page, pageSize)
}
func (s *clickhouseUserAccessLogStore) GetDailyTrend(ctx context.Context, days int) ([]analyticsrepo.DailyTrend, error) {
return analyticsrepo.GetDailyTrend(ctx, days)
}
func (s *clickhouseUserAccessLogStore) GetBrowserDistribution(ctx context.Context, startTime time.Time) ([]analyticsrepo.BrowserShare, error) {
return analyticsrepo.GetBrowserDistribution(ctx, startTime)
}
func (s *clickhouseUserAccessLogStore) GetTopActiveUsers(ctx context.Context, startTime time.Time, limit int) ([]analyticsrepo.TopUser, error) {
return analyticsrepo.GetTopActiveUsers(ctx, startTime, limit)
}
func (s *clickhouseUserAccessLogStore) EnsurePartitions(_ context.Context, _, _ time.Time) error {
return nil
}
func (s *clickhouseUserAccessLogStore) MigrationRange(ctx context.Context) (time.Time, time.Time, error) {
if db.ChConn == nil {
return time.Time{}, time.Time{}, fmt.Errorf("clickhouse connection is not initialized")
}
table := analyticsmodel.UserAccessLog{}.TableName()
var minTime, maxTime *time.Time
if err := db.ChConn.QueryRow(ctx, "SELECT min(created_at), max(created_at) FROM "+table).Scan(&minTime, &maxTime); err != nil {
return time.Time{}, time.Time{}, fmt.Errorf("query migration range %s: %w", table, err)
}
if minTime == nil || maxTime == nil {
return time.Time{}, time.Time{}, nil
}
return minTime.UTC(), maxTime.UTC(), nil
}
func (s *clickhouseUserAccessLogStore) ListForMigration(ctx context.Context, afterID uint64, limit int) ([]analyticsmodel.UserAccessLog, error) {
if db.ChConn == nil {
return nil, fmt.Errorf("clickhouse connection is not initialized")
}
if limit <= 0 {
limit = migrationPageSize
}
table := analyticsmodel.UserAccessLog{}.TableName()
columns := analyticsmodel.UserAccessLog{}.InsertColumns()
rows, err := db.ChConn.Query(ctx, fmt.Sprintf(
"SELECT %s FROM %s WHERE id > ? ORDER BY id ASC LIMIT ?",
columns, table,
), afterID, limit)
if err != nil {
return nil, fmt.Errorf("list user access logs for migration: %w", err)
}
defer func() { _ = rows.Close() }()
return scanUserAccessLogs(rows)
}
func scanUserAccessLogs(rows driver.Rows) ([]analyticsmodel.UserAccessLog, error) {
var result []analyticsmodel.UserAccessLog
for rows.Next() {
var item analyticsmodel.UserAccessLog
if err := rows.Scan(
&item.ID,
&item.UserID,
&item.Path,
&item.Method,
&item.IP,
&item.UserAgent,
&item.Headers,
&item.Status,
&item.Latency,
&item.CreatedAt,
); err != nil {
return nil, fmt.Errorf("scan user access log row: %w", err)
}
item.CreatedAt = item.CreatedAt.UTC()
result = append(result, item)
}
return result, nil
}
+327
View File
@@ -0,0 +1,327 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package logstore
import (
"context"
"errors"
"fmt"
"sort"
"strings"
"time"
"github.com/Rain-kl/Wavelet/internal/infra/persistence/idgen"
analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics"
analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics"
"gorm.io/gorm"
)
const (
insertBatchSize = 500
migrationPageSize = 100
defaultPageSize = 20
defaultTopN = 10
topUserAgents = 100
dayDuration = 24 * time.Hour
)
type gormLogStore struct {
db *gorm.DB
skipFreeze bool
}
func newGormStore(db *gorm.DB) *gormLogStore { return &gormLogStore{db: db} }
type userAccessLogGormStore struct {
*gormLogStore
}
func newUserAccessLogGormStore(db *gorm.DB) *userAccessLogGormStore {
return &userAccessLogGormStore{gormLogStore: newGormStore(db)}
}
var (
_ UserAccessLogStore = (*userAccessLogGormStore)(nil)
_ StatusStore = (*userAccessLogGormStore)(nil)
)
func (s *gormLogStore) ActiveDatabase(_ context.Context) (string, error) {
if isPostgresDialect(s.db) {
return dbNamePostgres, nil
}
return dbNameSQLite, nil
}
func (s *gormLogStore) ensureWritable(ctx context.Context) error {
if !s.skipFreeze && Migrating(ctx) {
return ErrMigrating
}
return nil
}
func (s *userAccessLogGormStore) BatchInsert(ctx context.Context, logs []analyticsmodel.UserAccessLog) error {
if len(logs) == 0 {
return nil
}
if err := s.ensureWritable(ctx); err != nil {
return err
}
for i := range logs {
if logs[i].ID == 0 {
logs[i].ID = idgen.NextUint64ID()
}
}
return s.db.WithContext(ctx).CreateInBatches(logs, insertBatchSize).Error
}
func (s *userAccessLogGormStore) DeleteAll(ctx context.Context) (int64, error) {
if err := s.ensureWritable(ctx); err != nil {
return 0, err
}
res := s.db.WithContext(ctx).Where("1 = 1").Delete(&analyticsmodel.UserAccessLog{})
return res.RowsAffected, res.Error
}
func (s *userAccessLogGormStore) DeleteBefore(ctx context.Context, cutoff time.Time) (int64, error) {
if err := s.ensureWritable(ctx); err != nil {
return 0, err
}
res := s.db.WithContext(ctx).Where("created_at < ?", cutoff).Delete(&analyticsmodel.UserAccessLog{})
if res.Error != nil && isMissingRelation(res.Error) {
return 0, nil
}
return res.RowsAffected, res.Error
}
func (s *userAccessLogGormStore) ListForMigration(ctx context.Context, afterID uint64, limit int) ([]analyticsmodel.UserAccessLog, error) {
var rows []analyticsmodel.UserAccessLog
q := s.db.WithContext(ctx).Model(&analyticsmodel.UserAccessLog{}).
Where("id > ?", afterID).
Order("id ASC").
Limit(limitOr(limit, migrationPageSize))
if err := q.Find(&rows).Error; err != nil {
return nil, err
}
return rows, nil
}
func (s *userAccessLogGormStore) MigrationRange(ctx context.Context) (time.Time, time.Time, error) {
return gormMigrationRange(ctx, s.db, "created_at", analyticsmodel.UserAccessLog{}, func(v *analyticsmodel.UserAccessLog) time.Time {
return v.CreatedAt
})
}
func (s *userAccessLogGormStore) Count(ctx context.Context, filter analyticsrepo.AccessLogFilter) (uint64, error) {
where, args, ok := buildUserAccessLogWhere(filter)
if !ok {
return 0, nil
}
var total int64
if err := s.db.WithContext(ctx).Model(&analyticsmodel.UserAccessLog{}).Where(where, args...).Count(&total).Error; err != nil {
return 0, err
}
return countToUint64(total), nil
}
func (s *userAccessLogGormStore) List(ctx context.Context, filter analyticsrepo.AccessLogFilter, page, pageSize int) ([]analyticsmodel.UserAccessLog, uint64, error) {
where, args, ok := buildUserAccessLogWhere(filter)
if !ok {
return []analyticsmodel.UserAccessLog{}, 0, nil
}
var total int64
if err := s.db.WithContext(ctx).Model(&analyticsmodel.UserAccessLog{}).Where(where, args...).Count(&total).Error; err != nil {
return nil, 0, err
}
if total == 0 {
return []analyticsmodel.UserAccessLog{}, 0, nil
}
var rows []analyticsmodel.UserAccessLog
q := s.db.WithContext(ctx).Where(where, args...).Order("created_at DESC, id DESC")
if err := q.Limit(limitOr(pageSize, defaultPageSize)).Offset(offsetOf(page, pageSize)).Find(&rows).Error; err != nil {
return nil, 0, err
}
return rows, countToUint64(total), nil
}
func buildUserAccessLogWhere(filter analyticsrepo.AccessLogFilter) (string, []any, bool) {
if filter.UserIDs != nil && len(filter.UserIDs) == 0 {
return "", nil, false
}
var parts []string
var args []any
if filter.UserIDs != nil {
parts = append(parts, "user_id IN ?")
args = append(args, filter.UserIDs)
}
if trimmed := strings.TrimSpace(filter.Path); trimmed != "" {
parts = append(parts, "path LIKE ?")
args = append(args, "%"+trimmed+"%")
}
if filter.StartTime != nil {
parts = append(parts, "created_at >= ?")
args = append(args, *filter.StartTime)
}
if filter.EndTime != nil {
parts = append(parts, "created_at <= ?")
args = append(args, *filter.EndTime)
}
if len(parts) == 0 {
return "1 = 1", args, true
}
return strings.Join(parts, " AND "), args, true
}
func (s *userAccessLogGormStore) GetDailyTrend(ctx context.Context, days int) ([]analyticsrepo.DailyTrend, error) {
if days <= 0 {
days = 7
}
start := time.Now().AddDate(0, 0, -(days - 1)).Truncate(dayDuration)
type row struct {
Date string
Cnt uint64
}
var rows []row
err := s.db.WithContext(ctx).Model(&analyticsmodel.UserAccessLog{}).
Select(dailyTrendDateSQL(s.db)+" AS date, COUNT(*) AS cnt").
Where("created_at >= ?", start).
Group("date").Order("date ASC").Scan(&rows).Error
if err != nil {
return nil, err
}
counts := make(map[string]uint64, len(rows))
for _, r := range rows {
counts[r.Date] = r.Cnt
}
out := make([]analyticsrepo.DailyTrend, 0, days)
for i := 0; i < days; i++ {
d := start.AddDate(0, 0, i).Format("2006-01-02")
out = append(out, analyticsrepo.DailyTrend{Date: d, Count: counts[d]})
}
return out, nil
}
func (s *userAccessLogGormStore) GetBrowserDistribution(ctx context.Context, startTime time.Time) ([]analyticsrepo.BrowserShare, error) {
type row struct {
UserAgent string
Cnt uint64
}
var rows []row
err := s.db.WithContext(ctx).Model(&analyticsmodel.UserAccessLog{}).
Select("user_agent, COUNT(*) AS cnt").
Where("created_at >= ?", startTime).
Group("user_agent").Order("cnt DESC").Limit(topUserAgents).Scan(&rows).Error
if err != nil {
return nil, err
}
counts := make(map[string]uint64)
for _, r := range rows {
counts[analyticsrepo.ParseBrowserName(r.UserAgent)] += r.Cnt
}
out := make([]analyticsrepo.BrowserShare, 0, len(counts))
for label, count := range counts {
out = append(out, analyticsrepo.BrowserShare{Browser: label, Count: count})
}
sort.Slice(out, func(i, j int) bool { return out[i].Count > out[j].Count })
return out, nil
}
func (s *userAccessLogGormStore) GetTopActiveUsers(ctx context.Context, startTime time.Time, limit int) ([]analyticsrepo.TopUser, error) {
type row struct {
UserID uint64
Cnt uint64
}
var rows []row
err := s.db.WithContext(ctx).Model(&analyticsmodel.UserAccessLog{}).
Select("user_id, COUNT(*) AS cnt").
Where("user_id <> 0 AND created_at >= ?", startTime).
Group("user_id").Order("cnt DESC").Limit(limitOr(limit, defaultTopN)).Scan(&rows).Error
if err != nil {
return nil, err
}
out := make([]analyticsrepo.TopUser, len(rows))
for i, r := range rows {
out[i] = analyticsrepo.TopUser{UserID: r.UserID, Count: r.Cnt}
}
return out, nil
}
func (s *userAccessLogGormStore) EnsurePartitions(ctx context.Context, from, to time.Time) error {
if !isPostgresDialect(s.db) {
return nil
}
for _, sql := range partitionStatementsRange(from, to) {
if err := s.db.WithContext(ctx).Exec(sql).Error; err != nil {
return fmt.Errorf("ensure partition: %w", err)
}
}
return nil
}
func gormMigrationRange[T any](
ctx context.Context,
gdb *gorm.DB,
column string,
model T,
timeOf func(*T) time.Time,
) (time.Time, time.Time, error) {
var first, last T
found := false
for _, order := range []string{"ASC", "DESC"} {
out := &first
if order == "DESC" {
out = &last
}
res := gdb.WithContext(ctx).Model(model).Order(column + " " + order).Limit(1).Take(out)
if res.Error != nil && !errors.Is(res.Error, gorm.ErrRecordNotFound) {
return time.Time{}, time.Time{}, fmt.Errorf("query migration range %s: %w", column, res.Error)
}
if res.Error == nil {
found = true
}
}
if !found {
return time.Time{}, time.Time{}, nil
}
return timeOf(&first).UTC(), timeOf(&last).UTC(), nil
}
func limitOr(v, def int) int {
if v <= 0 {
return def
}
return v
}
func offsetOf(page, pageSize int) int {
if page < 1 {
page = 1
}
return (page - 1) * limitOr(pageSize, defaultPageSize)
}
func countToUint64(v int64) uint64 {
if v < 0 {
return 0
}
return uint64(v)
}
func isPostgresDialect(db *gorm.DB) bool {
return db != nil && db.Dialector != nil && db.Name() == "postgres"
}
func dailyTrendDateSQL(db *gorm.DB) string {
if isPostgresDialect(db) {
return "to_char(created_at, 'YYYY-MM-DD')"
}
return "strftime('%Y-%m-%d', created_at)"
}
func isMissingRelation(err error) bool {
if err == nil {
return false
}
msg := strings.ToLower(err.Error())
return strings.Contains(msg, "no such table") || strings.Contains(msg, "does not exist")
}
+60
View File
@@ -0,0 +1,60 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package logstore
import (
"context"
"testing"
"time"
analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics"
analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics"
"github.com/glebarez/sqlite"
"github.com/stretchr/testify/require"
"gorm.io/gorm"
)
func newTestUserAccessStore(t *testing.T) *userAccessLogGormStore {
t.Helper()
gdb, err := gorm.Open(sqlite.Open("file:logstore-"+t.Name()+"?mode=memory&cache=shared"), &gorm.Config{})
require.NoError(t, err)
require.NoError(t, gdb.AutoMigrate(&analyticsmodel.UserAccessLog{}))
return newUserAccessLogGormStore(gdb)
}
func TestGormUserAccessLogCountList(t *testing.T) {
ua := newTestUserAccessStore(t)
ctx := context.Background()
now := time.Now().UTC().Truncate(time.Second)
require.NoError(t, ua.BatchInsert(ctx, []analyticsmodel.UserAccessLog{
{UserID: 10, Path: "/api/v1/users", Method: "GET", Status: 200, CreatedAt: now},
{UserID: 20, Path: "/api/v1/admin", Method: "GET", Status: 200, CreatedAt: now},
{UserID: 10, Path: "/api/v1/other", Method: "POST", Status: 201, CreatedAt: now},
}))
count, err := ua.Count(ctx, analyticsrepo.AccessLogFilter{UserIDs: []uint64{10}, Path: "users"})
require.NoError(t, err)
require.Equal(t, uint64(1), count)
rows, total, err := ua.List(ctx, analyticsrepo.AccessLogFilter{UserIDs: []uint64{10}, Path: "users"}, 1, 10)
require.NoError(t, err)
require.Equal(t, uint64(1), total)
require.Len(t, rows, 1)
require.Equal(t, "/api/v1/users", rows[0].Path)
require.NotZero(t, rows[0].ID)
}
func TestGormUserAccessLogFreeze(t *testing.T) {
ua := newTestUserAccessStore(t)
SetConfigReader(func(_ context.Context, key string) (string, error) {
if key == logMigrationKey {
return "migrating", nil
}
return "", nil
})
t.Cleanup(ResetForTest)
err := ua.BatchInsert(context.Background(), []analyticsmodel.UserAccessLog{{UserID: 1, CreatedAt: time.Now()}})
require.ErrorIs(t, err, ErrMigrating)
}
+43
View File
@@ -0,0 +1,43 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
// Package logstore abstracts user access-log storage across ClickHouse, PostgreSQL and SQLite.
package logstore
import (
"context"
"errors"
"time"
analyticsmodel "github.com/Rain-kl/Wavelet/internal/model/analytics"
analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics"
)
// ErrMigrating 表示日志数据库正在迁移,当前禁止写入。
var ErrMigrating = errors.New("log database is migrating, writes are disabled")
// UserAccessLogStore 用户访问日志(w_user_access_logs)。
type UserAccessLogStore interface {
BatchInsert(ctx context.Context, logs []analyticsmodel.UserAccessLog) error
DeleteAll(ctx context.Context) (int64, error)
DeleteBefore(ctx context.Context, cutoff time.Time) (int64, error)
Count(ctx context.Context, filter analyticsrepo.AccessLogFilter) (uint64, error)
List(ctx context.Context, filter analyticsrepo.AccessLogFilter, page, pageSize int) ([]analyticsmodel.UserAccessLog, uint64, error)
GetDailyTrend(ctx context.Context, days int) ([]analyticsrepo.DailyTrend, error)
GetBrowserDistribution(ctx context.Context, startTime time.Time) ([]analyticsrepo.BrowserShare, error)
GetTopActiveUsers(ctx context.Context, startTime time.Time, limit int) ([]analyticsrepo.TopUser, error)
ListForMigration(ctx context.Context, afterID uint64, limit int) ([]analyticsmodel.UserAccessLog, error)
MigrationRange(ctx context.Context) (from, to time.Time, err error)
EnsurePartitions(ctx context.Context, from, to time.Time) error
}
// StatusStore 日志库状态。
type StatusStore interface {
ActiveDatabase(ctx context.Context) (string, error)
}
// Store 当前生效日志库。
type Store struct {
UserAccessLogs UserAccessLogStore
Status StatusStore
}
+189
View File
@@ -0,0 +1,189 @@
// Copyright 2026 Arctel.net
// SPDX-License-Identifier: Apache-2.0
package logstore
import (
"context"
"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"
)
const (
logDatabaseKey = model.ConfigKeyLogDatabase
logMigrationKey = model.ConfigKeyLogDBMigration
)
const (
dbNamePostgres = "postgres"
dbNameSQLite = "sqlite"
dbNameClickHouse = "clickhouse"
)
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
lastResolveDB string
lastResolveTime time.Time
)
// SetConfigReader 注入系统配置读取函数(bootstrap 调用,测试可注入内存实现)。
func SetConfigReader(fn ConfigReader) { configReader = fn }
func getConfig(ctx context.Context, key string) (string, error) {
if configReader == nil {
return "", errConfigReaderNotWired
}
return configReader(ctx, key)
}
// Active 返回当前生效的日志库 Store。
func Active(ctx context.Context) (*Store, error) {
current, err := resolveDatabase(ctx)
if err != nil {
return nil, err
}
storeMu.RLock()
if active != nil && activeDB == current {
s := active
storeMu.RUnlock()
return s, nil
}
storeMu.RUnlock()
storeMu.Lock()
defer storeMu.Unlock()
if active != nil && activeDB == current {
return active, nil
}
s, err := buildStore(ctx, current, false)
if err != nil {
return nil, err
}
active = s
activeDB = current
return s, nil
}
// Build 直接按目标构造 store(不经 Active 缓存)。
func Build(ctx context.Context, database string) (*Store, error) {
return buildStore(ctx, database, false)
}
// BuildForMigration 构造迁移目标 store,跳过冻结检查。
func BuildForMigration(ctx context.Context, database string) (*Store, error) {
return buildStore(ctx, database, true)
}
func buildStore(ctx context.Context, database string, skipFreeze bool) (*Store, error) {
switch database {
case dbNameClickHouse:
ual := newClickHouseUserAccessLogStore()
ual.skipFreeze = skipFreeze
return &Store{UserAccessLogs: ual, Status: ual}, nil
case dbNamePostgres, dbNameSQLite:
gdb := db.DB(ctx)
ual := newUserAccessLogGormStore(gdb)
ual.skipFreeze = skipFreeze
return &Store{UserAccessLogs: ual, Status: ual}, nil
default:
return nil, fmt.Errorf("unsupported log database: %s", database)
}
}
// Migrating 返回日志库是否处于迁移冻结状态。
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"
}
// Init 预热激活 store,并兜底预建当前月及未来分区。
func Init(ctx context.Context) {
s, err := Active(ctx)
if err != nil {
return
}
now := time.Now().UTC()
if err := s.UserAccessLogs.EnsurePartitions(ctx, now, now.AddDate(0, partitionLeadMonths, 0)); err != nil {
logger.WarnF(ctx, "logstore: ensure startup partitions failed: %v", err)
}
}
// InvalidateCache 清空日志库解析缓存。
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
}
// ActiveDatabase 返回当前日志主库名。
func ActiveDatabase(ctx context.Context) (string, error) {
return resolveDatabase(ctx)
}
func resolveDatabase(ctx context.Context) (string, error) {
storeMu.RLock()
if active != nil && time.Since(lastResolveTime) < resolveCacheTTL {
name := lastResolveDB
storeMu.RUnlock()
return name, nil
}
storeMu.RUnlock()
v, err := getConfig(ctx, logDatabaseKey)
if err != nil && !errors.Is(err, errConfigReaderNotWired) {
return "", err
}
resolved := v
if resolved == "" {
resolved = dbNameSQLite
if config.Config.Database.Enabled {
resolved = dbNamePostgres
}
if config.Config.ClickHouse.Enabled {
resolved = dbNameClickHouse
}
}
storeMu.Lock()
lastResolveDB = resolved
lastResolveTime = time.Now()
storeMu.Unlock()
return resolved, nil
}