From f530cd4025e7330dc4a8cd0e503fdef4d0c54957 Mon Sep 17 00:00:00 2001 From: ryan Date: Sun, 9 Aug 2026 10:20:58 +0800 Subject: [PATCH] perf(log): optimize PG log store queries and expired partition cleanup MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit - Count/节点访问日志统计改为单次扫描聚合,WAF 按 IP 聚合由三次扫描合并为两次 - IPSummaries 归属地改为取过滤窗口内最新记录(对齐 ClickHouse argMax 口径), 子查询带窗口条件,可分区裁剪并命中索引 - 新增 goose 迁移:of_node_access_logs (logged_at DESC, id DESC) 前导索引与 lower(trim(host)) 表达式索引,加速列表排序与主机过滤 - 过期日志清理先按数据校验直接 DROP 完全过期整月分区,再对边界月逐行删除; 启动时兜底预建当月及未来 2 个月分区,跨月停机重启后首次写入不再报 "no partition of relation found" --- docs/changelog/index.md | 1 + .../postgres/202608090003_add_log_indexes.sql | 9 + .../sqlite/202608090003_add_log_indexes.sql | 9 + internal/repository/logstore/cleanup.go | 6 + internal/repository/logstore/cleanup_test.go | 32 +++ .../repository/logstore/clickhouse_store.go | 5 + internal/repository/logstore/logstore.go | 3 + .../repository/logstore/partition_cleanup.go | 66 ++++++ .../postgres_partition_integration_test.go | 214 ++++++++++++++++++ .../repository/logstore/postgres_store.go | 175 +++++++------- .../logstore/postgres_store_test.go | 70 ++++++ internal/repository/logstore/provider.go | 13 +- 12 files changed, 519 insertions(+), 84 deletions(-) create mode 100644 internal/infra/persistence/migrator/goose/postgres/202608090003_add_log_indexes.sql create mode 100644 internal/infra/persistence/migrator/goose/sqlite/202608090003_add_log_indexes.sql create mode 100644 internal/repository/logstore/partition_cleanup.go diff --git a/docs/changelog/index.md b/docs/changelog/index.md index 730b8615..b1f8b139 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -33,6 +33,7 @@ sidebar: false - 性能指标(CPU/内存/磁盘/网络)与访问日志的保留时长解耦:新增三库共用的 `metric_retention_days` 配置(默认 3 天),系统垃圾清理每日任务按独立短留存清理指标快照;访问日志仍按 `log_retention_days_*` 清理。 - ClickHouse 改为默认关闭:`clickhouse.enabled` 缺省或为 `false` 时不启用(此前会被强制置为 `true`),日志/指标由 PostgreSQL/SQLite 主库承担;显式 `true` 或设置 `CLICKHOUSE_HOST` / `CLICKHOUSE_ENABLED=true` 时启用。 - 系统垃圾清理任务在删除过期日志后,会同时清理旧月份空分区表:PostgreSQL 按月分区的访问日志表(节点/用户)在数据删除后若该月分区已无数据,则自动删除对应分区表,避免历史分区表无限累积;仅删除「当前月之前」且为空的月份分区,当月/未来月及仍有数据的分区保留。 +- 日志存储(PostgreSQL/SQLite 实现)性能优化:节点访问日志统计查询合并为单次扫描;IP 汇总的归属地改为取查询窗口内最新记录(与 ClickHouse 口径一致)并避免跨全部分区扫描;WAF 按 IP 聚合由三次扫描合并为两次;新增 `logged_at` 前导索引与主机名小写表达式索引,加速日志列表排序与主机过滤;过期日志清理改为按数据校验后直接删除完全过期的整月分区(比逐行删除快,且不会误删保留期内数据),并新增启动时分区兜底,跨月停机重启后首次写入不再报「no partition of relation found」。 ### 修复 diff --git a/internal/infra/persistence/migrator/goose/postgres/202608090003_add_log_indexes.sql b/internal/infra/persistence/migrator/goose/postgres/202608090003_add_log_indexes.sql new file mode 100644 index 00000000..4e257530 --- /dev/null +++ b/internal/infra/persistence/migrator/goose/postgres/202608090003_add_log_indexes.sql @@ -0,0 +1,9 @@ +-- +goose Up +-- 日志列表默认排序与 retention 清理的时间范围条件:补 logged_at 前导索引。 +CREATE INDEX IF NOT EXISTS idx_of_node_access_logs_logged_at ON of_node_access_logs (logged_at DESC, id DESC); +-- hosts 过滤为 lower(trim(host)) IN:补表达式索引,普通 (host,...) 索引无法命中函数表达式。 +CREATE INDEX IF NOT EXISTS idx_of_node_access_logs_host_lower ON of_node_access_logs (lower(trim(host))); + +-- +goose Down +DROP INDEX IF EXISTS idx_of_node_access_logs_host_lower; +DROP INDEX IF EXISTS idx_of_node_access_logs_logged_at; diff --git a/internal/infra/persistence/migrator/goose/sqlite/202608090003_add_log_indexes.sql b/internal/infra/persistence/migrator/goose/sqlite/202608090003_add_log_indexes.sql new file mode 100644 index 00000000..4e257530 --- /dev/null +++ b/internal/infra/persistence/migrator/goose/sqlite/202608090003_add_log_indexes.sql @@ -0,0 +1,9 @@ +-- +goose Up +-- 日志列表默认排序与 retention 清理的时间范围条件:补 logged_at 前导索引。 +CREATE INDEX IF NOT EXISTS idx_of_node_access_logs_logged_at ON of_node_access_logs (logged_at DESC, id DESC); +-- hosts 过滤为 lower(trim(host)) IN:补表达式索引,普通 (host,...) 索引无法命中函数表达式。 +CREATE INDEX IF NOT EXISTS idx_of_node_access_logs_host_lower ON of_node_access_logs (lower(trim(host))); + +-- +goose Down +DROP INDEX IF EXISTS idx_of_node_access_logs_host_lower; +DROP INDEX IF EXISTS idx_of_node_access_logs_logged_at; diff --git a/internal/repository/logstore/cleanup.go b/internal/repository/logstore/cleanup.go index d74e2b52..50c34e5d 100644 --- a/internal/repository/logstore/cleanup.go +++ b/internal/repository/logstore/cleanup.go @@ -107,6 +107,12 @@ func CleanupExpired(ctx context.Context) (*CleanupSummary, error) { return nil, fmt.Errorf("ensure partitions: %w", err) } + // 先直接删除完全过期的整月分区(比逐行 DELETE 快几个数量级、无 MVCC/WAL 负担), + // 再对边界月份执行 DeleteBefore(边界月仍可能含未过期数据,不可整表删)。 + if err := s.AccessLogs.DropExpiredPartitions(ctx, cutoff); err != nil { + return nil, fmt.Errorf("drop expired partitions: %w", err) + } + if err := cleanupTable("node_access_logs", func() (int64, error) { return s.AccessLogs.DeleteBefore(ctx, cutoff) }, summary); err != nil { diff --git a/internal/repository/logstore/cleanup_test.go b/internal/repository/logstore/cleanup_test.go index c2dab959..be80c49e 100644 --- a/internal/repository/logstore/cleanup_test.go +++ b/internal/repository/logstore/cleanup_test.go @@ -369,3 +369,35 @@ func TestDropEligiblePartitionNames(t *testing.T) { t.Fatalf("eligible at month boundary = %v, want empty", got) } } + +// TestDropExpiredPartitionsSQLiteNoop 验证 SQLite 下 DropExpiredPartitions 为 no-op: +// 直接返回 nil、不触碰任何分区 SQL(SQLite 无分区),数据不受影响。 +func TestDropExpiredPartitionsSQLiteNoop(t *testing.T) { + ResetForTest() + SetConfigReader(func(_ context.Context, key string) (string, error) { + if key == logDatabaseKey { + return "sqlite", nil + } + return "", nil + }) + defer ResetForTest() + + gdb := newCleanupTestDB(t) + ctx := context.Background() + if err := gdb.Create(&analyticsmodel.NodeAccessLog{ID: 1, NodeID: "n1", LoggedAt: time.Now().AddDate(0, 0, -100).UTC(), RemoteAddr: "1.1.1.1"}).Error; err != nil { + t.Fatalf("seed node access log: %v", err) + } + + store := newGormStore(gdb) + if err := store.DropExpiredPartitions(ctx, time.Now().AddDate(0, 0, -90)); err != nil { + t.Fatalf("DropExpiredPartitions on sqlite: %v", err) + } + + var n int64 + if err := gdb.Model(&analyticsmodel.NodeAccessLog{}).Count(&n).Error; err != nil { + t.Fatalf("count node access logs: %v", err) + } + if n != 1 { + t.Fatalf("node access log count = %d, want 1(no-op 不应删除任何行)", n) + } +} diff --git a/internal/repository/logstore/clickhouse_store.go b/internal/repository/logstore/clickhouse_store.go index 62794ce6..e7ef4cca 100644 --- a/internal/repository/logstore/clickhouse_store.go +++ b/internal/repository/logstore/clickhouse_store.go @@ -439,6 +439,11 @@ func (s *clickhouseLogStore) DropEmptyPartitions(_ context.Context, _ time.Time) return nil } +// DropExpiredPartitions 是 CH 分支 no-op(CH 无 PG 式分区,retention 仍走 DeleteBefore)。 +func (s *clickhouseLogStore) DropExpiredPartitions(_ context.Context, _ time.Time) error { + return nil +} + // chMigrationRange 查询 CH 表时间列 MIN/MAX;空表(NULL)返回零值。 func chMigrationRange(ctx context.Context, table, column string) (time.Time, time.Time, error) { if err := chConnErr(); err != nil { diff --git a/internal/repository/logstore/logstore.go b/internal/repository/logstore/logstore.go index 72369dba..ceac0ba3 100644 --- a/internal/repository/logstore/logstore.go +++ b/internal/repository/logstore/logstore.go @@ -51,6 +51,9 @@ type AccessLogStore interface { // DropEmptyPartitions 幂等清理 PG 空分区表:删除 before 月份之前、且无任何数据的按月分区; // CH/SQLite 为 no-op(CH 分区随数据删除自动消失、SQLite 无分区)。 DropEmptyPartitions(ctx context.Context, before time.Time) error + // DropExpiredPartitions 直接删除完全过期的 PG 整月分区(候选为月份早于 cutoff 月的分区, + // 删除前校验分区内无保留期内数据,避免时区偏移下误删;迁移冻结期间拒绝执行);CH/SQLite 为 no-op。 + DropExpiredPartitions(ctx context.Context, cutoff time.Time) error } // ObservabilityStore 可观测 4 表(metric snapshots / edge health / frps / frpc)。 diff --git a/internal/repository/logstore/partition_cleanup.go b/internal/repository/logstore/partition_cleanup.go new file mode 100644 index 00000000..5df40b02 --- /dev/null +++ b/internal/repository/logstore/partition_cleanup.go @@ -0,0 +1,66 @@ +// Copyright 2026 Arctel.net +// SPDX-License-Identifier: Apache-2.0 + +package logstore + +import ( + "context" + "fmt" + "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 +} + +// DropExpiredPartitions 直接删除完全过期的 PG 整月分区(避免 retention 清理逐行 DELETE): +// 候选 = 月份早于 cutoff 月(按 cutoff 的 UTC 时刻取月,避免本地时区偏移超前误删)的分区, +// 且删除前校验分区内不存在 logged_at >= cutoff 的行(分区边界随会话时区偏移, +// 名称月份只能粗筛,必须以数据为准);仅处理 of_node_access_logs +// (w_user_access_logs 无 retention 清理,刻意不删其分区);迁移冻结期间(ensureWritable) +// 直接返回 ErrMigrating,避免对冻结源库整月 DROP 丢数据;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, "of_node_access_logs") + 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("of_node_access_logs", name) + if !ok || !month.Before(cutoffMonth) { + continue // 非法命名或当月/未来月分区,必须保留 + } + // 数据校验:分区内仍有 logged_at >= cutoff 的行则保留(时区偏移下名称月份可能超前于真实边界)。 + var hasRetained int + if err := s.db.WithContext(ctx).Raw("SELECT 1 FROM "+name+" WHERE logged_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/postgres_partition_integration_test.go b/internal/repository/logstore/postgres_partition_integration_test.go index 16136ea1..43a4ba0e 100644 --- a/internal/repository/logstore/postgres_partition_integration_test.go +++ b/internal/repository/logstore/postgres_partition_integration_test.go @@ -245,6 +245,220 @@ func TestDropEmptyPartitionsPostgres(t *testing.T) { } } +// TestDropExpiredPartitionsPostgres 需要 TEST_POSTGRES_DSN(未设置时跳过): +// 验证直接删除完全早于 cutoff 月份的整月分区:早于 cutoff 月的分区(含其中全部数据)被整表 DROP、 +// 边界月分区保留且数据仍在;重复调用幂等;w_user_access_logs 分区不受影响(无 retention 清理)。 +func TestDropExpiredPartitionsPostgres(t *testing.T) { + dsn := strings.TrimSpace(os.Getenv("TEST_POSTGRES_DSN")) + if dsn == "" { + t.Skip("TEST_POSTGRES_DSN is not set") + } + + gdb, err := gorm.Open(postgres.Open(dsn), &gorm.Config{ + DisableForeignKeyConstraintWhenMigrating: true, + Logger: logger.Default.LogMode(logger.Silent), + }) + if err != nil { + t.Fatalf("open postgres: %v", err) + } + sqlDB, err := gdb.DB() + if err != nil { + t.Fatalf("sql db: %v", err) + } + sqlDB.SetMaxOpenConns(1) + + schema := fmt.Sprintf("logstore_drop_expired_%d", time.Now().UnixNano()) + if !regexp.MustCompile(`^[a-z0-9_]+$`).MatchString(schema) { + t.Fatalf("invalid schema: %s", schema) + } + if err := gdb.Exec(`CREATE SCHEMA "` + schema + `"`).Error; err != nil { + t.Fatalf("create schema: %v", err) + } + if err := gdb.Exec(`SET search_path TO "` + schema + `"`).Error; err != nil { + t.Fatalf("set search_path: %v", err) + } + t.Cleanup(func() { + _ = gdb.Exec("SET search_path TO public").Error + _ = gdb.Exec(`DROP SCHEMA IF EXISTS "` + schema + `" CASCADE`).Error + _ = sqlDB.Close() + }) + + for _, ddl := range []string{postgresNodeAccessLogsDDL, postgresUserAccessLogsDDL} { + if err := gdb.Exec(ddl).Error; err != nil { + t.Fatalf("create partitioned table: %v", err) + } + } + + ctx := context.Background() + store := newGormStore(gdb) + ua := newUserAccessLogGormStore(gdb) + + // 预建 202601..202604 分区;1/3 月有数据、2/4 月为空。 + if err := store.EnsurePartitions(ctx, + time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC), + time.Date(2026, 4, 20, 0, 0, 0, 0, time.UTC)); err != nil { + t.Fatalf("EnsurePartitions: %v", err) + } + if err := store.BatchInsertNodeAccessLogs(ctx, []analyticsmodel.NodeAccessLog{ + {ID: 1, NodeID: "n1", LoggedAt: time.Date(2026, 1, 15, 0, 0, 0, 0, time.UTC), RemoteAddr: "1.1.1.1"}, + {ID: 2, NodeID: "n1", LoggedAt: time.Date(2026, 1, 20, 0, 0, 0, 0, time.UTC), RemoteAddr: "1.1.1.2"}, + {ID: 3, NodeID: "n2", LoggedAt: time.Date(2026, 3, 5, 0, 0, 0, 0, time.UTC), RemoteAddr: "3.3.3.3"}, + {ID: 4, NodeID: "n2", LoggedAt: time.Date(2026, 3, 18, 0, 0, 0, 0, time.UTC), RemoteAddr: "3.3.3.4"}, + }); err != nil { + t.Fatalf("insert node access logs: %v", err) + } + if err := ua.BatchInsert(ctx, []analyticsmodel.UserAccessLog{ + {ID: 1, UserID: 101, Path: "/a", CreatedAt: time.Date(2026, 1, 16, 0, 0, 0, 0, time.UTC)}, + }); err != nil { + t.Fatalf("insert user access log: %v", err) + } + + // cutoff=2026-03-10:分区月份早于 2026-03 的(202601、202602)整表 DROP; + // 202603(边界月,可能含未过期数据)与 202604(未来月)保留。 + if err := store.DropExpiredPartitions(ctx, time.Date(2026, 3, 10, 0, 0, 0, 0, time.UTC)); err != nil { + t.Fatalf("DropExpiredPartitions: %v", err) + } + // 幂等:重复调用不报错、不额外删除。 + if err := store.DropExpiredPartitions(ctx, time.Date(2026, 3, 10, 0, 0, 0, 0, time.UTC)); err != nil { + t.Fatalf("DropExpiredPartitions idempotent: %v", err) + } + + assertPartitions := func(parent string, want int64) { + t.Helper() + var n int64 + if err := gdb.Raw( + "SELECT count(*) FROM pg_inherits WHERE inhparent = to_regclass(?)", + parent, + ).Scan(&n).Error; err != nil { + t.Fatalf("count partitions of %s: %v", parent, err) + } + if n != want { + t.Fatalf("%s partitions = %d, want %d", parent, n, want) + } + } + // of_node_access_logs 只剩边界月+未来月 2 个分区;w_user_access_logs 不受影响(仍 4 个)。 + assertPartitions("of_node_access_logs", 2) + assertPartitions("w_user_access_logs", 4) + + // 202601/202602 分区被整表 DROP:1 月数据随之消失,3 月数据保留。 + var nodeCount int64 + if err := gdb.Model(&analyticsmodel.NodeAccessLog{}).Count(&nodeCount).Error; err != nil { + t.Fatalf("count node access logs: %v", err) + } + if nodeCount != 2 { + t.Fatalf("node access log count = %d, want 2(仅剩 3 月数据)", nodeCount) + } +} + +// TestDropExpiredPartitionsTimezoneSafety 需要 TEST_POSTGRES_DSN(未设置时跳过): +// 覆盖本地时区偏移下 DropExpiredPartitions 的时区安全性:cutoff 为 UTC+8 本地时刻 +// (其实刻 = 2026-02-28T21:00Z),名称月份早于 cutoff 月但分区内仍含保留期行的 +// 202602 不得被误删(旧实现按本地月份取 cutoffMonth=2026-03 会整表 DROP 丢数据); +// 完全过期的 202601 正常整表 DROP;保留期行仍可查询到。 +func TestDropExpiredPartitionsTimezoneSafety(t *testing.T) { + dsn := strings.TrimSpace(os.Getenv("TEST_POSTGRES_DSN")) + if dsn == "" { + t.Skip("TEST_POSTGRES_DSN is not set") + } + + gdb, err := gorm.Open(postgres.Open(dsn), &gorm.Config{ + DisableForeignKeyConstraintWhenMigrating: true, + Logger: logger.Default.LogMode(logger.Silent), + }) + if err != nil { + t.Fatalf("open postgres: %v", err) + } + sqlDB, err := gdb.DB() + if err != nil { + t.Fatalf("sql db: %v", err) + } + sqlDB.SetMaxOpenConns(1) + + schema := fmt.Sprintf("logstore_drop_expired_tz_%d", time.Now().UnixNano()) + if !regexp.MustCompile(`^[a-z0-9_]+$`).MatchString(schema) { + t.Fatalf("invalid schema: %s", schema) + } + if err := gdb.Exec(`CREATE SCHEMA "` + schema + `"`).Error; err != nil { + t.Fatalf("create schema: %v", err) + } + if err := gdb.Exec(`SET search_path TO "` + schema + `"`).Error; err != nil { + t.Fatalf("set search_path: %v", err) + } + t.Cleanup(func() { + _ = gdb.Exec("SET search_path TO public").Error + _ = gdb.Exec(`DROP SCHEMA IF EXISTS "` + schema + `" CASCADE`).Error + _ = sqlDB.Close() + }) + + for _, ddl := range []string{postgresNodeAccessLogsDDL, postgresUserAccessLogsDDL} { + if err := gdb.Exec(ddl).Error; err != nil { + t.Fatalf("create partitioned table: %v", err) + } + } + + ctx := context.Background() + store := newGormStore(gdb) + + // 预建 202601..202602 分区。 + if err := store.EnsurePartitions(ctx, + time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC), + time.Date(2026, 2, 20, 0, 0, 0, 0, time.UTC)); err != nil { + t.Fatalf("EnsurePartitions: %v", err) + } + // 202601 仅含完全过期行;202602 含一条过期行(2026-02-10)与一条保留期行 + // (2026-02-28T21:00Z,恰等于 cutoff 其实刻,>= 语义下必须保留)。 + if err := store.BatchInsertNodeAccessLogs(ctx, []analyticsmodel.NodeAccessLog{ + {ID: 1, NodeID: "n1", LoggedAt: time.Date(2026, 1, 15, 0, 0, 0, 0, time.UTC), RemoteAddr: "1.1.1.1"}, + {ID: 2, NodeID: "n1", LoggedAt: time.Date(2026, 2, 10, 0, 0, 0, 0, time.UTC), RemoteAddr: "2.2.2.2"}, + {ID: 3, NodeID: "n1", LoggedAt: time.Date(2026, 2, 28, 21, 0, 0, 0, time.UTC), RemoteAddr: "3.3.3.3"}, + }); err != nil { + t.Fatalf("insert node access logs: %v", err) + } + + // cutoff 为 UTC+8 本地时刻 2026-03-01 05:00,其实刻 = 2026-02-28T21:00Z: + // 旧实现按本地月份取 cutoffMonth=2026-03 会把 202602 误判为完全过期整表 DROP。 + cutoff := time.Date(2026, 3, 1, 5, 0, 0, 0, time.FixedZone("UTC+8", 8*3600)) + if err := store.DropExpiredPartitions(ctx, cutoff); err != nil { + t.Fatalf("DropExpiredPartitions: %v", err) + } + + assertPartitions := func(parent string, want int64) { + t.Helper() + var n int64 + if err := gdb.Raw( + "SELECT count(*) FROM pg_inherits WHERE inhparent = to_regclass(?)", + parent, + ).Scan(&n).Error; err != nil { + t.Fatalf("count partitions of %s: %v", parent, err) + } + if n != want { + t.Fatalf("%s partitions = %d, want %d", parent, n, want) + } + } + // 202601 已整表 DROP,202602 保留;w_user_access_logs 不受影响(仍 2 个)。 + assertPartitions("of_node_access_logs", 1) + assertPartitions("w_user_access_logs", 2) + + // 202601 数据随之消失,202602 内保留期行(2026-02-28T21:00Z)仍可查询到。 + var nodeCount int64 + if err := gdb.Model(&analyticsmodel.NodeAccessLog{}).Count(&nodeCount).Error; err != nil { + t.Fatalf("count node access logs: %v", err) + } + if nodeCount != 2 { + t.Fatalf("node access log count = %d, want 2(仅剩 202602 两行)", nodeCount) + } + var retained int64 + if err := gdb.Raw( + "SELECT count(*) FROM of_node_access_logs WHERE logged_at >= ?", + cutoff.UTC(), + ).Scan(&retained).Error; err != nil { + t.Fatalf("count retained rows: %v", err) + } + if retained != 1 { + t.Fatalf("retained rows (logged_at >= cutoff) = %d, want 1", retained) + } +} + // postgresNodeAccessLogsDDL 与 goose/postgres/202608080001_create_log_tables.sql 对齐。 const postgresNodeAccessLogsDDL = ` CREATE TABLE IF NOT EXISTS of_node_access_logs ( diff --git a/internal/repository/logstore/postgres_store.go b/internal/repository/logstore/postgres_store.go index a4811d07..fbed86fc 100644 --- a/internal/repository/logstore/postgres_store.go +++ b/internal/repository/logstore/postgres_store.go @@ -171,22 +171,21 @@ func nodeAccessLogOrderClauseGORM(sortBy, sortOrder string) (string, error) { func (s *gormLogStore) Count(ctx context.Context, query model.OpenFlareAccessLogQuery) (int64, int64, int64, error) { f := toNodeAccessLogFilter(query) - var total, uniqIP, bytesSent int64 + // 对齐 CH CountNodeAccessLogs:单次扫描聚合 total/uniq IP/bytes_sent; + // distinct IP 排除空 remote_addr(uniqExactIf(remote_addr, remote_addr != ''), + // 由 distinctNonEmptyCountSQL 按方言处理 PG FILTER / SQLite CASE 差异)。 + type row struct { + Total int64 + UniqIP int64 + BytesSent int64 + } + var out row q := applyNodeAccessLogFilter(s.db.WithContext(ctx).Model(&analyticsmodel.NodeAccessLog{}), f) - if err := q.Count(&total).Error; err != nil { + q = q.Select("COUNT(*) AS total, " + distinctNonEmptyCountSQL(s.db, "remote_addr") + " AS uniq_ip, COALESCE(SUM(bytes_sent),0) AS bytes_sent") + if err := q.Scan(&out).Error; err != nil { return 0, 0, 0, err } - // 对齐 CH CountNodeAccessLogs:distinct IP 排除空 remote_addr(uniqExactIf(remote_addr, remote_addr != ''))。 - // 独立查询链,避免 remote_addr <> '' 条件泄漏到下面的 bytes_sent 求和。 - uniqQ := applyNodeAccessLogFilter(s.db.WithContext(ctx).Model(&analyticsmodel.NodeAccessLog{}), f) - uniqQ = uniqQ.Where("remote_addr <> ''").Distinct("remote_addr") - if err := uniqQ.Count(&uniqIP).Error; err != nil { - return 0, 0, 0, err - } - if err := q.Select("COALESCE(SUM(bytes_sent),0)").Scan(&bytesSent).Error; err != nil { - return 0, 0, 0, err - } - return total, uniqIP, bytesSent, nil + return out.Total, out.UniqIP, out.BytesSent, nil } func (s *gormLogStore) TrafficSummary(ctx context.Context, query model.OpenFlareAccessLogQuery) (model.OpenFlareAccessLogTrafficSummary, error) { @@ -417,7 +416,8 @@ func (s *gormLogStore) IPAggregates(ctx context.Context, query model.OpenFlareAc return out, nil } -// IPSummaries 按 IP 汇总(region 取该 IP 最近一条日志的 region;recent_requests 恒 0)。 +// IPSummaries 按 IP 汇总(region 取过滤窗口内该 IP 最近一条日志的 region,对齐 CH +// argMax(region, logged_at);recent_requests 恒 0)。 func (s *gormLogStore) IPSummaries(ctx context.Context, query model.OpenFlareAccessLogQuery, _ time.Time) ([]analyticsmodel.NodeAccessLogIPSummary, error) { f := toNodeAccessLogFilter(query) type row struct { @@ -432,15 +432,31 @@ func (s *gormLogStore) IPSummaries(ctx context.Context, query model.OpenFlareAcc } var rows []row q := applyNodeAccessLogFilter(s.db.WithContext(ctx).Model(&analyticsmodel.NodeAccessLog{}), f) - q = q.Select(` + // region 子查询携带与外层完全相同的过滤条件(复用 buildNodeAccessLogFilterParts), + // 只取窗口内该 IP 最近一条:避免每个 IP 组全分区扫最新(可命中 (remote_addr, logged_at DESC) + // 索引并分区裁剪);cond 为空时省略 AND。 + cond, condArgs := buildNodeAccessLogFilterParts(f) + regionExpr := `(SELECT t2.region FROM of_node_access_logs t2 WHERE t2.remote_addr = of_node_access_logs.remote_addr` + if cond != "" { + regionExpr += " AND " + cond + } + regionExpr += ` ORDER BY t2.logged_at DESC LIMIT 1)` + selectStr := ` remote_addr, - (SELECT t2.region FROM of_node_access_logs t2 WHERE t2.remote_addr = of_node_access_logs.remote_addr ORDER BY t2.logged_at DESC LIMIT 1) AS region, + ` + regionExpr + ` AS region, COUNT(*) AS total_requests, COUNT(*) FILTER (WHERE status_code >= 200 AND status_code < 300) AS success2xx_count, CASE WHEN COUNT(*) = 0 THEN 0.0 ELSE CAST(COUNT(*) FILTER (WHERE status_code >= 200 AND status_code < 300) AS REAL) / CAST(COUNT(*) AS REAL) END AS success_ratio, COALESCE(SUM(request_length),0) AS request_length, COALESCE(SUM(bytes_sent),0) AS bytes_sent, - MAX(` + epochSQL(s.db, "logged_at") + `) AS last_seen_epoch`) + MAX(` + epochSQL(s.db, "logged_at") + `) AS last_seen_epoch` + if len(condArgs) > 0 { + // GORM Select(query, args...) 在字符串含恰好 len(args) 个 "?" 时把 args 作为 SELECT 子句 + // 参数(绑定顺序在 WHERE 参数之前),与外层 applyNodeAccessLogFilter 的同一批参数不会错位。 + q = q.Select(selectStr, condArgs...) + } else { + q = q.Select(selectStr) + } q = q.Where("remote_addr != ''") q = q.Group("remote_addr").Order("total_requests DESC, last_seen_epoch DESC, remote_addr ASC") // 对齐 CH IPSummariesNodeAccessLogs 的 0-based 分页(仅 PageSize>0 时分页)。 @@ -519,10 +535,7 @@ func (s *gormLogStore) WAFIPAggregates(ctx context.Context, query model.OpenFlar order = append(order, remoteAddr) } if len(aggregates) > 0 { - if err := s.mergeWAFIPStatusCounts(ctx, f, aggregates); err != nil { - return nil, err - } - if err := s.mergeWAFIPHostCounts(ctx, f, aggregates); err != nil { + if err := s.mergeWAFIPStatusAndHostCounts(ctx, f, aggregates); err != nil { return nil, err } } @@ -535,18 +548,22 @@ func (s *gormLogStore) WAFIPAggregates(ctx context.Context, query model.OpenFlar return result, nil } -// mergeWAFIPStatusCounts 填充 WAF 每 IP 状态码分布。 -func (s *gormLogStore) mergeWAFIPStatusCounts(ctx context.Context, f analyticsmodel.NodeAccessLogFilter, aggregates map[string]*analyticsmodel.NodeAccessLogWAFIPAggregate) error { +// mergeWAFIPStatusAndHostCounts 一次扫描填充每 IP 状态码分布与 IP 字面量 host 计数 +// (GROUP BY remote_addr, status_code, host;CH 对应 countIf(hostIsIP) 折入主查询 + 状态码二次聚合)。 +// 不筛 host 非空以保持状态计数口径(含空 host 行);isIPLiteralHost 对空 host 返回 false, +// 故空 host 不计入 IPHostCount,与旧 mergeWAFIPHostCounts 的 host 非空过滤语义一致。 +func (s *gormLogStore) mergeWAFIPStatusAndHostCounts(ctx context.Context, f analyticsmodel.NodeAccessLogFilter, aggregates map[string]*analyticsmodel.NodeAccessLogWAFIPAggregate) error { type row struct { - RemoteAddr string - StatusCode int32 - StatusCount int64 + RemoteAddr string + StatusCode int32 + Host string + RowCount int64 } var rows []row q := applyNodeAccessLogFilter(s.db.WithContext(ctx).Model(&analyticsmodel.NodeAccessLog{}), f) - q = q.Select("remote_addr, status_code, COUNT(*) AS status_count"). + q = q.Select("remote_addr, status_code, host, COUNT(*) AS row_count"). Where("remote_addr != ''") - if err := q.Group("remote_addr, status_code").Scan(&rows).Error; err != nil { + if err := q.Group("remote_addr, status_code, host").Scan(&rows).Error; err != nil { return err } for _, r := range rows { @@ -557,30 +574,7 @@ func (s *gormLogStore) mergeWAFIPStatusCounts(ctx context.Context, f analyticsmo if a.StatusCounts == nil { a.StatusCounts = make(map[int]int64) } - a.StatusCounts[int(r.StatusCode)] += r.StatusCount - } - return nil -} - -// mergeWAFIPHostCounts 按 (remote_addr, host) 行数累加 IP 字面量 host 的访问行数。 -func (s *gormLogStore) mergeWAFIPHostCounts(ctx context.Context, f analyticsmodel.NodeAccessLogFilter, aggregates map[string]*analyticsmodel.NodeAccessLogWAFIPAggregate) error { - type row struct { - RemoteAddr string - Host string - RowCount int64 - } - var rows []row - q := applyNodeAccessLogFilter(s.db.WithContext(ctx).Model(&analyticsmodel.NodeAccessLog{}), f) - q = q.Select("remote_addr, host, COUNT(*) AS row_count"). - Where("remote_addr != '' AND trim(host) != ''") - if err := q.Group("remote_addr, host").Scan(&rows).Error; err != nil { - return err - } - for _, r := range rows { - a := aggregates[strings.TrimSpace(r.RemoteAddr)] - if a == nil { - continue - } + a.StatusCounts[int(r.StatusCode)] += r.RowCount if isIPLiteralHost(r.Host) { a.IPHostCount += r.RowCount } @@ -675,15 +669,9 @@ func (s *gormLogStore) DropEmptyPartitions(ctx context.Context, before time.Time return nil } for _, table := range accessLogPartitionTables { - var names []string - if err := s.db.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 fmt.Errorf("list partitions of %s: %w", table, err) + names, err := listPartitionNames(ctx, s.db, table) + if err != nil { + return err } for _, name := range dropEligiblePartitionNames(table, names, before) { var one int @@ -801,32 +789,55 @@ func toNodeAccessLogFilter(query model.OpenFlareAccessLogQuery) analyticsmodel.N } } +// buildNodeAccessLogFilterParts 拼装节点访问日志过滤条件(node_id 等值、remote_addr/host/path +// 前缀 LIKE、hosts 走 lower(trim(host)) IN、since/until 时间窗,顺序与 applyNodeAccessLogFilter +// 完全一致),返回 WHERE 片段与参数;无任何条件时返回空串与 nil。 +func buildNodeAccessLogFilterParts(f analyticsmodel.NodeAccessLogFilter) (string, []any) { + var parts []string + var args []any + if nodeID := strings.TrimSpace(f.NodeID); nodeID != "" { + parts = append(parts, "node_id = ?") + args = append(args, nodeID) + } + if remoteAddr := strings.TrimSpace(f.RemoteAddr); remoteAddr != "" { + parts = append(parts, "remote_addr LIKE ?") + args = append(args, remoteAddr+"%") + } + hosts := normalizeNodeAccessLogHosts(f.Hosts) + if len(hosts) > 0 { + parts = append(parts, "lower(trim(host)) IN ?") + args = append(args, hosts) + } else if host := strings.TrimSpace(f.Host); host != "" { + parts = append(parts, "host LIKE ?") + args = append(args, host+"%") + } + if path := strings.TrimSpace(f.Path); path != "" { + parts = append(parts, "path LIKE ?") + args = append(args, path+"%") + } + if !f.Since.IsZero() { + parts = append(parts, "logged_at >= ?") + args = append(args, f.Since) + } + if !f.Until.IsZero() { + parts = append(parts, "logged_at < ?") + args = append(args, f.Until) + } + if len(parts) == 0 { + return "", nil + } + return strings.Join(parts, " AND "), args +} + // applyNodeAccessLogFilter 对齐 CH 过滤语义(node_access_log_filter.go): // node_id trim 后等值;remote_addr/host/path 前缀 LIKE;hosts 走 lower(trim(host)) IN(参数已归一化); // since 闭区间 >=;until 开区间 <。 func applyNodeAccessLogFilter(q *gorm.DB, f analyticsmodel.NodeAccessLogFilter) *gorm.DB { - if nodeID := strings.TrimSpace(f.NodeID); nodeID != "" { - q = q.Where("node_id = ?", nodeID) + cond, args := buildNodeAccessLogFilterParts(f) + if cond == "" { + return q } - if remoteAddr := strings.TrimSpace(f.RemoteAddr); remoteAddr != "" { - q = q.Where("remote_addr LIKE ?", remoteAddr+"%") - } - hosts := normalizeNodeAccessLogHosts(f.Hosts) - if len(hosts) > 0 { - q = q.Where("lower(trim(host)) IN ?", hosts) - } else if host := strings.TrimSpace(f.Host); host != "" { - q = q.Where("host LIKE ?", host+"%") - } - if path := strings.TrimSpace(f.Path); path != "" { - q = q.Where("path LIKE ?", path+"%") - } - if !f.Since.IsZero() { - q = q.Where("logged_at >= ?", f.Since) - } - if !f.Until.IsZero() { - q = q.Where("logged_at < ?", f.Until) - } - return q + return q.Where(cond, args...) } // normalizeNodeAccessLogHosts 对 hosts 归一化:trim + lowercase + 去重去空。 diff --git a/internal/repository/logstore/postgres_store_test.go b/internal/repository/logstore/postgres_store_test.go index 33ec805d..da5b8c32 100644 --- a/internal/repository/logstore/postgres_store_test.go +++ b/internal/repository/logstore/postgres_store_test.go @@ -803,6 +803,76 @@ func TestGormWAFAndIPSummaries(t *testing.T) { } } +// TestGormIPSummariesRegionWithinFilterWindow 验证 IPSummaries 的 region 取过滤窗口内该 IP +// 最近一条(对齐 CH argMax(region, logged_at)):窗口外更新的 region 记录不参与, +// 旧实现(子查询不带窗口条件)会错误返回窗口外那条。 +func TestGormIPSummariesRegionWithinFilterWindow(t *testing.T) { + ResetForTest() + SetConfigReader(func(_ context.Context, _ string) (string, error) { return "", nil }) + s := newTestGormStore(t) + ctx := context.Background() + base := time.Now().UTC().Truncate(time.Minute) + rows := []analyticsmodel.NodeAccessLog{ + // 窗口外(晚于 Until):同一 IP 的更新 region,旧实现会误取。 + {ID: 1, NodeID: "n1", LoggedAt: base.Add(3 * time.Hour), RemoteAddr: "1.1.1.1", Region: "outside-new", StatusCode: 200}, + // 窗口内 [base, base+2h) 最新一条:region 应为 inside-old。 + {ID: 2, NodeID: "n1", LoggedAt: base.Add(time.Hour), RemoteAddr: "1.1.1.1", Region: "inside-old", StatusCode: 200}, + // 窗口内更早的一条,不应覆盖窗口内最新 region。 + {ID: 3, NodeID: "n1", LoggedAt: base, RemoteAddr: "1.1.1.1", Region: "inside-earlier", StatusCode: 200}, + } + if err := s.BatchInsertNodeAccessLogs(ctx, rows); err != nil { + t.Fatalf("insert: %v", err) + } + sums, err := s.IPSummaries(ctx, model.OpenFlareAccessLogQuery{ + NodeID: "n1", + Since: base, + Until: base.Add(2 * time.Hour), + }, time.Time{}) + if err != nil { + t.Fatalf("ip summaries: %v", err) + } + if len(sums) != 1 || sums[0].RemoteAddr != "1.1.1.1" || sums[0].Region != "inside-old" { + t.Fatalf("ip summaries = %+v, want 1.1.1.1 region inside-old (window-external outside-new excluded)", sums) + } +} + +// TestGormIPSummariesEmptyFilter 覆盖 IPSummaries 空过滤分支(cond=="" 时 region 子查询 +// 不带参数、Select 走无 args 路径):OpenFlareAccessLogQuery{} 不报错、返回全部 IP 分组, +// region 为各 IP 全部行中最新一条(无窗口即全部行)。 +func TestGormIPSummariesEmptyFilter(t *testing.T) { + ResetForTest() + SetConfigReader(func(_ context.Context, _ string) (string, error) { return "", nil }) + s := newTestGormStore(t) + ctx := context.Background() + base := time.Now().UTC().Truncate(time.Minute) + rows := []analyticsmodel.NodeAccessLog{ + {ID: 1, NodeID: "n1", LoggedAt: base, RemoteAddr: "1.1.1.1", Region: "cn", StatusCode: 200}, + {ID: 2, NodeID: "n1", LoggedAt: base.Add(time.Minute), RemoteAddr: "1.1.1.1", Region: "us", StatusCode: 200}, + {ID: 3, NodeID: "n2", LoggedAt: base, RemoteAddr: "2.2.2.2", Region: "jp", StatusCode: 404}, + } + if err := s.BatchInsertNodeAccessLogs(ctx, rows); err != nil { + t.Fatalf("insert: %v", err) + } + sums, err := s.IPSummaries(ctx, model.OpenFlareAccessLogQuery{}, time.Time{}) + if err != nil { + t.Fatalf("ip summaries empty filter: %v", err) + } + if len(sums) != 2 { + t.Fatalf("ip summaries = %d groups, want 2", len(sums)) + } + byIP := map[string]analyticsmodel.NodeAccessLogIPSummary{} + for _, x := range sums { + byIP[x.RemoteAddr] = x + } + // 无窗口即全部行:1.1.1.1 最新一条 region 为 us(region 子查询不带参数路径)。 + if byIP["1.1.1.1"].Region != "us" || byIP["1.1.1.1"].TotalRequests != 2 { + t.Fatalf("ip summaries 1.1.1.1 = %+v, want region us requests 2", byIP["1.1.1.1"]) + } + if byIP["2.2.2.2"].Region != "jp" || byIP["2.2.2.2"].TotalRequests != 1 { + t.Fatalf("ip summaries 2.2.2.2 = %+v, want region jp requests 1", byIP["2.2.2.2"]) + } +} + // TestGormListRejectsUnsupportedSortBy 验证 List 对不支持的 SortBy 直接报错, // 默认 logged_at 路径与 CH 支持的 status_code/remote_addr 正常可用。 func TestGormListRejectsUnsupportedSortBy(t *testing.T) { diff --git a/internal/repository/logstore/provider.go b/internal/repository/logstore/provider.go index d21f7fb4..551a47a0 100644 --- a/internal/repository/logstore/provider.go +++ b/internal/repository/logstore/provider.go @@ -142,9 +142,18 @@ func Migrating(ctx context.Context) bool { return v == "migrating" } -// Init 在 bootstrap 阶段预热一次激活 store(幂等,失败不致命——首次使用时再解析)。 +// Init 在 bootstrap 阶段预热一次激活 store(幂等,失败不致命——首次使用时再解析), +// 并兜底预建「当前月 + 未来 2 个月」分区:进程停机跨月边界、重启后每日 cleanup 之前 +// 首次写入不会报 "no partition of relation found"(CH/SQLite 分支 EnsurePartitions 为 no-op)。 func Init(ctx context.Context) { - _, _ = Active(ctx) + s, err := Active(ctx) + if err != nil { + return + } + now := time.Now().UTC() + if err := s.AccessLogs.EnsurePartitions(ctx, now, now.AddDate(0, partitionLeadMonths, 0)); err != nil { + logger.WarnF(ctx, "logstore: ensure startup partitions failed: %v", err) + } } // InvalidateCache 清空日志库解析缓存(在修改 log_database 配置后显式调用)。