diff --git a/.agents/skills/logstore/SKILL.md b/.agents/skills/logstore/SKILL.md index 64335519..8e991a8a 100644 --- a/.agents/skills/logstore/SKILL.md +++ b/.agents/skills/logstore/SKILL.md @@ -54,7 +54,7 @@ description: "Wavelet 项目专用:当新增或修改日志/分析用途表( - 写入:`BatchInsert`(flush 目标;内调 `ensureWritable`) - 查询:业务需要的 List/Count/聚合 - 迁移:`ListForMigration(afterID, limit)`、`MigrationRange`、`DeleteAll`、`EnsurePartitions`(PG 按月预建,CH/SQLite no-op) - - 清理:`DeleteBefore(cutoff)` + - 清理:`DeleteBefore(cutoff)`、`DropEmptyPartitions`、`DropExpiredPartitions`(仅 PG;CH/SQLite no-op) 4. **双实现** - CH:委托 `analyticsrepo`,零额外查询路径。 @@ -70,7 +70,7 @@ description: "Wavelet 项目专用:当新增或修改日志/分析用途表( 在 `copy*` 流程增加该表:`DeleteAll` 目标 → `MigrationRange` + `EnsurePartitions` → 按 id 分页复制。不要改切换协议(仍冻结写入、源数据不删、成功才翻转)。 8. **清理** - `CleanupExpired` 对该表 `DeleteBefore`;保留天数用已有 `log_retention_days_*`,不要为单表再发明一套 key,除非产品明确要求独立 TTL。 + `CleanupExpired`:PG 先 `DropExpiredPartitions`(整月过期分区),再 `DeleteBefore`(边界月),最后 `DropEmptyPartitions`。保留天数用已有 `log_retention_days_*`。apps 禁止 import `repository/analytics`(`imports_test.go`)。 ## 禁止 diff --git a/docs/LOGSTORE.md b/docs/LOGSTORE.md index 012a7333..caf8daa6 100644 --- a/docs/LOGSTORE.md +++ b/docs/LOGSTORE.md @@ -26,7 +26,7 @@ Wavelet 的访问审计等日志表不绑死 ClickHouse。`internal/repository/l | 主库实现 | `logstore` GORM | PostgreSQL 按月分区;SQLite 普通表 | | 入队 | `risk_control` + `batchwriter` | `FlushFunc` → `logstore.Active` | | 切换 | `logs:db_switch` | 冻结 → 排空 → 复制 → 翻转 | -| 清理 | `logstore.CleanupExpired` | `system:cleanup` 按库读 `log_retention_days_*` 后 `DeleteBefore` | +| 清理 | `logstore.CleanupExpired` | `system:cleanup` 按库读 `log_retention_days_*`:PG 先 `DropExpiredPartitions` 再 `DeleteBefore`,最后 `DropEmptyPartitions` | `log_database` 只能是「随业务主库」或 `clickhouse`。`log_database` / `log_db_migration` 受保护,管理端不可改。 diff --git a/internal/apps/admin/logs/routers.go b/internal/apps/admin/logs/routers.go index 73a10181..1d115e3c 100644 --- a/internal/apps/admin/logs/routers.go +++ b/internal/apps/admin/logs/routers.go @@ -15,7 +15,6 @@ import ( "github.com/Rain-kl/Wavelet/internal/apps/admin" "github.com/Rain-kl/Wavelet/internal/repository" - analyticsrepo "github.com/Rain-kl/Wavelet/internal/repository/analytics" "github.com/Rain-kl/Wavelet/internal/repository/logstore" "github.com/Rain-kl/Wavelet/pkg/logger" "github.com/gin-gonic/gin" @@ -157,8 +156,8 @@ type accessLogsResponse struct { List []accessLogItem `json:"list"` } -func buildAccessLogFilter(ctx context.Context, c *gin.Context) (analyticsrepo.AccessLogFilter, error) { - filter := analyticsrepo.AccessLogFilter{} +func buildAccessLogFilter(ctx context.Context, c *gin.Context) (logstore.AccessLogFilter, error) { + filter := logstore.AccessLogFilter{} username := c.Query("username") if username != "" { diff --git a/internal/repository/logstore/cleanup.go b/internal/repository/logstore/cleanup.go index 78439f69..8f347dec 100644 --- a/internal/repository/logstore/cleanup.go +++ b/internal/repository/logstore/cleanup.go @@ -45,11 +45,18 @@ func CleanupExpired(ctx context.Context) (CleanupSummary, error) { logger.WarnF(ctx, "logstore: ensure partitions during cleanup failed: %v", err) } cutoff := now.AddDate(0, 0, -days) + // 先 DROP 完全过期的整月分区,再对边界月逐行 DeleteBefore。 + if err := store.UserAccessLogs.DropExpiredPartitions(ctx, cutoff); err != nil { + return summary, fmt.Errorf("drop expired partitions: %w", err) + } deleted, err := store.UserAccessLogs.DeleteBefore(ctx, cutoff) if err != nil { return summary, fmt.Errorf("delete expired user access logs: %w", err) } summary.Deleted = deleted + if err := store.UserAccessLogs.DropEmptyPartitions(ctx, now); err != nil { + logger.WarnF(ctx, "drop empty log partitions failed: %v", err) + } return summary, nil } diff --git a/internal/repository/logstore/clickhouse.go b/internal/repository/logstore/clickhouse.go index 1e151cc5..8d0e0020 100644 --- a/internal/repository/logstore/clickhouse.go +++ b/internal/repository/logstore/clickhouse.go @@ -86,6 +86,14 @@ func (s *clickhouseUserAccessLogStore) EnsurePartitions(_ context.Context, _, _ return nil } +func (s *clickhouseUserAccessLogStore) DropEmptyPartitions(_ context.Context, _ time.Time) error { + return nil +} + +func (s *clickhouseUserAccessLogStore) DropExpiredPartitions(_ 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") diff --git a/internal/repository/logstore/imports_test.go b/internal/repository/logstore/imports_test.go new file mode 100644 index 00000000..73c1a745 --- /dev/null +++ b/internal/repository/logstore/imports_test.go @@ -0,0 +1,45 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package logstore + +import ( + "os/exec" + "strings" + "testing" +) + +// apps 禁止直连 analytics 做日志读写;查询过滤器请用 logstore.AccessLogFilter。 +var forbiddenImports = []string{ + "github.com/Rain-kl/Wavelet/internal/repository/analytics", +} + +// logstore 的 CH 实现按设计委托 analyticsrepo。 +var allowedAnalyticsDelegation = map[string]bool{ + "github.com/Rain-kl/Wavelet/internal/repository/logstore": true, +} + +func TestAppsMustNotImportLogBackendDirectly(t *testing.T) { + t.Chdir("../../..") + out, err := exec.Command("go", "list", "-test", "-f", `{{.ImportPath}} {{join .Imports " "}}`, "./internal/apps/...").Output() + if err != nil { + t.Fatalf("go list: %v", err) + } + for _, line := range strings.Split(string(out), "\n") { + fields := strings.Fields(line) + if len(fields) == 0 { + continue + } + pkg := fields[0] + if !strings.HasPrefix(pkg, "github.com/Rain-kl/Wavelet/internal/apps") { + continue + } + 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) + } + } + } + } +} diff --git a/internal/repository/logstore/logstore.go b/internal/repository/logstore/logstore.go index ef583010..19f63545 100644 --- a/internal/repository/logstore/logstore.go +++ b/internal/repository/logstore/logstore.go @@ -29,8 +29,17 @@ type UserAccessLogStore interface { 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 + // DropEmptyPartitions 幂等清理 PG 空分区表:删除 before 月份之前、且无任何数据的按月分区; + // CH/SQLite 为 no-op。 + DropEmptyPartitions(ctx context.Context, before time.Time) error + // DropExpiredPartitions 直接删除完全过期的 PG 整月分区(候选为月份早于 cutoff 月的分区, + // 删除前校验分区内无保留期内数据,避免时区偏移下误删;迁移冻结期间拒绝执行);CH/SQLite 为 no-op。 + DropExpiredPartitions(ctx context.Context, cutoff time.Time) error } +// AccessLogFilter 是查询过滤器的别名,供 apps 使用,避免 import repository/analytics。 +type AccessLogFilter = analyticsrepo.AccessLogFilter + // StatusStore 日志库状态。 type StatusStore interface { ActiveDatabase(ctx context.Context) (string, error) diff --git a/internal/repository/logstore/partition_cleanup.go b/internal/repository/logstore/partition_cleanup.go new file mode 100644 index 00000000..aa9d663f --- /dev/null +++ b/internal/repository/logstore/partition_cleanup.go @@ -0,0 +1,115 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package logstore + +import ( + "context" + "fmt" + "strings" + "time" + + "gorm.io/gorm" +) + +// listPartitionNames 列出 table 在当前 schema 下的全部直接分区表名(pg_inherits)。 +func listPartitionNames(ctx context.Context, gdb *gorm.DB, table string) ([]string, error) { + var names []string + if err := gdb.WithContext(ctx).Raw(` +SELECT c.relname +FROM pg_inherits i +JOIN pg_class c ON c.oid = i.inhrelid +JOIN pg_class p ON p.oid = i.inhparent +JOIN pg_namespace n ON n.oid = p.relnamespace AND n.nspname = current_schema() +WHERE p.relname = ?`, table).Scan(&names).Error; err != nil { + return nil, fmt.Errorf("list partitions of %s: %w", table, err) + } + return names, nil +} + +// partitionNameMonth 解析按月分区表名 _YYYYMM 的所属月份;命名不匹配返回 (零值, false)。 +func partitionNameMonth(table, name string) (time.Time, bool) { + suffix, ok := strings.CutPrefix(name, table+"_") + if !ok || len(suffix) != 6 { + return time.Time{}, false + } + m, err := time.Parse("200601", suffix) + if err != nil { + return time.Time{}, false + } + return m, true +} + +// dropEligiblePartitionNames 返回 before 月份之前、命名合法的分区表名(是否为空由调用方校验)。 +func dropEligiblePartitionNames(table string, names []string, before time.Time) []string { + beforeMonth := time.Date(before.Year(), before.Month(), 1, 0, 0, 0, 0, time.UTC) + out := make([]string, 0, len(names)) + for _, name := range names { + month, ok := partitionNameMonth(table, name) + if !ok || !month.Before(beforeMonth) { + continue + } + out = append(out, name) + } + return out +} + +// DropEmptyPartitions 幂等清理 PG 空分区表:仅删除 before 月份之前、且无任何数据的分区。 +// 非 PG 方言为 no-op。 +func (s *gormLogStore) DropEmptyPartitions(ctx context.Context, before time.Time) error { + if !isPostgresDialect(s.db) { + return nil + } + names, err := listPartitionNames(ctx, s.db, userAccessLogTable) + if err != nil { + return err + } + for _, name := range dropEligiblePartitionNames(userAccessLogTable, names, before) { + var one int + if err := s.db.WithContext(ctx).Raw("SELECT 1 FROM " + name + " LIMIT 1").Scan(&one).Error; err != nil { + return fmt.Errorf("check partition %s empty: %w", name, err) + } + if one == 1 { + continue + } + if err := s.db.WithContext(ctx).Exec("DROP TABLE IF EXISTS " + name).Error; err != nil { + return fmt.Errorf("drop empty partition %s: %w", name, err) + } + } + return nil +} + +// DropExpiredPartitions 直接删除完全过期的 PG 整月分区(避免 retention 清理逐行 DELETE)。 +// 候选 = 月份早于 cutoff 月(按 cutoff 的 UTC 时刻取月);删除前校验分区内不存在 created_at >= cutoff 的行。 +// 迁移冻结期间返回 ErrMigrating。CH/SQLite 为 no-op。 +func (s *gormLogStore) DropExpiredPartitions(ctx context.Context, cutoff time.Time) error { + if !isPostgresDialect(s.db) { + return nil + } + if err := s.ensureWritable(ctx); err != nil { + return err + } + names, err := listPartitionNames(ctx, s.db, userAccessLogTable) + if err != nil { + return err + } + cu := cutoff.UTC() + cutoffMonth := time.Date(cu.Year(), cu.Month(), 1, 0, 0, 0, 0, time.UTC) + for _, name := range names { + month, ok := partitionNameMonth(userAccessLogTable, name) + if !ok || !month.Before(cutoffMonth) { + continue + } + var hasRetained int + if err := s.db.WithContext(ctx).Raw("SELECT 1 FROM "+name+" WHERE created_at >= ? LIMIT 1", cu).Scan(&hasRetained).Error; err != nil { + return fmt.Errorf("check partition %s retained rows: %w", name, err) + } + if hasRetained == 1 { + continue + } + if err := s.db.WithContext(ctx).Exec("DROP TABLE IF EXISTS " + name).Error; err != nil { + return fmt.Errorf("drop expired partition %s: %w", name, err) + } + } + return nil +} diff --git a/internal/repository/logstore/partition_cleanup_test.go b/internal/repository/logstore/partition_cleanup_test.go new file mode 100644 index 00000000..6dfba9d0 --- /dev/null +++ b/internal/repository/logstore/partition_cleanup_test.go @@ -0,0 +1,70 @@ +// 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" + "github.com/stretchr/testify/require" +) + +func TestPartitionNameMonth(t *testing.T) { + cases := []struct { + table string + name string + want string + }{ + {"w_user_access_logs", "w_user_access_logs_202612", "2026-12"}, + {"w_user_access_logs", "w_user_access_logs_202608", "2026-08"}, + {"w_user_access_logs", "of_node_access_logs_202608", ""}, + {"w_user_access_logs", "w_user_access_logs_20268", ""}, + {"w_user_access_logs", "w_user_access_logs_202613", ""}, + {"w_user_access_logs", "w_user_access_logs_default", ""}, + } + for _, c := range cases { + got, ok := partitionNameMonth(c.table, c.name) + if c.want == "" { + if ok { + t.Fatalf("partitionNameMonth(%q, %q) ok = true, want false", c.table, c.name) + } + continue + } + if !ok || got.Format("2006-01") != c.want { + t.Fatalf("partitionNameMonth(%q, %q) = %v, want %s", c.table, c.name, got, c.want) + } + } +} + +func TestDropEligiblePartitionNames(t *testing.T) { + before := time.Date(2026, 10, 15, 0, 0, 0, 0, time.UTC) + names := []string{ + "w_user_access_logs_202608", + "w_user_access_logs_202609", + "w_user_access_logs_202610", + "w_user_access_logs_202611", + "w_user_access_logs_default", + } + got := dropEligiblePartitionNames(userAccessLogTable, names, before) + want := []string{"w_user_access_logs_202608", "w_user_access_logs_202609"} + require.Equal(t, want, got) + + first := time.Date(2026, 10, 1, 0, 0, 0, 0, time.UTC) + require.Empty(t, dropEligiblePartitionNames(userAccessLogTable, []string{"w_user_access_logs_202610"}, first)) +} + +func TestDropPartitionHelpersSQLiteNoop(t *testing.T) { + ua := newTestUserAccessStore(t) + ctx := context.Background() + require.NoError(t, ua.BatchInsert(ctx, []analyticsmodel.UserAccessLog{ + {UserID: 1, Path: "/x", CreatedAt: time.Now().UTC()}, + })) + require.NoError(t, ua.DropExpiredPartitions(ctx, time.Now().AddDate(0, 0, -90))) + require.NoError(t, ua.DropEmptyPartitions(ctx, time.Now())) + count, err := ua.Count(ctx, AccessLogFilter{}) + require.NoError(t, err) + require.Equal(t, uint64(1), count) +}