feat(log): PG 分区清理与 logstore import-lint

CleanupExpired 先按月 DROP 过期分区,再删边界行并清理空分区;apps 禁止直连 analytics。
This commit is contained in:
ryan
2026-08-16 16:48:15 +08:00
parent a8fcf6087a
commit b66cf3ae9c
9 changed files with 259 additions and 6 deletions
+2 -2
View File
@@ -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`)。
## 禁止
+1 -1
View File
@@ -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` 受保护,管理端不可改。
+2 -3
View File
@@ -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 != "" {
+7
View File
@@ -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
}
@@ -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")
@@ -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)
}
}
}
}
}
+9
View File
@@ -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)
@@ -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 解析按月分区表名 <table>_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
}
@@ -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)
}