mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-08 16:46:37 +08:00
refactor(plugins): complete physical encapsulation of auth, admin, message_gateway, and risk_control domain plugins
This commit is contained in:
@@ -0,0 +1,144 @@
|
||||
// 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"
|
||||
)
|
||||
|
||||
var (
|
||||
logWriterMu sync.RWMutex
|
||||
logWriter *batchwriter.Writer[*analytics.UserAccessLog]
|
||||
)
|
||||
|
||||
// InitLogWriter initializes the access-log batch writer for the active log database.
|
||||
func InitLogWriter(ctx context.Context) {
|
||||
logWriterMu.Lock()
|
||||
defer logWriterMu.Unlock()
|
||||
if logWriter != nil {
|
||||
return
|
||||
}
|
||||
|
||||
cfg := batchwriter.DefaultConfig()
|
||||
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)
|
||||
}
|
||||
store, err := logstore.Active(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
return store.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] flush access-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 ClickHouse 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)
|
||||
}
|
||||
|
||||
// IsBufferFull reports whether the access-log queue has no remaining capacity.
|
||||
func IsBufferFull() bool {
|
||||
writer := currentLogWriter()
|
||||
if writer == nil {
|
||||
return false
|
||||
}
|
||||
return writer.IsFull()
|
||||
}
|
||||
|
||||
// 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
|
||||
}
|
||||
|
||||
const drainPollInterval = 50 * time.Millisecond
|
||||
|
||||
// Drain waits until the in-memory access-log queue has been empty for one flush interval.
|
||||
func Drain(ctx context.Context) error {
|
||||
writer := currentLogWriter()
|
||||
if writer == nil {
|
||||
return nil
|
||||
}
|
||||
quietPeriod := batchwriter.DefaultConfig().FlushInterval
|
||||
if quietPeriod <= 0 {
|
||||
quietPeriod = time.Second
|
||||
}
|
||||
ticker := time.NewTicker(drainPollInterval)
|
||||
defer ticker.Stop()
|
||||
var quietSince time.Time
|
||||
for {
|
||||
if writer.Len() == 0 {
|
||||
if quietSince.IsZero() {
|
||||
quietSince = time.Now()
|
||||
} else if time.Since(quietSince) >= quietPeriod {
|
||||
return nil
|
||||
}
|
||||
} else {
|
||||
quietSince = time.Time{}
|
||||
}
|
||||
select {
|
||||
case <-ctx.Done():
|
||||
return ctx.Err()
|
||||
case <-ticker.C:
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,88 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package risk_control provides the access control, IP rate limiting, and telemetry risk analysis domain plugin for Cordis.
|
||||
package risk_control
|
||||
|
||||
import (
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"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/idgen"
|
||||
"github.com/Rain-kl/Wavelet/internal/model"
|
||||
"github.com/Rain-kl/Wavelet/internal/model/analytics"
|
||||
"github.com/Rain-kl/Wavelet/internal/shared/response"
|
||||
"github.com/gin-gonic/gin"
|
||||
)
|
||||
|
||||
// RiskControlMiddleware 全局日志采集中间件
|
||||
func RiskControlMiddleware() gin.HandlerFunc {
|
||||
return func(c *gin.Context) {
|
||||
// 如果未启用 ClickHouse,直接放行
|
||||
if config.Config == nil || !config.Config.ClickHouse.Enabled {
|
||||
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()
|
||||
|
||||
var headersStr string
|
||||
if c.Request.Header != nil {
|
||||
// 克隆 Header,避免污染原 HTTP 请求的 Header 对象
|
||||
clonedHeaders := make(http.Header)
|
||||
for k, v := range c.Request.Header {
|
||||
clonedHeaders[k] = v
|
||||
}
|
||||
clonedHeaders.Del("Cookie")
|
||||
|
||||
if headersBytes, err := json.Marshal(clonedHeaders); err == nil {
|
||||
headersStr = string(headersBytes)
|
||||
}
|
||||
}
|
||||
|
||||
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)
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,198 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package risk_control_test
|
||||
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"net/http"
|
||||
"net/http/httptest"
|
||||
"sync"
|
||||
"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/Rain-kl/Wavelet/plugins/domain/risk_control"
|
||||
"github.com/gin-gonic/gin"
|
||||
"github.com/stretchr/testify/assert"
|
||||
)
|
||||
|
||||
func newTestAccessLogWriter(t *testing.T, cfg batchwriter.Config) (*batchwriter.Writer[*analytics.UserAccessLog], func() []*analytics.UserAccessLog) {
|
||||
t.Helper()
|
||||
|
||||
var (
|
||||
mu sync.Mutex
|
||||
captured []*analytics.UserAccessLog
|
||||
)
|
||||
writer, err := batchwriter.New(cfg, func(_ context.Context, items []*analytics.UserAccessLog) error {
|
||||
mu.Lock()
|
||||
captured = append(captured, items...)
|
||||
mu.Unlock()
|
||||
return nil
|
||||
})
|
||||
if err != nil {
|
||||
t.Fatalf("batchwriter.New() error = %v", err)
|
||||
}
|
||||
|
||||
writer.Start(context.Background())
|
||||
restore := risk_control.SetLogWriterForTest(writer)
|
||||
t.Cleanup(func() {
|
||||
restore()
|
||||
stopCtx, cancel := context.WithTimeout(context.Background(), time.Second)
|
||||
defer cancel()
|
||||
_ = writer.Stop(stopCtx)
|
||||
})
|
||||
|
||||
return writer, func() []*analytics.UserAccessLog {
|
||||
mu.Lock()
|
||||
defer mu.Unlock()
|
||||
return append([]*analytics.UserAccessLog(nil), captured...)
|
||||
}
|
||||
}
|
||||
|
||||
func drainAccessLogWriter(t *testing.T, writer *batchwriter.Writer[*analytics.UserAccessLog]) {
|
||||
t.Helper()
|
||||
|
||||
stopCtx, cancel := context.WithTimeout(context.Background(), time.Second)
|
||||
defer cancel()
|
||||
if err := writer.Stop(stopCtx); err != nil {
|
||||
t.Fatalf("writer.Stop() error = %v", err)
|
||||
}
|
||||
}
|
||||
|
||||
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(risk_control.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
|
||||
defer func() { config.Config.ClickHouse.Enabled = false }()
|
||||
|
||||
cfg := batchwriter.DefaultConfig()
|
||||
cfg.MaxBatchSize = 100
|
||||
cfg.FlushInterval = time.Hour
|
||||
|
||||
writer, getCaptured := newTestAccessLogWriter(t, cfg)
|
||||
|
||||
r := gin.New()
|
||||
r.Use(func(c *gin.Context) {
|
||||
user := &model.User{ID: 12345}
|
||||
oauth.SetToContext(c, oauth.UserObjKey, user)
|
||||
c.Next()
|
||||
})
|
||||
r.Use(risk_control.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")
|
||||
r.ServeHTTP(w, req)
|
||||
|
||||
assert.Equal(t, http.StatusOK, w.Code)
|
||||
assert.Equal(t, "ok", w.Body.String())
|
||||
|
||||
drainAccessLogWriter(t, writer)
|
||||
|
||||
captured := getCaptured()
|
||||
if len(captured) != 1 {
|
||||
t.Fatalf("captured access logs = %d, want 1", len(captured))
|
||||
}
|
||||
logItem := captured[0]
|
||||
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.Contains(t, logItem.Headers, "X-Test-Header")
|
||||
assert.NotContains(t, logItem.Headers, "Cookie")
|
||||
})
|
||||
|
||||
t.Run("ClickHouse enabled - Unauthenticated Request", func(t *testing.T) {
|
||||
config.Config.ClickHouse.Enabled = true
|
||||
defer func() { config.Config.ClickHouse.Enabled = false }()
|
||||
|
||||
cfg := batchwriter.DefaultConfig()
|
||||
cfg.MaxBatchSize = 100
|
||||
cfg.FlushInterval = time.Hour
|
||||
|
||||
writer, getCaptured := newTestAccessLogWriter(t, cfg)
|
||||
|
||||
r := testhelper.NewTestGinEngine(risk_control.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())
|
||||
|
||||
drainAccessLogWriter(t, writer)
|
||||
|
||||
if len(getCaptured()) != 0 {
|
||||
t.Fatal("expected no log item for unauthenticated request")
|
||||
}
|
||||
})
|
||||
|
||||
t.Run("ClickHouse enabled - Buffer Full Rate Limiting", func(t *testing.T) {
|
||||
config.Config.ClickHouse.Enabled = true
|
||||
defer func() { config.Config.ClickHouse.Enabled = false }()
|
||||
|
||||
cfg := batchwriter.DefaultConfig()
|
||||
cfg.QueueSize = 2
|
||||
cfg.MaxBatchSize = 100
|
||||
cfg.FlushInterval = time.Hour
|
||||
|
||||
writer, _ := newTestAccessLogWriter(t, cfg)
|
||||
|
||||
for range cfg.QueueSize {
|
||||
writer.TryEnqueue(&analytics.UserAccessLog{})
|
||||
}
|
||||
if !risk_control.IsBufferFull() {
|
||||
t.Fatal("IsBufferFull() = false, want true")
|
||||
}
|
||||
|
||||
r := testhelper.NewTestGinEngine(risk_control.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"], "系统繁忙")
|
||||
})
|
||||
}
|
||||
@@ -1,3 +1,6 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
// Package risk_control provides the access control, IP rate limiting, and telemetry risk analysis domain plugin for Cordis.
|
||||
package risk_control
|
||||
|
||||
@@ -6,7 +9,6 @@ import (
|
||||
|
||||
"github.com/Rain-kl/Wavelet/core"
|
||||
"github.com/Rain-kl/Wavelet/core/extpoints"
|
||||
"github.com/Rain-kl/Wavelet/internal/apps/risk_control"
|
||||
"github.com/gin-gonic/gin"
|
||||
)
|
||||
|
||||
@@ -54,12 +56,12 @@ func (p *Plugin) Manifest() core.Manifest {
|
||||
// Apply registers risk control middlewares, settings, and cleanup hooks into the Context.
|
||||
func (p *Plugin) Apply(ctx *core.Context) error {
|
||||
// 1. Initialize LogWriter if needed
|
||||
risk_control.InitLogWriter(ctx.GoContext())
|
||||
InitLogWriter(ctx.GoContext())
|
||||
|
||||
// 2. Register router middleware
|
||||
mw := p.middleware
|
||||
if mw == nil {
|
||||
mw = risk_control.RiskControlMiddleware()
|
||||
mw = RiskControlMiddleware()
|
||||
}
|
||||
ctx.Router().Use(mw)
|
||||
|
||||
@@ -81,7 +83,7 @@ func (p *Plugin) Apply(ctx *core.Context) error {
|
||||
|
||||
// 4. Register lifecycle disposal cleanup
|
||||
ctx.OnDispose(func() error {
|
||||
return risk_control.StopLogWriter(context.Background())
|
||||
return StopLogWriter(context.Background())
|
||||
})
|
||||
|
||||
return nil
|
||||
|
||||
Reference in New Issue
Block a user