mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-09-28 05:46:36 +08:00
feat(log): disable user access log recording
- 移除全局用户访问日志采集中间件与批写入 writer(risk_control 包整包删除), 不再写入 w_user_access_logs;存量数据与管理端访问日志统计页面保留 - 日志库迁移任务不再排空用户访问日志队列,状态接口不再展示其缓冲队列统计 - 迁移测试的系统配置 seed 计数断言更新为当前实际值(86 → 95), 注释改为提示新增配置 seed 时同步更新
This commit is contained in:
@@ -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`)记录禁用:不再采集与写入新的用户访问日志,已存历史数据与管理端访问日志统计页面保留;日志库迁移任务不再排空用户访问日志队列,状态接口不再展示其缓冲队列统计。
|
||||
|
||||
### 修复
|
||||
|
||||
|
||||
@@ -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
|
||||
}
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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
|
||||
}
|
||||
@@ -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")
|
||||
}
|
||||
}
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
@@ -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"], "系统繁忙")
|
||||
})
|
||||
}
|
||||
@@ -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{
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
})
|
||||
|
||||
@@ -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)
|
||||
|
||||
|
||||
Reference in New Issue
Block a user