mirror of
https://github.com/Rain-kl/OpenFlare.git
synced 2026-10-07 08:06:37 +08:00
refactor(plugins): restructure admin and message_gateway into standard layered sub-packages
This commit is contained in:
@@ -0,0 +1,10 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package risk_control
|
||||
|
||||
// 模块内专用错误文案常量,集中在此文件维护,禁止在 handler 中内联。
|
||||
const (
|
||||
// errSystemBusy 访问日志缓冲队列满载时的限流提示,刻意使用模糊文案避免泄露内部容量细节。
|
||||
errSystemBusy = "系统繁忙,请稍后再试"
|
||||
)
|
||||
@@ -32,7 +32,7 @@ func RiskControlMiddleware() gin.HandlerFunc {
|
||||
|
||||
// 1. 限流背压检测(检测本地缓冲队列是否已满)
|
||||
if IsBufferFull() {
|
||||
response.AbortTooManyRequests(c, "系统繁忙,请稍后再试")
|
||||
response.AbortTooManyRequests(c, errSystemBusy)
|
||||
return
|
||||
}
|
||||
|
||||
|
||||
@@ -122,80 +122,3 @@ func (p *Plugin) Apply(ctx *core.Context) error {
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
type riskControlServiceImpl struct{}
|
||||
|
||||
func (s *riskControlServiceImpl) QueryAccessLogs(ctx context.Context, filter contracts.AccessLogFilterDTO, page, pageSize int) ([]contracts.AccessLogDTO, uint64, error) {
|
||||
store, err := logstore.Active(ctx)
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
f := logstore.AccessLogFilter{
|
||||
UserIDs: filter.UserIDs,
|
||||
Path: filter.Path,
|
||||
StartTime: filter.StartTime,
|
||||
EndTime: filter.EndTime,
|
||||
}
|
||||
list, total, err := store.UserAccessLogs.List(ctx, f, page, pageSize)
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
items := make([]contracts.AccessLogDTO, len(list))
|
||||
for i, item := range list {
|
||||
items[i] = contracts.AccessLogDTO{
|
||||
ID: item.ID,
|
||||
UserID: item.UserID,
|
||||
IP: item.IP,
|
||||
UserAgent: item.UserAgent,
|
||||
Method: item.Method,
|
||||
Path: item.Path,
|
||||
Status: item.Status,
|
||||
Latency: item.Latency,
|
||||
CreatedAt: item.CreatedAt,
|
||||
}
|
||||
}
|
||||
return items, total, nil
|
||||
}
|
||||
|
||||
func (s *riskControlServiceImpl) QueryAccessLogStats(ctx context.Context, days int) ([]contracts.AccessLogDailyStatsDTO, error) {
|
||||
store, err := logstore.Active(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
trend, err := store.UserAccessLogs.GetDailyTrend(ctx, days)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
res := make([]contracts.AccessLogDailyStatsDTO, len(trend))
|
||||
for i, t := range trend {
|
||||
res[i] = contracts.AccessLogDailyStatsDTO{
|
||||
Date: t.Date,
|
||||
PV: t.Count,
|
||||
}
|
||||
}
|
||||
return res, nil
|
||||
}
|
||||
|
||||
func (s *riskControlServiceImpl) ActiveLogEngine(ctx context.Context) string {
|
||||
store, err := logstore.Active(ctx)
|
||||
if err != nil {
|
||||
return "sqlite"
|
||||
}
|
||||
active, err := store.Status.ActiveDatabase(ctx)
|
||||
if err != nil {
|
||||
return "sqlite"
|
||||
}
|
||||
return active
|
||||
}
|
||||
|
||||
func (s *riskControlServiceImpl) IsLogEngineMigrating(ctx context.Context) bool {
|
||||
return logstore.Migrating(ctx)
|
||||
}
|
||||
|
||||
func (s *riskControlServiceImpl) Drain(ctx context.Context) error {
|
||||
return Drain(ctx)
|
||||
}
|
||||
|
||||
func (s *riskControlServiceImpl) SwitchLogEngine(ctx context.Context, targetEngine string) error {
|
||||
return MigrateAndSwitchEngine(ctx, targetEngine, nil)
|
||||
}
|
||||
|
||||
@@ -0,0 +1,124 @@
|
||||
// Copyright 2026 Arctel.net
|
||||
// SPDX-License-Identifier: Apache-2.0
|
||||
|
||||
package risk_control
|
||||
|
||||
import (
|
||||
"Wavelet/plugins/domain/risk_control/logstore"
|
||||
"context"
|
||||
)
|
||||
|
||||
// Repository 层:本插件根包内唯一的持久化访问入口。
|
||||
//
|
||||
// 真正的 SQL / 驱动实现位于 logstore 子包(受 logstore skill 约束的存储抽象),
|
||||
// 本文件负责解析当前生效日志库并转发读写、迁移与查询,使 service.go 只做用例编排。
|
||||
|
||||
// accessLogMigrationBatchSize 是日志引擎迁移时单批搬运的行数。
|
||||
const accessLogMigrationBatchSize = 1000
|
||||
|
||||
// writeAccessLogBatch 持久化批写缓冲队列中取出的一批访问日志。
|
||||
func writeAccessLogBatch(ctx context.Context, items []*logstore.UserAccessLog) error {
|
||||
rows := make([]logstore.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)
|
||||
}
|
||||
|
||||
// listAccessLogs 按过滤条件读取一页访问日志。
|
||||
func listAccessLogs(ctx context.Context, filter logstore.AccessLogFilter, page, pageSize int) ([]logstore.UserAccessLog, uint64, error) {
|
||||
store, err := logstore.Active(ctx)
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
return store.UserAccessLogs.List(ctx, filter, page, pageSize)
|
||||
}
|
||||
|
||||
// accessLogDailyTrend 读取最近 days 天的按天访问趋势。
|
||||
func accessLogDailyTrend(ctx context.Context, days int) ([]logstore.DailyTrend, error) {
|
||||
store, err := logstore.Active(ctx)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
return store.UserAccessLogs.GetDailyTrend(ctx, days)
|
||||
}
|
||||
|
||||
// activeLogDatabase 返回当前生效日志库的引擎标识。
|
||||
func activeLogDatabase(ctx context.Context) (string, error) {
|
||||
store, err := logstore.Active(ctx)
|
||||
if err != nil {
|
||||
return "", err
|
||||
}
|
||||
return store.Status.ActiveDatabase(ctx)
|
||||
}
|
||||
|
||||
// logStoreMigrating 报告日志库是否处于迁移冻结期。
|
||||
func logStoreMigrating(ctx context.Context) bool {
|
||||
return logstore.Migrating(ctx)
|
||||
}
|
||||
|
||||
// loadMigrationStores 解析迁移源(当前生效库)与目标引擎库。
|
||||
func loadMigrationStores(ctx context.Context, targetEngine string) (src, dst *logstore.Store, err error) {
|
||||
src, err = logstore.Active(ctx)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
dst, err = logstore.BuildForMigration(ctx, targetEngine)
|
||||
if err != nil {
|
||||
return nil, nil, err
|
||||
}
|
||||
return src, dst, nil
|
||||
}
|
||||
|
||||
// copyAccessLogs 清空目标库、按源库时间范围预建分区后分批搬运全部源数据,
|
||||
// 并在每批完成后通过 reportProgress 回调累计已搬运行数。
|
||||
func copyAccessLogs(ctx context.Context, src, dst *logstore.Store, reportProgress func(copied int)) error {
|
||||
if _, err := dst.UserAccessLogs.DeleteAll(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
from, to, err := src.UserAccessLogs.MigrationRange(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if !from.IsZero() && !to.IsZero() {
|
||||
if err := dst.UserAccessLogs.EnsurePartitions(ctx, from, to); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
|
||||
var afterID uint64
|
||||
var copied int
|
||||
for {
|
||||
rows, err := src.UserAccessLogs.ListForMigration(ctx, afterID, accessLogMigrationBatchSize)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if len(rows) == 0 {
|
||||
break
|
||||
}
|
||||
if err := dst.UserAccessLogs.BatchInsert(ctx, rows); err != nil {
|
||||
return err
|
||||
}
|
||||
afterID = rows[len(rows)-1].ID
|
||||
copied += len(rows)
|
||||
if reportProgress != nil {
|
||||
reportProgress(copied)
|
||||
}
|
||||
if len(rows) < accessLogMigrationBatchSize {
|
||||
break
|
||||
}
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// resetLogStoreCache 丢弃缓存的生效日志库,使下一次访问重新解析。
|
||||
func resetLogStoreCache() {
|
||||
logstore.InvalidateCache()
|
||||
}
|
||||
+74
-53
@@ -4,6 +4,7 @@
|
||||
package risk_control
|
||||
|
||||
import (
|
||||
"Wavelet/core/contracts"
|
||||
"Wavelet/pkg/batchwriter"
|
||||
"Wavelet/pkg/logger"
|
||||
"Wavelet/plugins/domain/risk_control/logstore"
|
||||
@@ -12,6 +13,9 @@ import (
|
||||
"time"
|
||||
)
|
||||
|
||||
// fallbackLogEngine 是日志库状态不可得时对外暴露的引擎标识。
|
||||
const fallbackLogEngine = "sqlite"
|
||||
|
||||
var (
|
||||
logWriterMu sync.RWMutex
|
||||
logWriter *batchwriter.Writer[*logstore.UserAccessLog]
|
||||
@@ -26,20 +30,7 @@ func InitLogWriter(ctx context.Context) {
|
||||
}
|
||||
|
||||
cfg := batchwriter.DefaultConfig()
|
||||
writer, err := batchwriter.New[*logstore.UserAccessLog](cfg, func(ctx context.Context, items []*logstore.UserAccessLog) error {
|
||||
rows := make([]logstore.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)
|
||||
},
|
||||
writer, err := batchwriter.New[*logstore.UserAccessLog](cfg, writeAccessLogBatch,
|
||||
batchwriter.WithDropHandler[*logstore.UserAccessLog](func(item *logstore.UserAccessLog) {
|
||||
path := ""
|
||||
if item != nil {
|
||||
@@ -144,49 +135,79 @@ func MigrateAndSwitchEngine(ctx context.Context, targetEngine string, reportProg
|
||||
if err := Drain(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
src, err := logstore.Active(ctx)
|
||||
src, dst, err := loadMigrationStores(ctx, targetEngine)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
dst, err := logstore.BuildForMigration(ctx, targetEngine)
|
||||
if err != nil {
|
||||
if err := copyAccessLogs(ctx, src, dst, reportProgress); err != nil {
|
||||
return err
|
||||
}
|
||||
if _, err := dst.UserAccessLogs.DeleteAll(ctx); err != nil {
|
||||
return err
|
||||
}
|
||||
from, to, err := src.UserAccessLogs.MigrationRange(ctx)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if !from.IsZero() && !to.IsZero() {
|
||||
if err := dst.UserAccessLogs.EnsurePartitions(ctx, from, to); err != nil {
|
||||
return err
|
||||
}
|
||||
}
|
||||
var afterID uint64
|
||||
var copied int
|
||||
const copyBatchSize = 1000
|
||||
for {
|
||||
rows, err := src.UserAccessLogs.ListForMigration(ctx, afterID, copyBatchSize)
|
||||
if err != nil {
|
||||
return err
|
||||
}
|
||||
if len(rows) == 0 {
|
||||
break
|
||||
}
|
||||
if err := dst.UserAccessLogs.BatchInsert(ctx, rows); err != nil {
|
||||
return err
|
||||
}
|
||||
afterID = rows[len(rows)-1].ID
|
||||
copied += len(rows)
|
||||
if reportProgress != nil {
|
||||
reportProgress(copied)
|
||||
}
|
||||
if len(rows) < copyBatchSize {
|
||||
break
|
||||
}
|
||||
}
|
||||
logstore.InvalidateCache()
|
||||
resetLogStoreCache()
|
||||
return nil
|
||||
}
|
||||
|
||||
// riskControlServiceImpl implements contracts.RiskControlService by orchestrating
|
||||
// the repository layer and mapping persistence rows into contract DTOs.
|
||||
type riskControlServiceImpl struct{}
|
||||
|
||||
func (s *riskControlServiceImpl) QueryAccessLogs(ctx context.Context, filter contracts.AccessLogFilterDTO, page, pageSize int) ([]contracts.AccessLogDTO, uint64, error) {
|
||||
list, total, err := listAccessLogs(ctx, logstore.AccessLogFilter{
|
||||
UserIDs: filter.UserIDs,
|
||||
Path: filter.Path,
|
||||
StartTime: filter.StartTime,
|
||||
EndTime: filter.EndTime,
|
||||
}, page, pageSize)
|
||||
if err != nil {
|
||||
return nil, 0, err
|
||||
}
|
||||
items := make([]contracts.AccessLogDTO, len(list))
|
||||
for i, item := range list {
|
||||
items[i] = contracts.AccessLogDTO{
|
||||
ID: item.ID,
|
||||
UserID: item.UserID,
|
||||
IP: item.IP,
|
||||
UserAgent: item.UserAgent,
|
||||
Method: item.Method,
|
||||
Path: item.Path,
|
||||
Status: item.Status,
|
||||
Latency: item.Latency,
|
||||
CreatedAt: item.CreatedAt,
|
||||
}
|
||||
}
|
||||
return items, total, nil
|
||||
}
|
||||
|
||||
func (s *riskControlServiceImpl) QueryAccessLogStats(ctx context.Context, days int) ([]contracts.AccessLogDailyStatsDTO, error) {
|
||||
trend, err := accessLogDailyTrend(ctx, days)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
res := make([]contracts.AccessLogDailyStatsDTO, len(trend))
|
||||
for i, t := range trend {
|
||||
res[i] = contracts.AccessLogDailyStatsDTO{
|
||||
Date: t.Date,
|
||||
PV: t.Count,
|
||||
}
|
||||
}
|
||||
return res, nil
|
||||
}
|
||||
|
||||
func (s *riskControlServiceImpl) ActiveLogEngine(ctx context.Context) string {
|
||||
engine, err := activeLogDatabase(ctx)
|
||||
if err != nil {
|
||||
return fallbackLogEngine
|
||||
}
|
||||
return engine
|
||||
}
|
||||
|
||||
func (s *riskControlServiceImpl) IsLogEngineMigrating(ctx context.Context) bool {
|
||||
return logStoreMigrating(ctx)
|
||||
}
|
||||
|
||||
func (s *riskControlServiceImpl) Drain(ctx context.Context) error {
|
||||
return Drain(ctx)
|
||||
}
|
||||
|
||||
func (s *riskControlServiceImpl) SwitchLogEngine(ctx context.Context, targetEngine string) error {
|
||||
return MigrateAndSwitchEngine(ctx, targetEngine, nil)
|
||||
}
|
||||
Reference in New Issue
Block a user