diff --git a/docs/changelog/index.md b/docs/changelog/index.md index b1f8b139..aeff802a 100644 --- a/docs/changelog/index.md +++ b/docs/changelog/index.md @@ -34,6 +34,7 @@ sidebar: false - 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」。 +- 用户访问日志(`w_user_access_logs`)记录禁用:不再采集与写入新的用户访问日志,已存历史数据与管理端访问日志统计页面保留;日志库迁移任务不再排空用户访问日志队列,状态接口不再展示其缓冲队列统计。 ### 修复 diff --git a/internal/apps/admin/status/clickhouse.go b/internal/apps/admin/status/clickhouse.go index 8c063d10..3785466c 100644 --- a/internal/apps/admin/status/clickhouse.go +++ b/internal/apps/admin/status/clickhouse.go @@ -9,7 +9,6 @@ import ( "net/http" "github.com/Rain-kl/Wavelet/internal/apps/openflare/chwriter" - "github.com/Rain-kl/Wavelet/internal/apps/risk_control" "github.com/Rain-kl/Wavelet/internal/infra/config" "github.com/Rain-kl/Wavelet/internal/infra/persistence/batchwriter" "github.com/Rain-kl/Wavelet/internal/model" @@ -132,6 +131,5 @@ func collectBatchWriterStats() []batchwriter.Stats { if out == nil { out = make([]batchwriter.Stats, 0, 1) } - out = append(out, risk_control.LogWriterStats()) return out } diff --git a/internal/apps/openflare/tasks/log_db_switch.go b/internal/apps/openflare/tasks/log_db_switch.go index 12f2053b..78d193d0 100644 --- a/internal/apps/openflare/tasks/log_db_switch.go +++ b/internal/apps/openflare/tasks/log_db_switch.go @@ -11,7 +11,6 @@ import ( "time" "github.com/Rain-kl/Wavelet/internal/apps/openflare/chwriter" - "github.com/Rain-kl/Wavelet/internal/apps/risk_control" "github.com/Rain-kl/Wavelet/internal/infra/config" "github.com/Rain-kl/Wavelet/internal/infra/task" "github.com/Rain-kl/Wavelet/internal/model" @@ -179,13 +178,10 @@ func currentLogDatabase(ctx context.Context) (string, error) { return cfg.Value, nil } -// drainLogWriters 等待 chwriter(节点访问日志 + 可观测 4 表)与 risk_control -// (用户访问日志)的在途批次全部落库。见设计 §7.2:先排空再冻结。 +// drainLogWriters 等待 chwriter(节点访问日志 + 可观测 4 表)的在途批次全部落库。 +// 见设计 §7.2:先排空再冻结。用户访问日志(w_user_access_logs)记录已禁用,无在途批次。 func drainLogWriters(ctx context.Context) error { - if err := chwriter.Drain(ctx); err != nil { - return err - } - return risk_control.DrainLogWriter(ctx) + return chwriter.Drain(ctx) } // setMigrationFlag 写入迁移冻结标记。用 SaveOrUpdateSystemConfig:行缺失时 upsert, diff --git a/internal/apps/risk_control/logics.go b/internal/apps/risk_control/logics.go deleted file mode 100644 index 5ebd33b3..00000000 --- a/internal/apps/risk_control/logics.go +++ /dev/null @@ -1,162 +0,0 @@ -// Copyright 2026 Arctel.net -// SPDX-License-Identifier: Apache-2.0 - -package risk_control - -import ( - "context" - "sync" - "time" - - "github.com/Rain-kl/Wavelet/internal/infra/persistence/batchwriter" - "github.com/Rain-kl/Wavelet/internal/model/analytics" - "github.com/Rain-kl/Wavelet/internal/platform/lifecycle" - "github.com/Rain-kl/Wavelet/internal/repository/logstore" - "github.com/Rain-kl/Wavelet/pkg/logger" -) - -const ( - // Bound visibility lag for sparse access-log traffic when MinBatchSize is not met. - accessLogMaxFlushWait = 3 * time.Second -) - -var ( - logWriterMu sync.RWMutex - logWriter *batchwriter.Writer[*analytics.UserAccessLog] -) - -// InitLogWriter initializes the user access-log batch writer. -// The active log store is resolved via logstore at flush time, so the writer -// runs for PG/SQLite as well as ClickHouse. -func InitLogWriter(ctx context.Context) { - logWriterMu.Lock() - defer logWriterMu.Unlock() - if logWriter != nil { - return - } - - cfg := batchwriter.DefaultConfig() - cfg.Name = "user_access_logs" - cfg.MaxFlushWait = accessLogMaxFlushWait - writer, err := batchwriter.New[*analytics.UserAccessLog](cfg, func(ctx context.Context, items []*analytics.UserAccessLog) error { - rows := make([]analytics.UserAccessLog, 0, len(items)) - for _, item := range items { - if item == nil { - continue - } - rows = append(rows, *item) - } - s, err := logstore.Active(ctx) - if err != nil { - return err - } - return s.UserAccessLogs.BatchInsert(ctx, rows) - }, - batchwriter.WithDropHandler[*analytics.UserAccessLog](func(item *analytics.UserAccessLog) { - path := "" - if item != nil { - path = item.Path - } - logger.WarnF(context.Background(), "[RiskControl] Log queue full, dropping log item for path: %s", path) - }), - batchwriter.WithFlushErrorHandler[*analytics.UserAccessLog](func(ctx context.Context, items []*analytics.UserAccessLog, err error) { - logger.ErrorF(ctx, "[RiskControl] Send log batch failed (batch=%d): %v", len(items), err) - }), - ) - if err != nil { - logger.ErrorF(ctx, "[RiskControl] init log writer failed: %v", err) - return - } - - writer.Start(ctx) - logWriter = writer - lifecycle.OnShutdown("risk_control_log_writer", StopLogWriter) -} - -// StopLogWriter stops the user access-log batch writer and drains pending logs. -func StopLogWriter(ctx context.Context) error { - writer := currentLogWriter() - if writer == nil { - return nil - } - return writer.Stop(ctx) -} - -// DrainLogWriter 等待用户访问日志 writer 的在途批次落库:队列 Depth 归零后 -// 再保持一个 flush 周期(1s)持续为空才返回;不停止 writer(迁移冻结后由 -// ensureWritable 拒绝新写入)。writer 未初始化时直接返回 nil。 -func DrainLogWriter(ctx context.Context) error { - writer := currentLogWriter() - if writer == nil { - return nil - } - ticker := time.NewTicker(drainPollInterval) - defer ticker.Stop() - var quietSince time.Time - for { - if writer.Stats().Depth == 0 { - if quietSince.IsZero() { - quietSince = time.Now() - } else if time.Since(quietSince) >= batchwriter.DefaultConfig().FlushInterval { - return nil - } - } else { - quietSince = time.Time{} - } - select { - case <-ctx.Done(): - return ctx.Err() - case <-ticker.C: - } - } -} - -// drainPollInterval 用户访问日志队列轮询间隔。 -const drainPollInterval = 50 * time.Millisecond - -// IsBufferFull reports whether the access-log queue has no remaining capacity. -func IsBufferFull() bool { - writer := currentLogWriter() - if writer == nil { - return false - } - return writer.IsFull() -} - -// LogWriterStats returns queue depth and failure counters for the access-log writer. -// When the writer is not initialized, it returns a zero-value Stats with the expected name. -func LogWriterStats() batchwriter.Stats { - writer := currentLogWriter() - if writer == nil { - return batchwriter.Stats{Name: "user_access_logs"} - } - return writer.Stats() -} - -// QueueAccessLog enqueues an access log without blocking. -func QueueAccessLog(logItem *analytics.UserAccessLog) { - writer := currentLogWriter() - if writer == nil || logItem == nil { - return - } - writer.TryEnqueue(logItem) -} - -// SetLogWriterForTest swaps the access-log writer for unit tests. -func SetLogWriterForTest(writer *batchwriter.Writer[*analytics.UserAccessLog]) func() { - logWriterMu.Lock() - previous := logWriter - logWriter = writer - logWriterMu.Unlock() - return func() { - logWriterMu.Lock() - logWriter = previous - logWriterMu.Unlock() - } -} - -func currentLogWriter() *batchwriter.Writer[*analytics.UserAccessLog] { - logWriterMu.RLock() - defer logWriterMu.RUnlock() - return logWriter -} diff --git a/internal/apps/risk_control/logics_test.go b/internal/apps/risk_control/logics_test.go deleted file mode 100644 index 69fe2490..00000000 --- a/internal/apps/risk_control/logics_test.go +++ /dev/null @@ -1,30 +0,0 @@ -// Copyright 2026 Arctel.net -// SPDX-License-Identifier: Apache-2.0 - -package risk_control - -import ( - "testing" - "time" -) - -func TestAccessLogMaxFlushWaitInRange(t *testing.T) { - t.Parallel() - if accessLogMaxFlushWait < 2*time.Second || accessLogMaxFlushWait > 5*time.Second { - t.Fatalf("accessLogMaxFlushWait = %v, want in [2s, 5s]", accessLogMaxFlushWait) - } -} - -func TestLogWriterStatsWhenNil(t *testing.T) { - t.Parallel() - reset := SetLogWriterForTest(nil) - t.Cleanup(reset) - - stats := LogWriterStats() - if stats.Name != "user_access_logs" { - t.Fatalf("LogWriterStats().Name = %q, want user_access_logs", stats.Name) - } - if stats.Running { - t.Fatal("LogWriterStats().Running = true for nil writer, want false") - } -} diff --git a/internal/apps/risk_control/middleware.go b/internal/apps/risk_control/middleware.go deleted file mode 100644 index 43b7f2fc..00000000 --- a/internal/apps/risk_control/middleware.go +++ /dev/null @@ -1,136 +0,0 @@ -// Copyright 2026 Arctel.net -// SPDX-License-Identifier: Apache-2.0 - -// Package risk_control 提供风险控制中间件 -package risk_control - -import ( - "crypto/sha256" - "encoding/hex" - "encoding/json" - "net/http" - "strings" - "time" - - "github.com/Rain-kl/Wavelet/internal/apps/oauth" - "github.com/Rain-kl/Wavelet/internal/infra/persistence/idgen" - "github.com/Rain-kl/Wavelet/internal/model" - "github.com/Rain-kl/Wavelet/internal/model/analytics" - "github.com/Rain-kl/Wavelet/internal/repository/logstore" - "github.com/Rain-kl/Wavelet/internal/shared/response" - "github.com/Rain-kl/Wavelet/pkg/logger" - "github.com/gin-gonic/gin" -) - -const maxAuditLogHeadersBytes = 2 * 1024 - -var auditLogHeaderAllowlist = map[string]struct{}{ - "Authorization": {}, - "Cookie": {}, - "X-Forwarded-For": {}, - "X-Real-Ip": {}, - "User-Agent": {}, - "Content-Type": {}, -} - -func marshalAuditLogHeaders(headers http.Header) string { - if headers == nil { - return "" - } - - filtered := make(http.Header) - for key, values := range headers { - if _, ok := auditLogHeaderAllowlist[key]; !ok { - continue - } - filtered[key] = redactAuditLogHeaderValues(key, values) - } - - headersBytes, err := json.Marshal(filtered) - if err != nil { - return "" - } - if len(headersBytes) <= maxAuditLogHeadersBytes { - return string(headersBytes) - } - return string(headersBytes[:maxAuditLogHeadersBytes]) -} - -func redactAuditLogHeaderValues(key string, values []string) []string { - switch key { - case "Authorization", "Cookie": - redacted := make([]string, len(values)) - for i, value := range values { - redacted[i] = hashAuditLogSensitiveValue(value) - } - return redacted - default: - return values - } -} - -func hashAuditLogSensitiveValue(value string) string { - if strings.TrimSpace(value) == "" { - return "" - } - sum := sha256.Sum256([]byte(value)) - return "sha256:" + hex.EncodeToString(sum[:8]) -} - -// RiskControlMiddleware 全局日志采集中间件 -func RiskControlMiddleware() gin.HandlerFunc { - return func(c *gin.Context) { - // 日志库迁移冻结期跳过采集,不阻断业务请求。 - if logstore.Migrating(c.Request.Context()) { - logger.WarnF(c.Request.Context(), "[RiskControl] log DB migrating, skip audit log") - c.Next() - return - } - - // 1. 限流背压检测(检测本地缓冲队列是否已满) - if IsBufferFull() { - response.AbortTooManyRequests(c, "系统繁忙,请稍后再试") - return - } - - start := time.Now() - - // 2. 执行后续请求(穿过业务处理和认证中间件) - c.Next() - - // 3. 后置身份检查:仅记录通过认证的请求 - userObj, exists := oauth.GetFromContext[*model.User](c, oauth.UserObjKey) - if !exists || userObj == nil { - return - } - - // 4. 计算耗时并异步推送到缓冲队列 - latency := time.Since(start).Milliseconds() - - headersStr := marshalAuditLogHeaders(c.Request.Header) - - const maxHTTPStatus = 999 - status := c.Writer.Status() - if status < 0 { - status = 0 - } else if status > maxHTTPStatus { - status = maxHTTPStatus - } - - logItem := &analytics.UserAccessLog{ - ID: idgen.NextUint64ID(), - UserID: userObj.ID, // 直接从 Context 获取已登录用户ID,避免数据库查询 - Path: c.Request.URL.Path, - Method: c.Request.Method, - IP: c.ClientIP(), - UserAgent: c.Request.UserAgent(), - Headers: headersStr, - Status: int32(status), - Latency: latency, - CreatedAt: time.Now(), - } - - // 非阻塞地推入缓存队列 - QueueAccessLog(logItem) - } -} diff --git a/internal/apps/risk_control/middleware_test.go b/internal/apps/risk_control/middleware_test.go deleted file mode 100644 index 15641078..00000000 --- a/internal/apps/risk_control/middleware_test.go +++ /dev/null @@ -1,178 +0,0 @@ -// Copyright 2026 Arctel.net -// SPDX-License-Identifier: Apache-2.0 - -package risk_control - -import ( - "context" - "encoding/json" - "net/http" - "net/http/httptest" - "testing" - "time" - - "github.com/Rain-kl/Wavelet/internal/apps/oauth" - "github.com/Rain-kl/Wavelet/internal/infra/config" - "github.com/Rain-kl/Wavelet/internal/infra/persistence/batchwriter" - "github.com/Rain-kl/Wavelet/internal/model" - "github.com/Rain-kl/Wavelet/internal/model/analytics" - "github.com/Rain-kl/Wavelet/internal/testhelper" - "github.com/gin-gonic/gin" - "github.com/stretchr/testify/assert" -) - -const testLogQueueSize = 10_000 - -func newTestLogWriter(t *testing.T, queueSize int) (*batchwriter.Writer[*analytics.UserAccessLog], chan *analytics.UserAccessLog) { - t.Helper() - - received := make(chan *analytics.UserAccessLog, queueSize) - cfg := batchwriter.Config{ - QueueSize: queueSize, - MaxBatchSize: 1, - FlushInterval: 5 * time.Millisecond, - } - writer, err := batchwriter.New[*analytics.UserAccessLog](cfg, func(_ context.Context, items []*analytics.UserAccessLog) error { - for _, item := range items { - if item != nil { - received <- item - } - } - return nil - }) - assert.NoError(t, err) - writer.Start(context.Background()) - t.Cleanup(func() { - stopCtx, cancel := context.WithTimeout(context.Background(), time.Second) - defer cancel() - _ = writer.Stop(stopCtx) - }) - return writer, received -} - -func TestRiskControlMiddleware(t *testing.T) { - gin.SetMode(gin.TestMode) - - t.Run("ClickHouse disabled", func(t *testing.T) { - config.Config.ClickHouse.Enabled = false - defer func() { config.Config.ClickHouse.Enabled = false }() - - r := testhelper.NewTestGinEngine(RiskControlMiddleware()) - r.GET("/test", func(c *gin.Context) { - c.String(http.StatusOK, "ok") - }) - - w := httptest.NewRecorder() - req, _ := http.NewRequest(http.MethodGet, "/test", nil) - r.ServeHTTP(w, req) - - assert.Equal(t, http.StatusOK, w.Code) - assert.Equal(t, "ok", w.Body.String()) - }) - - t.Run("ClickHouse enabled - Normal Authenticated Request", func(t *testing.T) { - config.Config.ClickHouse.Enabled = true - writer, received := newTestLogWriter(t, testLogQueueSize) - resetWriter := SetLogWriterForTest(writer) - defer func() { - config.Config.ClickHouse.Enabled = false - resetWriter() - }() - - r := gin.New() - r.Use(func(c *gin.Context) { - user := &model.User{ID: 12345} - oauth.SetToContext(c, oauth.UserObjKey, user) - c.Next() - }) - r.Use(RiskControlMiddleware()) - r.GET("/test", func(c *gin.Context) { - c.String(http.StatusOK, "ok") - }) - - w := httptest.NewRecorder() - req, _ := http.NewRequest(http.MethodGet, "/test", nil) - req.Header.Set("X-Test-Header", "hello") - req.Header.Set("Cookie", "session_id=abcdef123456") - req.Header.Set("Authorization", "Bearer secret-token") - req.Header.Set("Content-Type", "application/json") - r.ServeHTTP(w, req) - - assert.Equal(t, http.StatusOK, w.Code) - assert.Equal(t, "ok", w.Body.String()) - - select { - case logItem := <-received: - assert.Equal(t, uint64(12345), logItem.UserID) - assert.Equal(t, "/test", logItem.Path) - assert.Equal(t, http.MethodGet, logItem.Method) - assert.Equal(t, int32(http.StatusOK), logItem.Status) - assert.NotEmpty(t, logItem.Headers) - assert.NotContains(t, logItem.Headers, "X-Test-Header") - assert.Contains(t, logItem.Headers, "Content-Type") - assert.Contains(t, logItem.Headers, "sha256:") - assert.NotContains(t, logItem.Headers, "secret-token") - assert.NotContains(t, logItem.Headers, "session_id=abcdef123456") - case <-time.After(200 * time.Millisecond): - t.Fatal("expected flushed log item, but got none") - } - }) - - t.Run("ClickHouse enabled - Unauthenticated Request", func(t *testing.T) { - config.Config.ClickHouse.Enabled = true - writer, received := newTestLogWriter(t, testLogQueueSize) - resetWriter := SetLogWriterForTest(writer) - defer func() { - config.Config.ClickHouse.Enabled = false - resetWriter() - }() - - r := testhelper.NewTestGinEngine(RiskControlMiddleware()) - r.GET("/test", func(c *gin.Context) { - c.String(http.StatusOK, "ok") - }) - - w := httptest.NewRecorder() - req, _ := http.NewRequest(http.MethodGet, "/test", nil) - r.ServeHTTP(w, req) - - assert.Equal(t, http.StatusOK, w.Code) - assert.Equal(t, "ok", w.Body.String()) - - select { - case <-received: - t.Fatal("expected no log item for unauthenticated request") - case <-time.After(50 * time.Millisecond): - } - }) - - t.Run("ClickHouse enabled - Buffer Full Rate Limiting", func(t *testing.T) { - config.Config.ClickHouse.Enabled = true - writer, _ := newTestLogWriter(t, 2) - resetWriter := SetLogWriterForTest(writer) - defer func() { - config.Config.ClickHouse.Enabled = false - resetWriter() - }() - - for range 2 { - assert.True(t, writer.TryEnqueue(&analytics.UserAccessLog{})) - } - - r := testhelper.NewTestGinEngine(RiskControlMiddleware()) - r.GET("/test", func(c *gin.Context) { - c.String(http.StatusOK, "ok") - }) - - w := httptest.NewRecorder() - req, _ := http.NewRequest(http.MethodGet, "/test", nil) - r.ServeHTTP(w, req) - - assert.Equal(t, http.StatusTooManyRequests, w.Code) - - var resp map[string]interface{} - err := json.Unmarshal(w.Body.Bytes(), &resp) - assert.NoError(t, err) - assert.Contains(t, resp["error_msg"], "系统繁忙") - }) -} diff --git a/internal/infra/persistence/migrator/migrator_test.go b/internal/infra/persistence/migrator/migrator_test.go index 33d3e867..cd55aded 100644 --- a/internal/infra/persistence/migrator/migrator_test.go +++ b/internal/infra/persistence/migrator/migrator_test.go @@ -19,10 +19,11 @@ import ( "gorm.io/gorm" ) -// expectedMigratedSystemConfigCount 包含初始 32 项系统配置、202606220004 -// 从 of_options 迁移过来的 48 项业务配置、Pages 的 2 项业务配置、 -// OpenResty 默认限流的 3 项业务配置,以及单 IP 请求频率限制 1 项业务配置。 -const expectedMigratedSystemConfigCount = 86 +// expectedMigratedSystemConfigCount 为全新库执行全部迁移后 w_system_configs 的行数 +// (初始系统配置 + 各期配置迁移/新增 seed:of_options 迁移、文件白名单、磁盘缓存、 +// 登录会话 TTL、升级源、存储、FRPS Web UI、Pages、OpenResty 限流、单 IP 限频、 +// 错误页、SW 离线、日志保留期、指标保留期等);新增配置 seed 迁移时需同步更新本常量。 +const expectedMigratedSystemConfigCount = 95 func TestMigrateInitializesSQLiteDatabase(t *testing.T) { sqliteDB, err := gorm.Open(sqlite.Open(":memory:"), &gorm.Config{ diff --git a/internal/platform/bootstrap/bootstrap.go b/internal/platform/bootstrap/bootstrap.go index 2fb1206d..259f5b75 100644 --- a/internal/platform/bootstrap/bootstrap.go +++ b/internal/platform/bootstrap/bootstrap.go @@ -16,7 +16,6 @@ import ( "github.com/Rain-kl/Wavelet/internal/apps/admin/push/custom_events" "github.com/Rain-kl/Wavelet/internal/apps/openflare/chwriter" ofgeoip "github.com/Rain-kl/Wavelet/internal/apps/openflare/geoip" - "github.com/Rain-kl/Wavelet/internal/apps/risk_control" "github.com/Rain-kl/Wavelet/internal/infra/config" taskhandlers "github.com/Rain-kl/Wavelet/internal/infra/task/handlers" "github.com/Rain-kl/Wavelet/internal/model" @@ -168,7 +167,6 @@ func Init(ctx context.Context, opts Options) { logger.ErrorF(ctx, "[Bootstrap] sync push events failed: %v", err) } if opts.API { - risk_control.InitLogWriter(ctx) chwriter.Init(ctx) } }) diff --git a/internal/router/router.go b/internal/router/router.go index d429aff3..ef821e72 100644 --- a/internal/router/router.go +++ b/internal/router/router.go @@ -16,7 +16,6 @@ import ( "syscall" "time" - "github.com/Rain-kl/Wavelet/internal/apps/risk_control" "github.com/Rain-kl/Wavelet/internal/platform/bootstrap" router_root "github.com/Rain-kl/Wavelet/internal/router/root" v1 "github.com/Rain-kl/Wavelet/internal/router/v1" @@ -76,7 +75,7 @@ func Serve(onStarted func()) { r.Use(sessions.Sessions(config.Config.App.SessionCookieName, sessionStore)) // 补充中间件 - r.Use(otelgin.Middleware(config.Config.App.AppName), errorHandlerMiddleware(), loggerMiddleware(), risk_control.RiskControlMiddleware()) + r.Use(otelgin.Middleware(config.Config.App.AppName), errorHandlerMiddleware(), loggerMiddleware()) registerRoutes(r)