feat:add weblog to kafka

#12 #IAUI3L
This commit is contained in:
samwaf
2024-10-23 08:08:16 +08:00
parent 231f286fa7
commit 98f8bd5ce2
68 changed files with 945 additions and 470 deletions
+3 -3
View File
@@ -27,7 +27,7 @@ func (w *WafHostAPi) AddApi(c *gin.Context) {
_, svrOk := globalobj.GWAF_RUNTIME_OBJ_WAF_ENGINE.ServerOnline[req.Port]
if !svrOk && utils.PortCheck(req.Port) == false {
//发送websocket 推送消息
global.GQEQUE_MESSAGE_DB.PushBack(innerbean.OpResultMessageInfo{
global.GQEQUE_MESSAGE_DB.Enqueue(innerbean.OpResultMessageInfo{
BaseMessageInfo: innerbean.BaseMessageInfo{OperaType: "提示信息", Server: global.GWAF_CUSTOM_SERVER_NAME},
Msg: "端口被其他应用占用不能使用,如果使用的宝塔请在Samwaf系统管理-一键修改进行操作",
Success: "true",
@@ -129,7 +129,7 @@ func (w *WafHostAPi) ModifyHostApi(c *gin.Context) {
_, svrOk := globalobj.GWAF_RUNTIME_OBJ_WAF_ENGINE.ServerOnline[req.Port]
if !svrOk && utils.PortCheck(req.Port) == false {
//发送websocket 推送消息
global.GQEQUE_MESSAGE_DB.PushBack(innerbean.OpResultMessageInfo{
global.GQEQUE_MESSAGE_DB.Enqueue(innerbean.OpResultMessageInfo{
BaseMessageInfo: innerbean.BaseMessageInfo{OperaType: "提示信息", Server: global.GWAF_CUSTOM_SERVER_NAME},
Msg: "端口被其他应用占用不能使用,如果使用的宝塔请在Samwaf系统管理-一键修改进行操作",
Success: "true",
@@ -181,7 +181,7 @@ func (w *WafHostAPi) ModifyStartStatusApi(c *gin.Context) {
if req.START_STATUS == 0 && !svrOk && utils.PortCheck(wafHostOld.Port) == false {
//发送websocket 推送消息
global.GQEQUE_MESSAGE_DB.PushBack(innerbean.OpResultMessageInfo{
global.GQEQUE_MESSAGE_DB.Enqueue(innerbean.OpResultMessageInfo{
BaseMessageInfo: innerbean.BaseMessageInfo{OperaType: "提示信息", Server: global.GWAF_CUSTOM_SERVER_NAME},
Msg: "端口被其他应用占用不能使用,如果使用的宝塔请在Samwaf系统管理-一键修改进行操作",
Success: "true",
+1 -1
View File
@@ -1,11 +1,11 @@
package api
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/model/common/response"
"SamWaf/model/request"
response2 "SamWaf/model/response"
"SamWaf/utils/zlog"
"SamWaf/wafreg"
"github.com/gin-gonic/gin"
"io/ioutil"
+3 -3
View File
@@ -1,13 +1,13 @@
package api
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/innerbean"
"SamWaf/model/common/response"
"SamWaf/model/request"
response2 "SamWaf/model/response"
"SamWaf/utils"
"SamWaf/utils/zlog"
"SamWaf/wafdb"
"fmt"
"github.com/gin-gonic/gin"
@@ -92,7 +92,7 @@ func (w *WafLogAPi) ExportDBApi(c *gin.Context) {
downloadFilePath := filepath.Join(downLoadDir, downloadFileName)
err := wafdb.BackupDatabase(global.GWAF_LOCAL_LOG_DB, downloadFilePath)
if err != nil {
global.GQEQUE_MESSAGE_DB.PushBack(innerbean.OpResultMessageInfo{
global.GQEQUE_MESSAGE_DB.Enqueue(innerbean.OpResultMessageInfo{
BaseMessageInfo: innerbean.BaseMessageInfo{OperaType: "DOWNLOAD_LOG", Server: global.GWAF_CUSTOM_SERVER_NAME},
Msg: "导出失败",
Success: "true",
@@ -100,7 +100,7 @@ func (w *WafLogAPi) ExportDBApi(c *gin.Context) {
} else {
global.GWAF_RUNTIME_CURRENT_EXPORT_DB_LOG_FILE_PATH = downloadFilePath
//发送websocket 推送消息
global.GQEQUE_MESSAGE_DB.PushBack(innerbean.ExportResultMessageInfo{
global.GQEQUE_MESSAGE_DB.Enqueue(innerbean.ExportResultMessageInfo{
BaseMessageInfo: innerbean.BaseMessageInfo{OperaType: "DOWNLOAD_LOG", Server: global.GWAF_CUSTOM_SERVER_NAME},
Msg: "导出完毕",
Success: "true",
+1 -1
View File
@@ -61,7 +61,7 @@ func (w *WafLoginApi) LoginApi(c *gin.Context) {
//通知信息
noticeStr := fmt.Sprintf("登录IP:%s 归属地区:%s", c.ClientIP(), utils.GetCountry(c.ClientIP()))
global.GQEQUE_MESSAGE_DB.PushBack(innerbean.OperatorMessageInfo{
global.GQEQUE_MESSAGE_DB.Enqueue(innerbean.OperatorMessageInfo{
BaseMessageInfo: innerbean.BaseMessageInfo{OperaType: "登录信息"},
OperaCnt: noticeStr,
})
+1 -1
View File
@@ -1,13 +1,13 @@
package api
import (
"SamWaf/common/zlog"
"SamWaf/enums"
"SamWaf/global"
"SamWaf/model"
"SamWaf/model/common/response"
"SamWaf/model/request"
"SamWaf/model/spec"
"SamWaf/utils/zlog"
"errors"
"github.com/gin-gonic/gin"
"gorm.io/gorm"
+3 -3
View File
@@ -1,11 +1,11 @@
package api
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/innerbean"
"SamWaf/model"
"SamWaf/model/common/response"
"SamWaf/utils/zlog"
"SamWaf/wafupdate"
"github.com/gin-gonic/gin"
)
@@ -84,7 +84,7 @@ func (w *WafSysInfoApi) UpdateApi(c *gin.Context) {
wafDelayMsgService.Add("升级结果", "升级结果", "升级成功,当前版本为:"+global.GWAF_RUNTIME_NEW_VERSION+" 版本说明:"+global.GWAF_RUNTIME_NEW_VERSION_DESC)
global.GWAF_CHAN_UPDATE <- 1
//发送websocket 推送消息
global.GQEQUE_MESSAGE_DB.PushBack(innerbean.UpdateResultMessageInfo{
global.GQEQUE_MESSAGE_DB.Enqueue(innerbean.UpdateResultMessageInfo{
BaseMessageInfo: innerbean.BaseMessageInfo{OperaType: "系统即将重启", Server: global.GWAF_CUSTOM_SERVER_NAME},
Msg: "升级成功,等待重启",
Success: "true",
@@ -98,7 +98,7 @@ func (w *WafSysInfoApi) UpdateApi(c *gin.Context) {
global.GWAF_RUNTIME_IS_UPDATETING = false
//发送websocket 推送消息
global.GQEQUE_MESSAGE_DB.PushBack(innerbean.UpdateResultMessageInfo{
global.GQEQUE_MESSAGE_DB.Enqueue(innerbean.UpdateResultMessageInfo{
BaseMessageInfo: innerbean.BaseMessageInfo{OperaType: "升级结果", Server: global.GWAF_CUSTOM_SERVER_NAME},
Msg: "升级错误",
Success: "False",
+2
View File
@@ -3,6 +3,7 @@ package api
import (
"SamWaf/model/common/response"
"SamWaf/model/request"
"SamWaf/waftask"
"errors"
"github.com/gin-gonic/gin"
"gorm.io/gorm"
@@ -85,6 +86,7 @@ func (w *WafSystemConfigApi) ModifyApi(c *gin.Context) {
if err != nil {
response.FailWithMessage("编辑发生错误", c)
} else {
waftask.TaskLoadSetting(true)
response.OkWithMessage("编辑成功", c)
}
+1 -1
View File
@@ -1,9 +1,9 @@
package api
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/model"
"SamWaf/utils/zlog"
"encoding/json"
"github.com/gin-gonic/gin"
"github.com/gorilla/websocket"
@@ -1,4 +1,4 @@
package model
package gwebsocket
import (
"errors"
@@ -40,6 +40,17 @@ func (wafWebsocket *WebSocketOnline) GetWebSocket(key string) *Wssocket.Conn {
return nil
}
}
func (wafWebsocket *WebSocketOnline) GetAllWebSocket() map[string]*Wssocket.Conn {
wafWebsocket.Mux.Lock()
defer wafWebsocket.Mux.Unlock()
// 创建一个副本,防止外部修改
sockets := make(map[string]*Wssocket.Conn)
for k, v := range wafWebsocket.SocketMap {
sockets[k] = v
}
return sockets
}
func (wafWebsocket *WebSocketOnline) DelWebSocket(key string) error {
wafWebsocket.Mux.Lock()
defer wafWebsocket.Mux.Unlock()
+51
View File
@@ -0,0 +1,51 @@
package queue
import (
"container/list"
"sync"
)
type Queue struct {
data *list.List
mutex sync.RWMutex
}
func NewQueue() *Queue {
q := &Queue{data: list.New(), mutex: sync.RWMutex{}}
return q
}
func (q *Queue) Enqueue(v interface{}) {
q.mutex.Lock()
defer q.mutex.Unlock()
q.data.PushBack(v)
}
func (q *Queue) Dequeue() (interface{}, bool) {
q.mutex.Lock()
defer q.mutex.Unlock()
if q.data.Len() > 0 {
element := q.data.Front()
v := element.Value
q.data.Remove(element)
return v, true
} else {
return nil, false
}
}
func (q *Queue) Size() int {
q.mutex.RLock()
defer q.mutex.RUnlock()
return q.data.Len()
}
func (q *Queue) Empty() bool {
q.mutex.RLock()
defer q.mutex.RUnlock()
return q.data.Len() == 0
}
+2 -3
View File
@@ -1,7 +1,6 @@
package zlog
import (
"SamWaf/global"
"fmt"
"go.uber.org/zap"
"go.uber.org/zap/zapcore"
@@ -20,7 +19,7 @@ import (
// zlog.Warn("hello", zap.String("name", "Kevin"), zap.Any("arbitraryObj", dummyObject))
var logger *zap.Logger
func init() {
func InitZLog(releaseFlag string) {
encoderConfig := zap.NewProductionEncoderConfig()
encoderConfig.EncodeTime = zapcore.ISO8601TimeEncoder
encoder := zapcore.NewJSONEncoder(encoderConfig)
@@ -29,7 +28,7 @@ func init() {
//fileWriteSyncer = zapcore.AddSync(file)
fileWriteSyncer := getFileLogWriter()
if global.GWAF_RELEASE == "false" {
if releaseFlag == "false" {
core := zapcore.NewTee(
// 同时向控制台和文件写日志, 生产环境记得把控制台写入去掉
zapcore.NewCore(encoder, zapcore.AddSync(os.Stdout), zapcore.DebugLevel),
+10
View File
@@ -22,3 +22,13 @@ int64转成string
string := strconv.FormatInt(int64,10)
依赖检测:
C:\huawei\goproject\SamWafTools\go-cyclic\go-cyclic.exe run --dir .
依赖图:
C:\huawei\goproject\SamWafTools\godepgraph\godepgraph.exe ./
查询索引触发情况:
EXPLAIN QUERY PLAN select * from web_logs
+17 -8
View File
@@ -2,11 +2,13 @@ package global
import (
"SamWaf/cache"
"SamWaf/common/gwebsocket"
"SamWaf/common/queue"
"SamWaf/model"
"SamWaf/model/spec"
"SamWaf/wafnotify"
"SamWaf/wafsnowflake"
"github.com/bytedance/godlp/dlpheader"
Dequelib "github.com/edwingeng/deque"
"gorm.io/gorm"
"strconv"
"time"
@@ -76,14 +78,17 @@ var (
GCACHE_WECHAT_ACCESS string = "" //微信访问密钥
GCACHE_IP_CBUFF []byte // IP相关缓存
GDATA_DELETE_INTERVAL = 180 // 删除180天前的数据
GDATA_DELETE_INTERVAL int64 = 180 // 删除180天前的数据
/****队列相关*****/
GQEQUE_DB Dequelib.Deque //正常DB队列
GQEQUE_LOG_DB Dequelib.Deque //日志DB队列
GQEQUE_STATS_DB Dequelib.Deque //统计DB队列
GQEQUE_STATS_UPDATE_DB Dequelib.Deque //统计更新DB队列
GQEQUE_MESSAGE_DB Dequelib.Deque //发送消息队列
GQEQUE_DB *queue.Queue //正常DB队列
GQEQUE_LOG_DB *queue.Queue //日志DB队列
GQEQUE_STATS_DB *queue.Queue //统计DB队列
GQEQUE_STATS_UPDATE_DB *queue.Queue //统计更新DB队列
GQEQUE_MESSAGE_DB *queue.Queue //发送消息队列
/*******通知相关*************/
GNOTIFY_KAKFA_SERVICE *wafnotify.WafNotifyService //通知服务
/******数据库处理参数*****/
GDATA_BATCH_INSERT int = 1000 //最大批量插入
@@ -91,7 +96,7 @@ var (
GDATA_CURRENT_CHANGE bool = false //当前是否正在切换
GDATA_CURRENT_LOG_DB_MAP map[string]*gorm.DB //当前自定义的数据库连接 TODO 如果用户打开了多个 会不会影响内存
/******WebSocket*********/
GWebSocket *model.WebSocketOnline
GWebSocket *gwebsocket.WebSocketOnline
/******记录参数配置****************/
GCONFIG_RECORD_MAX_BODY_LENGTH int64 = 1024 * 2 //限制记录最大请求的body长度 record_max_req_body_length
@@ -99,6 +104,10 @@ var (
GCONFIG_RECORD_RESP int64 = 0 // 是否记录响应记录 record_resp
GCONFIG_RECORD_PROXY_HEADER string = "X-Forwarded-For,X-Real-IP" //配置获取IP头信息
GCONFIG_RECORD_AUTO_LOAD_SSL int64 = 1 //是否每天凌晨3点自动加载ssl证书
GCONFIG_RECORD_KAFKA_ENABLE int64 = 0 //kafka 是否激活
GCONFIG_RECORD_KAFKA_URL string = "127.0.0.1:9092" //kafka url地址
GCONFIG_RECORD_KAFKA_TOPIC string = "samwaf_logs_topic" //kafka topic
//升级相关
GUPDATE_VERSION_URL string = "https://update.samwaf.com/" //
+1
View File
@@ -11,4 +11,5 @@ var (
*/
GWAF_RUNTIME_OBJ_WAF_ENGINE *wafenginecore.WafEngine //当前引擎对象
GWAF_RUNTIME_OBJ_WAF_CRON *gocron.Scheduler //定时器
)
+13 -7
View File
@@ -1,6 +1,8 @@
module SamWaf
go 1.19
go 1.21
toolchain go1.21.4
require (
github.com/360EntSecGroup-Skylar/excelize v1.4.1
@@ -9,22 +11,25 @@ require (
github.com/denisbrodbeck/machineid v1.0.1
github.com/dsnet/compress v0.0.1
github.com/edwingeng/deque v1.0.3
github.com/emirpasic/gods v1.18.1
github.com/gin-gonic/gin v1.8.1
github.com/go-co-op/gocron v1.17.1
github.com/gorilla/websocket v1.5.0
github.com/hyperjumptech/grule-rule-engine v1.11.0
github.com/kardianos/service v1.2.2
github.com/lionsoul2014/ip2region/binding/golang v0.0.0-20220907060842-b2ba5d58e48d
github.com/pengge/go-wxsqlite3 v0.0.0-20231127082057-d869bc67f783
github.com/pengge/sqlitedriver v0.0.0-20231127095117-b0f000e40c2c
github.com/satori/go.uuid v1.2.0
github.com/shirou/gopsutil v3.21.11+incompatible
github.com/spf13/viper v1.13.0
github.com/stretchr/testify v1.8.4
github.com/twmb/franz-go v1.18.0
go.uber.org/zap v1.21.0
golang.org/x/mod v0.8.0
golang.org/x/net v0.9.0
golang.org/x/sys v0.15.0
golang.org/x/text v0.9.0
golang.org/x/net v0.21.0
golang.org/x/sys v0.20.0
golang.org/x/text v0.15.0
golang.org/x/time v0.0.0-20191024005414-555d28b269f0
gopkg.in/natefinch/lumberjack.v2 v2.0.0
gorm.io/gorm v1.25.2-0.20230530020048-26663ab9bf55
@@ -35,7 +40,6 @@ require (
github.com/antlr/antlr4/runtime/Go/antlr v0.0.0-20220527190237-ee62e23da966 // indirect
github.com/bmatcuk/doublestar v1.3.2 // indirect
github.com/davecgh/go-spew v1.1.1 // indirect
github.com/emirpasic/gods v1.12.0 // indirect
github.com/fsnotify/fsnotify v1.5.4 // indirect
github.com/gin-contrib/sse v0.1.0 // indirect
github.com/go-ole/go-ole v1.2.6 // indirect
@@ -51,6 +55,7 @@ require (
github.com/jinzhu/now v1.1.5 // indirect
github.com/json-iterator/go v1.1.12 // indirect
github.com/kevinburke/ssh_config v0.0.0-20190725054713-01f96b0aa0cd // indirect
github.com/klauspost/compress v1.17.8 // indirect
github.com/leodido/go-urn v1.2.1 // indirect
github.com/magiconair/properties v1.8.6 // indirect
github.com/mattn/go-isatty v0.0.14 // indirect
@@ -62,7 +67,7 @@ require (
github.com/mohae/deepcopy v0.0.0-20170929034955-c48cc78d4826 // indirect
github.com/pelletier/go-toml v1.9.5 // indirect
github.com/pelletier/go-toml/v2 v2.0.5 // indirect
github.com/pengge/go-wxsqlite3 v0.0.0-20231127082057-d869bc67f783 // indirect
github.com/pierrec/lz4/v4 v4.1.21 // indirect
github.com/pmezard/go-difflib v1.0.0 // indirect
github.com/robfig/cron/v3 v3.0.1 // indirect
github.com/sergi/go-diff v1.0.0 // indirect
@@ -75,12 +80,13 @@ require (
github.com/subosito/gotenv v1.4.1 // indirect
github.com/tklauser/go-sysconf v0.3.12 // indirect
github.com/tklauser/numcpus v0.6.1 // indirect
github.com/twmb/franz-go/pkg/kmsg v1.9.0 // indirect
github.com/ugorji/go/codec v1.2.7 // indirect
github.com/xanzy/ssh-agent v0.2.1 // indirect
github.com/yusufpapurcu/wmi v1.2.3 // indirect
go.uber.org/atomic v1.7.0 // indirect
go.uber.org/multierr v1.6.0 // indirect
golang.org/x/crypto v0.8.0 // indirect
golang.org/x/crypto v0.23.0 // indirect
golang.org/x/sync v0.0.0-20210220032951-036812b2e83c // indirect
google.golang.org/protobuf v1.28.0 // indirect
gopkg.in/ini.v1 v1.67.0 // indirect
+21 -10
View File
@@ -79,8 +79,9 @@ github.com/dsnet/compress v0.0.1/go.mod h1:Aw8dCMJ7RioblQeTqt88akK31OvO8Dhf5Jflh
github.com/dsnet/golib v0.0.0-20171103203638-1ea166775780/go.mod h1:Lj+Z9rebOhdfkVLjJ8T6VcRQv3SXugXy999NBtR9aFY=
github.com/edwingeng/deque v1.0.3 h1:gx+5OnQK8qMcYNUxcD/M76BT/LujVNVAZEP/MGF8WNw=
github.com/edwingeng/deque v1.0.3/go.mod h1:3Ys1pJhyVaB6iWigv4o2r6Ug1GZmfDWqvqmO6bjojg0=
github.com/emirpasic/gods v1.12.0 h1:QAUIPSaCu4G+POclxeqb3F+WPpdKqFGlw36+yOzGlrg=
github.com/emirpasic/gods v1.12.0/go.mod h1:YfzfFFoVP/catgzJb4IKIqXjX78Ha8FMSDh3ymbK86o=
github.com/emirpasic/gods v1.18.1 h1:FXtiHYKDGKCW2KzwZKx0iC0PQmdlorYgdFG9jPXJ1Bc=
github.com/emirpasic/gods v1.18.1/go.mod h1:8tpGGwCnJ5H4r6BWwaV6OrWmMoPhUl5jm/FMNAnJvWQ=
github.com/envoyproxy/go-control-plane v0.9.0/go.mod h1:YTl/9mNaCwkRvm6d1a2C3ymFceY/DCBVvsKhRF0iEA4=
github.com/envoyproxy/go-control-plane v0.9.1-0.20191026205805-5f8ba28d4473/go.mod h1:YTl/9mNaCwkRvm6d1a2C3ymFceY/DCBVvsKhRF0iEA4=
github.com/envoyproxy/go-control-plane v0.9.4/go.mod h1:6rpuAdCZL397s3pYoYcLgu1mIlRU8Am5FuJP05cCM98=
@@ -89,6 +90,7 @@ github.com/envoyproxy/go-control-plane v0.9.9-0.20201210154907-fd9021fe5dad/go.m
github.com/envoyproxy/protoc-gen-validate v0.1.0/go.mod h1:iSmxcyjqTsJpI2R4NaDN7+kN2VEUnK/pcBlmesArF7c=
github.com/flynn/go-shlex v0.0.0-20150515145356-3f9db97f8568/go.mod h1:xEzjJPgXI435gkrCt3MPfRiAkVrwSbHsst4LCFVfpJc=
github.com/frankban/quicktest v1.14.3 h1:FJKSZTDHjyhriyC81FLQ0LY93eSai0ZyR/ZIkd3ZUKE=
github.com/frankban/quicktest v1.14.3/go.mod h1:mgiwOwqx65TmIk1wJ6Q7wvnVMocbUorkibMOrVTHZps=
github.com/fsnotify/fsnotify v1.5.4 h1:jRbGcIw6P2Meqdwuo0H1p6JVLbL5DHKAKlYndzMwVZI=
github.com/fsnotify/fsnotify v1.5.4/go.mod h1:OVB6XrOHzAwXMpEM7uPOzcehqUV2UqJxmVXmkdnm1bU=
github.com/gin-contrib/sse v0.1.0 h1:Y/yl/+YNO8GZSjAhjMsSuLt29uWRFHdHYUb5lYOV9qE=
@@ -202,6 +204,8 @@ github.com/kevinburke/ssh_config v0.0.0-20190725054713-01f96b0aa0cd h1:Coekwdh0v
github.com/kevinburke/ssh_config v0.0.0-20190725054713-01f96b0aa0cd/go.mod h1:CT57kijsi8u/K/BOFA39wgDQJ9CxiF4nAY/ojJ6r6mM=
github.com/kisielk/gotool v1.0.0/go.mod h1:XhKaO+MFFWcvkIS/tQcRk01m1F5IRFswLeQ+oQHNcck=
github.com/klauspost/compress v1.4.1/go.mod h1:RyIbtBH6LamlWaDj8nUwkbUhJ87Yi3uG0guNDohfE1A=
github.com/klauspost/compress v1.17.8 h1:YcnTYrq7MikUT7k0Yb5eceMmALQPYBW/Xltxn0NAMnU=
github.com/klauspost/compress v1.17.8/go.mod h1:Di0epgTjJY877eYKx5yC51cX2A2Vl2ibi7bDH9ttBbw=
github.com/klauspost/cpuid v1.2.0/go.mod h1:Pj4uuM528wm8OyEC2QMXAi2YiTZ96dNQPGgoMS4s3ek=
github.com/kr/fs v0.1.0/go.mod h1:FFnZGqtBN9Gxj7eW1uZ42v5BccTP0vu6NEaFoC2HwRg=
github.com/kr/pretty v0.1.0/go.mod h1:dAy3ld7l9f0ibDNOQOHHMYYIIbhfbHSm3C4ZsoJORNo=
@@ -243,6 +247,8 @@ github.com/pengge/go-wxsqlite3 v0.0.0-20231127082057-d869bc67f783 h1:lHxAJqsiNfr
github.com/pengge/go-wxsqlite3 v0.0.0-20231127082057-d869bc67f783/go.mod h1:jlvEL0pxxqcz0f0N7bfNcsBWVhtGKmL73FdmLI9BofM=
github.com/pengge/sqlitedriver v0.0.0-20231127095117-b0f000e40c2c h1:NLpfD+rTkCVnGzNsXVq7crZ4RfspIDN/z19cVxUf0D0=
github.com/pengge/sqlitedriver v0.0.0-20231127095117-b0f000e40c2c/go.mod h1:TfNjocYjZ8S5JqYGrZKL4Pq7KLLPP4gMcBCtOjCMWvg=
github.com/pierrec/lz4/v4 v4.1.21 h1:yOVMLb6qSIDP67pl/5F7RepeKYu/VmTyEXvuMI5d9mQ=
github.com/pierrec/lz4/v4 v4.1.21/go.mod h1:gZWDp/Ze/IJXGXf23ltt2EXimqmTUXEy0GFuRQyBid4=
github.com/pkg/diff v0.0.0-20210226163009-20ebb0f2a09e/go.mod h1:pJLUxLENpZxwdsKMEsNbx1VGcRFpLqf3715MtcvvzbA=
github.com/pkg/errors v0.8.1/go.mod h1:bwawxfHBFNV+L2hUp1rHADufV3IMtnDRdf1r5NINEl0=
github.com/pkg/errors v0.9.1 h1:FEBLx1zS214owpjy7qsBeixbURkuhQAwrK5UwLGTwt4=
@@ -297,6 +303,10 @@ github.com/tklauser/go-sysconf v0.3.12 h1:0QaGUFOdQaIVdPgfITYzaTegZvdCjmYO52cSFA
github.com/tklauser/go-sysconf v0.3.12/go.mod h1:Ho14jnntGE1fpdOqQEEaiKRpvIavV0hSfmBq8nJbHYI=
github.com/tklauser/numcpus v0.6.1 h1:ng9scYS7az0Bk4OZLvrNXNSAO2Pxr1XXRAPyjhIx+Fk=
github.com/tklauser/numcpus v0.6.1/go.mod h1:1XfjsgE2zo8GVw7POkMbHENHzVg3GzmoZ9fESEdAacY=
github.com/twmb/franz-go v1.18.0 h1:25FjMZfdozBywVX+5xrWC2W+W76i0xykKjTdEeD2ejw=
github.com/twmb/franz-go v1.18.0/go.mod h1:zXCGy74M0p5FbXsLeASdyvfLFsBvTubVqctIaa5wQ+I=
github.com/twmb/franz-go/pkg/kmsg v1.9.0 h1:JojYUph2TKAau6SBtErXpXGC7E3gg4vGZMv9xFU/B6M=
github.com/twmb/franz-go/pkg/kmsg v1.9.0/go.mod h1:CMbfazviCyY6HM0SXuG5t9vOwYDHRCSrJJyBAe5paqg=
github.com/ugorji/go v1.2.7/go.mod h1:nF9osbDWLy6bDVv/Rtoh6QgnvNDpmCalQV5urGCCS6M=
github.com/ugorji/go/codec v1.2.7 h1:YPXUKf7fYbp/y8xloBqZOw2qaVggbfwMlI8WM3wZUJ0=
github.com/ugorji/go/codec v1.2.7/go.mod h1:WGN1fab3R1fzQlVQTkfxVtIBhWDRqOviHU95kRgeqEY=
@@ -334,8 +344,8 @@ golang.org/x/crypto v0.0.0-20200622213623-75b288015ac9/go.mod h1:LzIPMQfyMNhhGPh
golang.org/x/crypto v0.0.0-20210421170649-83a5a9bb288b/go.mod h1:T9bdIzuCu7OtxOm1hfPfRQxPLYneinmdGuTeoZ9dtd4=
golang.org/x/crypto v0.0.0-20210711020723-a769d52b0f97/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc=
golang.org/x/crypto v0.0.0-20211108221036-ceb1ce70b4fa/go.mod h1:GvvjBRRGRdwPK5ydBHafDWAxML/pGHZbMvKqRZ5+Abc=
golang.org/x/crypto v0.8.0 h1:pd9TJtTueMTVQXzk8E2XESSMQDj/U7OUu0PqJqPXQjQ=
golang.org/x/crypto v0.8.0/go.mod h1:mRqEX+O9/h5TFCrQhkgjo2yKi0yYA+9ecGkdQoHrywE=
golang.org/x/crypto v0.23.0 h1:dIJU/v2J8Mdglj/8rJ6UUOM3Zc9zLZxVZwwxMooUSAI=
golang.org/x/crypto v0.23.0/go.mod h1:CKFgDieR+mRhux2Lsu27y0fO304Db0wZe70UKqHu0v8=
golang.org/x/exp v0.0.0-20190121172915-509febef88a4/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA=
golang.org/x/exp v0.0.0-20190306152737-a1d7652674e8/go.mod h1:CJ0aWSM057203Lf6IL+f9T1iT9GByDxfZKAQTCR3kQA=
golang.org/x/exp v0.0.0-20190510132918-efd6b22b2522/go.mod h1:ZjyILWgesfNpC6sMxTJOJm9Kp84zZh5NQWvqDGG3Qr8=
@@ -405,8 +415,8 @@ golang.org/x/net v0.0.0-20201224014010-6772e930b67b/go.mod h1:m0MpNAwzfU5UDzcl9v
golang.org/x/net v0.0.0-20210226172049-e18ecbb05110/go.mod h1:m0MpNAwzfU5UDzcl9v0D8zg8gWTRqZa9RBIspLL5mdg=
golang.org/x/net v0.0.0-20210405180319-a5a99cb37ef4/go.mod h1:p54w0d4576C0XHj96bSt6lcn1PtDYWL6XObtHCRCNQM=
golang.org/x/net v0.0.0-20210805182204-aaa1db679c0d/go.mod h1:9nx3DQGgdP8bBQD5qxJ1jj9UTztislL4KSBs9R2vV5Y=
golang.org/x/net v0.9.0 h1:aWJ/m6xSmxWBx+V0XRHTlrYrPG56jKsLdTFmsSsCzOM=
golang.org/x/net v0.9.0/go.mod h1:d48xBJpPfHeWQsugry2m+kC02ZBRGRgulfHnEXEuWns=
golang.org/x/net v0.21.0 h1:AQyQV4dYCvJ7vGmJyKki9+PBdyvhkSd8EIx/qb0AYv4=
golang.org/x/net v0.21.0/go.mod h1:bIjVDfnllIU7BJ2DNgfnXvpSvtn8VRwhlsaeUTyUS44=
golang.org/x/oauth2 v0.0.0-20180821212333-d2e6202438be/go.mod h1:N/0e6XlmueqKjAGxoOufVs8QHGRruUQn6yWY3a++T0U=
golang.org/x/oauth2 v0.0.0-20190226205417-e64efc72b421/go.mod h1:gOpvHmFTYa4IltrdGE7lF6nIHvwfUNPOp7c8zoXwtLw=
golang.org/x/oauth2 v0.0.0-20190604053449-0f29369cfe45/go.mod h1:gOpvHmFTYa4IltrdGE7lF6nIHvwfUNPOp7c8zoXwtLw=
@@ -475,10 +485,11 @@ golang.org/x/sys v0.0.0-20210809222454-d867a43fc93e/go.mod h1:oPkhp1MJrh7nUepCBc
golang.org/x/sys v0.0.0-20220412211240-33da011f77ad/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.8.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.11.0/go.mod h1:oPkhp1MJrh7nUepCBck5+mAzfO9JrbApNNgaTdGDITg=
golang.org/x/sys v0.15.0 h1:h48lPFYpsTvQJZF4EKyI4aLHaev3CxivZmv7yZig9pc=
golang.org/x/sys v0.15.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
golang.org/x/sys v0.20.0 h1:Od9JTbYCk261bKm4M/mw7AklTlFYIa0bIp9BgSm1S8Y=
golang.org/x/sys v0.20.0/go.mod h1:/VUhepiaJMQUp4+oa/7Zr1D23ma6VTLIYjOOTFZPUcA=
golang.org/x/term v0.0.0-20201126162022-7de9c90e9dd1/go.mod h1:bj7SfCRtBDWHUb9snDiAeCFNEtKQo2Wmx5Cou7ajbmo=
golang.org/x/term v0.7.0 h1:BEvjmm5fURWqcfbSKTdpkDXYBrUS1c0m8agp14W48vQ=
golang.org/x/term v0.20.0 h1:VnkxpohqXaOBYJtBmEppKUG6mXpi+4O6purfc2+sMhw=
golang.org/x/term v0.20.0/go.mod h1:8UkIAJTvZgivsXaD6/pH6U9ecQzZ45awqEOzuCvwpFY=
golang.org/x/text v0.0.0-20170915032832-14c0d48ead0c/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
golang.org/x/text v0.3.0/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
golang.org/x/text v0.3.1-0.20180807135948-17ff2d5776d2/go.mod h1:NqM8EUOU14njkJ3fqMW+pc6Ldnwhi/IjpwHt7yyuwOQ=
@@ -486,8 +497,8 @@ golang.org/x/text v0.3.2/go.mod h1:bEr9sfX3Q8Zfm5fL9x+3itogRgK3+ptLWKqgva+5dAk=
golang.org/x/text v0.3.3/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
golang.org/x/text v0.3.4/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
golang.org/x/text v0.3.6/go.mod h1:5Zoc/QRtKVWzQhOtBMvqHzDpF6irO9z98xDceosuGiQ=
golang.org/x/text v0.9.0 h1:2sjJmO8cDvYveuX97RDLsxlyUxLl+GHoLxBiRdHllBE=
golang.org/x/text v0.9.0/go.mod h1:e1OnstbJyHTd6l/uOt8jFFHp6TRDWZR/bV3emEE/zU8=
golang.org/x/text v0.15.0 h1:h1V/4gjBv8v9cjcR6+AR5+/cIYK5N/WAgiv4xlsEtAk=
golang.org/x/text v0.15.0/go.mod h1:18ZOQIKpY8NJVqYksKHtTdi31H5itFRjB5/qKTNYzSU=
golang.org/x/time v0.0.0-20181108054448-85acf8d2951c/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ=
golang.org/x/time v0.0.0-20190308202827-9d24e82272b4/go.mod h1:tRJNPiyCQ0inRvYxbN9jk5I+vvW/OXSQhTDSoE431IQ=
golang.org/x/time v0.0.0-20191024005414-555d28b269f0 h1:/5xXl8Y5W96D+TtHSlonuFqGHIWVuyCkGJLwGh9JJFs=
+1 -1
View File
@@ -1,7 +1,7 @@
package innerbean
import (
"SamWaf/wechat"
"SamWaf/utils/wechat"
)
type BaseInfo interface {
+39 -27
View File
@@ -1,40 +1,52 @@
package innerbean
import "gorm.io/gorm"
type WebLog struct {
WafInnerDFlag string `json:"waf_inner_dflag"` //日志队列处理方式
HOST string `json:"host"`
URL string `json:"url"`
REFERER string `json:"referer"`
USER_AGENT string `json:"user_agent"`
METHOD string `json:"method"`
HEADER string `json:"header"`
SRC_IP string `json:"src_ip"`
SRC_PORT string `json:"src_port"`
COUNTRY string `json:"country"`
PROVINCE string `json:"province"`
CITY string `json:"city"`
CREATE_TIME string `gorm:"index:idx_weblog_time" json:"create_time"`
//CREATE_TIME string ` json:"create_time"`
WafInnerDFlag string `json:"waf_inner_dflag"` //日志队列处理方式
HOST string `json:"host"`
URL string `json:"url"`
REFERER string `json:"referer"`
USER_AGENT string `json:"user_agent"`
METHOD string `json:"method"`
HEADER string `json:"header"`
SRC_IP string `json:"src_ip"`
SRC_PORT string `json:"src_port"`
COUNTRY string `json:"country"`
PROVINCE string `json:"province"`
CITY string `json:"city"`
CREATE_TIME string `gorm:"index:idx_weblog_time" json:"create_time"`
CONTENT_LENGTH int64 `json:"content_length"`
COOKIES string `json:"cookies"`
BODY string `json:"body"`
REQ_UUID string `json:"req_uuid"`
USER_CODE string `json:"user_code"`
TenantId string `json:"tenant_id"` //租户ID(主要键)
HOST_CODE string `json:"host_code"` //主机ID (主要键)
Day int `json:"day"` //日 (主要键)
USER_CODE string `json:"user_code" gorm:"index"`
TenantId string `json:"tenant_id" gorm:"index"` //租户ID(主要键)
HOST_CODE string `json:"host_code" ` //主机ID (主要键)
Day int `json:"day"` //日 (主要键)
ACTION string `json:"action"`
RULE string `json:"rule"`
STATUS string `json:"status"` //状态
STATUS_CODE int `json:"status_code"` //状态编码
RES_BODY string `json:"res_body"` //返回信息
POST_FORM string `json:"post_form"` //提交的表单数据
TASK_FLAG int `json:"task_flag" gorm:"default:-1"` //任务处理标记 -1 等待处理 1 可以进行处理 (定时任务处理),2 处理完毕
UNIX_ADD_TIME int64 `json:"unix_add_time"` //添加日期unix
RISK_LEVEL int `json:"risk_level"` //危险等级 0:正常 1: 轻微 2: 有害 3: 严重 4: 特别严重
GUEST_IDENTIFICATION string `json:"guest_identification"` //访客身份识别
TimeSpent int64 `json:"time_spent"` //用时
STATUS string `json:"status"` //状态
STATUS_CODE int `json:"status_code"` //状态编码
RES_BODY string `json:"res_body"` //返回信息
POST_FORM string `json:"post_form"` //提交的表单数据
TASK_FLAG int `json:"task_flag" gorm:"default:-1;index"` //任务处理标记 -1 等待处理;1 可以进行处理2 处理完毕
UNIX_ADD_TIME int64 `json:"unix_add_time" gorm:"index"` //添加日期unix
RISK_LEVEL int `json:"risk_level"` //危险等级 0:正常 1:轻微 2:有害 3:严重 4:特别严重
GUEST_IDENTIFICATION string `json:"guest_identification"` //访客身份识别
TimeSpent int64 `json:"time_spent"` //用时
}
// 在 GORM 的 Model 方法中定义复合索引
func (WebLog) TableName() string {
return "web_logs"
}
func (WebLog) BeforeCreate(tx *gorm.DB) (err error) {
tx.Exec("CREATE INDEX IF NOT EXISTS idx_web_logs_task_flag_time ON web_logs (task_flag, unix_add_time)")
return
}
type WAFLog struct {
REQ_UUID string `json:"req_uuid"`
ACTION string `json:"action"`
+24 -10
View File
@@ -2,6 +2,8 @@ package main
import (
"SamWaf/cache"
"SamWaf/common/gwebsocket"
"SamWaf/common/zlog"
"SamWaf/enums"
"SamWaf/global"
"SamWaf/globalobj"
@@ -9,11 +11,12 @@ import (
"SamWaf/model"
"SamWaf/model/wafenginmodel"
"SamWaf/utils"
"SamWaf/utils/zlog"
"SamWaf/wafconfig"
"SamWaf/wafdb"
"SamWaf/wafenginecore"
"SamWaf/wafmangeweb"
"SamWaf/wafnotify"
"SamWaf/wafqueue"
"SamWaf/wafreg"
"SamWaf/wafsafeclear"
"SamWaf/wafsnowflake"
@@ -36,6 +39,7 @@ import (
"os/signal"
"path/filepath"
"runtime"
"runtime/debug"
"strconv"
"strings"
"sync/atomic"
@@ -76,15 +80,16 @@ func (m *wafSystenService) Stop(s service.Service) error {
// 守护协程
func NeverExit(name string, f func()) {
defer func() {
if v := recover(); v != nil { // 侦测到一个恐慌
stackBuf := make([]byte, 1024)
stackSize := runtime.Stack(stackBuf, false)
stackTrace := string(stackBuf[:stackSize])
if v := recover(); v != nil {
zlog.Info(fmt.Sprintf("协程%s崩溃了,准备重启一个。 : %v, Stack Trace: %s", name, v, debug.Stack()))
if global.GWAF_RELEASE == "false" {
debug.PrintStack()
}
zlog.Info(fmt.Sprintf("协程%s崩溃了,准备重启一个。 : %v, Stack Trace: %s", name, v, stackTrace))
go NeverExit(name, f) // 重启一个同功能协程
}
}()
zlog.Info(name + " start")
f()
}
@@ -191,10 +196,17 @@ func (m *wafSystenService) run() {
wafdb.InitLogDb("")
wafdb.InitStatsDb("")
//初始化队列引擎
wafenginecore.InitDequeEngine()
wafqueue.InitDequeEngine()
//启动队列消费
go NeverExit("ProcessDequeEngine", wafenginecore.ProcessDequeEngine)
go NeverExit("ProcessCoreDequeEngine", wafqueue.ProcessCoreDequeEngine)
go NeverExit("ProcessMessageDequeEngine", wafqueue.ProcessMessageDequeEngine)
go NeverExit("ProcessStatDequeEngine", wafqueue.ProcessStatDequeEngine)
go NeverExit("ProcessLogDequeEngine", wafqueue.ProcessLogDequeEngine)
//初始化一次系统参数信息
waftask.TaskLoadSetting(true)
//启动通知相关程序
global.GNOTIFY_KAKFA_SERVICE = wafnotify.InitNotifyKafkaEngine(global.GCONFIG_RECORD_KAFKA_ENABLE, global.GCONFIG_RECORD_KAFKA_URL, global.GCONFIG_RECORD_KAFKA_TOPIC) //kafka
//启动waf
globalobj.GWAF_RUNTIME_OBJ_WAF_ENGINE = &wafenginecore.WafEngine{
HostTarget: map[string]*wafenginmodel.HostSafe{},
@@ -218,7 +230,7 @@ func (m *wafSystenService) run() {
}()
//启动websocket
global.GWebSocket = model.InitWafWebSocket()
global.GWebSocket = gwebsocket.InitWafWebSocket()
//定时取规则并更新(考虑后期定时拉取公共规则 待定,可能会影响实际生产)
//定时器 (后期考虑是否独立包处理)
@@ -257,7 +269,7 @@ func (m *wafSystenService) run() {
})
// 获取参数
globalobj.GWAF_RUNTIME_OBJ_WAF_CRON.Every(1).Minutes().Do(func() {
go waftask.TaskLoadSetting()
go waftask.TaskLoadSetting(false)
})
@@ -565,6 +577,8 @@ func (m *wafSystenService) Graceful() {
}
func main() {
//初始化日志
zlog.InitZLog(global.GWAF_RELEASE)
if v := recover(); v != nil { // 侦测到一个恐慌
zlog.Info("主流程上被异常了")
}
+1 -1
View File
@@ -1,9 +1,9 @@
package middleware
import (
"SamWaf/common/zlog"
"SamWaf/model/common/response"
"SamWaf/service/waf_service"
"SamWaf/utils/zlog"
"github.com/gin-gonic/gin"
"strings"
)
+1 -1
View File
@@ -1,9 +1,9 @@
package middleware
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/service/waf_service"
"SamWaf/utils/zlog"
"bytes"
"github.com/gin-gonic/gin"
"io/ioutil"
@@ -1,9 +0,0 @@
package request
type WafSystemConfigAddReq struct {
Item string `json:"item" form:"item"` //item
Value string `json:"value" form:"value"` //value
ItemType string `json:"item_type" form:"item_type"` //item_type
Options string `json:"options" form:"options"` //options
Remarks string `json:"remarks" form:"remarks"` //备注
}
@@ -1,5 +0,0 @@
package request
type WafSystemConfigDelReq struct {
Id string `json:"id" form:"id"` //唯一键
}
@@ -1,5 +0,0 @@
package request
type WafSystemConfigDetailReq struct {
Id string `json:"id" form:"id"` //唯一键
}
@@ -1,10 +0,0 @@
package request
type WafSystemConfigEditReq struct {
Id string `json:"id"`
Item string `json:"item" form:"item"` //item
Value string `json:"value" form:"value"` //value
ItemType string `json:"item_type" form:"item_type"` //item_type
Options string `json:"options" form:"options"` //options
Remarks string `json:"remarks" form:"remarks"` //备注
}
+32
View File
@@ -0,0 +1,32 @@
package request
import "SamWaf/model/common/request"
type WafSystemConfigAddReq struct {
ItemClass string `json:"item_class"` //所属分类
Item string `json:"item" form:"item"` //item
Value string `json:"value" form:"value"` //value
ItemType string `json:"item_type" form:"item_type"` //item_type
Options string `json:"options" form:"options"` //options
Remarks string `json:"remarks" form:"remarks"` //备注
}
type WafSystemConfigDelReq struct {
Id string `json:"id" form:"id"` //唯一键
}
type WafSystemConfigDetailReq struct {
Id string `json:"id" form:"id"` //唯一键
}
type WafSystemConfigEditReq struct {
Id string `json:"id"`
ItemClass string `json:"item_class"` //所属分类
Item string `json:"item" form:"item"` //item
Value string `json:"value" form:"value"` //value
ItemType string `json:"item_type" form:"item_type"` //item_type
Options string `json:"options" form:"options"` //options
Remarks string `json:"remarks" form:"remarks"` //备注
}
type WafSystemConfigSearchReq struct {
Item string `json:"item" form:"item"`
Remarks string `json:"remarks" form:"remarks"`
request.PageInfo
}
@@ -1,9 +0,0 @@
package request
import "SamWaf/model/common/request"
type WafSystemConfigSearchReq struct {
Item string `json:"item" form:"item"`
Remarks string `json:"remarks" form:"remarks"`
request.PageInfo
}
+8 -7
View File
@@ -10,11 +10,12 @@ import (
*/
type SystemConfig struct {
baseorm.BaseOrm
IsSystem string `json:"is_system"` //是否是系统值
Item string `json:"item"`
Value string `json:"value"`
ItemType string `json:"item_type"` //配置类型 普通字符串,可选项
Options string `json:"options"` //如果ItemType == 可选项 这个地方有数据
HashInfo string `json:"hash_info"`
Remarks string `json:"remarks"` //备注
IsSystem string `json:"is_system"` //是否是系统值
ItemClass string `json:"item_class"` //所属分类
Item string `json:"item"`
Value string `json:"value"`
ItemType string `json:"item_type"` //配置类型 普通字符串,可选项
Options string `json:"options"` //如果ItemType == 可选项 这个地方有数据
HashInfo string `json:"hash_info"`
Remarks string `json:"remarks"` //备注
}
+15 -5
View File
@@ -37,6 +37,8 @@ func (receiver *WafLogService) GetListApi(req request.WafAttackLogSearch) ([]inn
splitFilterBys := strings.Split(req.FilterBy, "|")
splitFilterValues := strings.Split(req.FilterValue, "|")
/*强制索引*/
var forceIndex = "web_logs"
/*where条件*/
var whereField = ""
var whereValues []interface{}
@@ -96,6 +98,12 @@ func (receiver *WafLogService) GetListApi(req request.WafAttackLogSearch) ([]inn
}
}
}
//强制索引
{
if strings.Contains(whereField, "unix_add_time") {
forceIndex = "web_logs INDEXED BY idx_web_time_tenant_user_code"
}
}
//where字段赋值
{
@@ -141,13 +149,15 @@ func (receiver *WafLogService) GetListApi(req request.WafAttackLogSearch) ([]inn
return nil, 0, errors.New("输入排序字段不合法")
}
if len(req.CurrrentDbName) == 0 || req.CurrrentDbName == "local_log.db" {
global.GWAF_LOCAL_LOG_DB.Limit(req.PageSize).Where(whereField, whereValues...).Offset(req.PageSize * (req.PageIndex - 1)).Order(orderInfo).Find(&weblogs)
global.GWAF_LOCAL_LOG_DB.Model(&innerbean.WebLog{}).Where(whereField, whereValues...).Count(&total)
global.GWAF_LOCAL_LOG_DB.Table(forceIndex).Limit(req.PageSize).Where(whereField, whereValues...).Offset(req.PageSize * (req.PageIndex - 1)).Order(orderInfo).Find(&weblogs)
global.GWAF_LOCAL_LOG_DB.Table(forceIndex).Where(whereField, whereValues...).Count(&total)
/*global.GWAF_LOCAL_LOG_DB.Table("web_logs INDEXED BY idx_web_time_tenant_user_code ").Limit(req.PageSize).Where(whereField, whereValues...).Offset(req.PageSize * (req.PageIndex - 1)).Order(orderInfo).Find(&weblogs)
global.GWAF_LOCAL_LOG_DB.Table("web_logs INDEXED BY idx_web_time_tenant_user_code ").Model(&innerbean.WebLog{}).Where(whereField, whereValues...).Count(&total)
*/
} else {
wafdb.InitManaulLogDb("", req.CurrrentDbName)
global.GDATA_CURRENT_LOG_DB_MAP[req.CurrrentDbName].Debug().Limit(req.PageSize).Where(whereField, whereValues...).Offset(req.PageSize * (req.PageIndex - 1)).Order(orderInfo).Find(&weblogs)
global.GDATA_CURRENT_LOG_DB_MAP[req.CurrrentDbName].Model(&innerbean.WebLog{}).Where(whereField, whereValues...).Count(&total)
global.GDATA_CURRENT_LOG_DB_MAP[req.CurrrentDbName].Table(forceIndex).Limit(req.PageSize).Where(whereField, whereValues...).Offset(req.PageSize * (req.PageIndex - 1)).Order(orderInfo).Find(&weblogs)
global.GDATA_CURRENT_LOG_DB_MAP[req.CurrrentDbName].Table(forceIndex).Where(whereField, whereValues...).Count(&total)
}
return weblogs, total, nil
+1 -1
View File
@@ -1,12 +1,12 @@
package waf_service
import (
"SamWaf/common/zlog"
"SamWaf/customtype"
"SamWaf/global"
"SamWaf/model"
"SamWaf/model/baseorm"
"SamWaf/model/request"
"SamWaf/utils/zlog"
"errors"
uuid "github.com/satori/go.uuid"
"time"
+1 -1
View File
@@ -204,7 +204,7 @@ func (receiver *WafStatService) StatHomeRumtimeSysinfo() []response2.WafNameValu
data = append(data, response2.WafNameValue{Name: "软件版本Code", Value: fmt.Sprintf("%v", global.GWAF_RELEASE_VERSION)})
data = append(data, response2.WafNameValue{Name: "当前QPS", Value: fmt.Sprintf("%v", atomic.LoadUint64(&global.GWAF_RUNTIME_QPS))})
data = append(data, response2.WafNameValue{Name: "当前队列数", Value: fmt.Sprintf("主数据:%v 日志数据:%v 统计数据:%v 消息队列:%v", global.GQEQUE_DB.Len(), global.GQEQUE_LOG_DB.Len(), global.GQEQUE_STATS_DB.Len(), global.GQEQUE_MESSAGE_DB.Len())})
data = append(data, response2.WafNameValue{Name: "当前队列数", Value: fmt.Sprintf("主数据:%v 日志数据:%v 统计数据:%v 消息队列:%v", global.GQEQUE_DB.Size(), global.GQEQUE_LOG_DB.Size(), global.GQEQUE_STATS_DB.Size(), global.GQEQUE_MESSAGE_DB.Size())})
data = append(data, response2.WafNameValue{Name: "当前日志队列处理QPS", Value: fmt.Sprintf("%v", atomic.LoadUint64(&global.GWAF_RUNTIME_LOG_PROCESS))})
data = append(data, response2.WafNameValue{Name: "当前web端口使用列表", Value: fmt.Sprintf("%v", global.GWAF_RUNTIME_CURRENT_WEBPORT)})
+7 -5
View File
@@ -24,11 +24,12 @@ func (receiver *WafSystemConfigService) AddApi(wafSystemConfigAddReq request.Waf
CREATE_TIME: customtype.JsonTime(time.Now()),
UPDATE_TIME: customtype.JsonTime(time.Now()),
},
Item: wafSystemConfigAddReq.Item,
Value: wafSystemConfigAddReq.Value,
IsSystem: "0",
Remarks: wafSystemConfigAddReq.Remarks,
HashInfo: "",
ItemClass: wafSystemConfigAddReq.ItemClass,
Item: wafSystemConfigAddReq.Item,
Value: wafSystemConfigAddReq.Value,
IsSystem: "0",
Remarks: wafSystemConfigAddReq.Remarks,
HashInfo: "",
}
global.GWAF_LOCAL_DB.Create(bean)
return nil
@@ -45,6 +46,7 @@ func (receiver *WafSystemConfigService) ModifyApi(req request.WafSystemConfigEdi
}
editMap := map[string]interface{}{
"Item": req.Item,
"ItemClass": req.ItemClass,
"Value": req.Value,
"Remarks": req.Remarks,
"ItemType": req.ItemType,
+1 -1
View File
@@ -1,9 +1,9 @@
package utils
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/model"
"SamWaf/utils/zlog"
"fmt"
"github.com/lionsoul2014/ip2region/binding/golang/xdb"
"io/ioutil"
+2 -2
View File
@@ -1,10 +1,10 @@
package utils
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/innerbean"
"SamWaf/utils/zlog"
"SamWaf/wechat"
"SamWaf/utils/wechat"
)
type NotifyHelper struct {
+1 -1
View File
@@ -1,9 +1,9 @@
package utils
import (
"SamWaf/common/zlog"
"SamWaf/innerbean"
"SamWaf/model"
"SamWaf/utils/zlog"
"errors"
"github.com/hyperjumptech/grule-rule-engine/ast"
"github.com/hyperjumptech/grule-rule-engine/builder"
@@ -1,7 +1,7 @@
package wechat
import (
"SamWaf/utils/zlog"
"SamWaf/common/zlog"
"encoding/json"
"fmt"
"io/ioutil"
+1 -1
View File
@@ -1,8 +1,8 @@
package wafbot
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/utils/zlog"
"context"
"fmt"
"net"
+1 -1
View File
@@ -1,9 +1,9 @@
package wafconfig
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/utils"
"SamWaf/utils/zlog"
"github.com/denisbrodbeck/machineid"
uuid "github.com/satori/go.uuid"
"github.com/spf13/viper"
+4 -1
View File
@@ -1,13 +1,13 @@
package wafdb
import (
"SamWaf/common/zlog"
"SamWaf/customtype"
"SamWaf/global"
"SamWaf/innerbean"
"SamWaf/model"
"SamWaf/model/baseorm"
"SamWaf/utils"
"SamWaf/utils/zlog"
"context"
"database/sql"
"fmt"
@@ -139,12 +139,14 @@ func InitCoreDb(currentDir string) {
//SSL证书
db.AutoMigrate(&model.SslConfig{})
global.GWAF_LOCAL_DB.Callback().Query().Before("gorm:query").Register("tenant_plugin:before_query", before_query)
global.GWAF_LOCAL_DB.Callback().Query().Before("gorm:update").Register("tenant_plugin:before_update", before_update)
//重启需要删除无效规则
db.Where("user_code = ? and rule_status = 999", global.GWAF_USER_CODE).Delete(model.Rules{})
pathCoreSql(db)
}
}
@@ -181,6 +183,7 @@ func InitLogDb(currentDir string) {
global.GWAF_LOCAL_LOG_DB.Callback().Query().Before("gorm:query").Register("tenant_plugin:before_query", before_query)
global.GWAF_LOCAL_LOG_DB.Callback().Query().Before("gorm:update").Register("tenant_plugin:before_update", before_update)
pathLogSql(db)
var total int64 = 0
global.GWAF_LOCAL_DB.Model(&model.ShareDb{}).Count(&total)
if total == 0 {
+36
View File
@@ -0,0 +1,36 @@
package wafdb
import (
"SamWaf/common/zlog"
"gorm.io/gorm"
)
/*
*
一些后续补丁
*/
func pathLogSql(db *gorm.DB) {
// 20241018 创建联合索引 weblog
err := db.Exec("CREATE INDEX IF NOT EXISTS idx_web_logs_task_flag_time ON web_logs (task_flag, unix_add_time)").Error
if err != nil {
panic("failed to create index: " + err.Error())
} else {
zlog.Info("db", "idx_web_logs_task_flag_time created")
}
// 创建联合索引
err = db.Exec("CREATE INDEX IF NOT EXISTS idx_web_time_tenant_user_code ON web_logs (unix_add_time, tenant_id,user_code)").Error
if err != nil {
panic("failed to create index: " + err.Error())
} else {
zlog.Info("db", "idx_web_time_tenant_user_code created")
}
}
func pathCoreSql(db *gorm.DB) {
// 20241018 创建联合索引 weblog
err := db.Exec("UPDATE system_configs SET item_class = 'system' WHERE item_class IS NULL or item_class='' ").Error
if err != nil {
panic("failed to system_config :item_class " + err.Error())
} else {
zlog.Info("db", "system_config :item_class init successfully")
}
}
+1 -1
View File
@@ -1,10 +1,10 @@
package wafenginecore
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/innerbean"
"SamWaf/model/detection"
"SamWaf/utils/zlog"
"net/url"
)
+1 -1
View File
@@ -1,7 +1,7 @@
package wafenginecore
import (
"SamWaf/utils/zlog"
"SamWaf/common/zlog"
"strconv"
)
+1 -1
View File
@@ -1,9 +1,9 @@
package wafenginecore
import (
"SamWaf/common/zlog"
"SamWaf/innerbean"
"SamWaf/model"
"SamWaf/utils/zlog"
"SamWaf/wafproxy"
"context"
"crypto/tls"
+25 -21
View File
@@ -1,6 +1,7 @@
package wafenginecore
import (
"SamWaf/common/zlog"
"SamWaf/customtype"
"SamWaf/enums"
"SamWaf/global"
@@ -11,7 +12,6 @@ import (
"SamWaf/model/wafenginmodel"
"SamWaf/service/waf_service"
"SamWaf/utils"
"SamWaf/utils/zlog"
"SamWaf/wafenginecore/loadbalance"
"SamWaf/wafproxy"
"SamWaf/webplugin"
@@ -179,10 +179,11 @@ func (waf *WafEngine) ServeHTTP(w http.ResponseWriter, r *http.Request) {
if strings.Contains(r.Header.Get("Content-Type"), "application/x-www-form-urlencoded") {
// 解码 x-www-form-urlencoded 数据
formValuest, err := url.ParseQuery(weblogbean.BODY)
if err != nil {
fmt.Println("解码失败:", err)
} else {
if err == nil {
formValues = formValuest
} else {
fmt.Println("解码失败:", err)
fmt.Println("解码失败:", weblogbean.BODY)
}
}
@@ -350,7 +351,7 @@ func (waf *WafEngine) ServeHTTP(w http.ResponseWriter, r *http.Request) {
//记录响应body
weblogbean.RES_BODY = string(resBytes)
weblogbean.ACTION = "禁止"
global.GQEQUE_LOG_DB.PushBack(weblogbean)
global.GQEQUE_LOG_DB.Enqueue(weblogbean)
}
}
func (waf *WafEngine) getClientIP(r *http.Request, headers ...string) (error, string, string) {
@@ -379,13 +380,16 @@ func (waf *WafEngine) getClientIP(r *http.Request, headers ...string) (error, st
return nil, ip, port
}
func EchoErrorInfo(w http.ResponseWriter, r *http.Request, weblogbean innerbean.WebLog, ruleName string, blockInfo string) {
//发送微信推送消息
global.GQEQUE_MESSAGE_DB.PushBack(innerbean.RuleMessageInfo{
BaseMessageInfo: innerbean.BaseMessageInfo{OperaType: "命中保护规则", Server: global.GWAF_CUSTOM_SERVER_NAME},
Domain: weblogbean.HOST,
RuleInfo: ruleName,
Ip: fmt.Sprintf("%s (%s)", weblogbean.SRC_IP, utils.GetCountry(weblogbean.SRC_IP)),
})
go func() {
//发送推送消息
global.GQEQUE_MESSAGE_DB.Enqueue(innerbean.RuleMessageInfo{
BaseMessageInfo: innerbean.BaseMessageInfo{OperaType: "命中保护规则", Server: global.GWAF_CUSTOM_SERVER_NAME},
Domain: weblogbean.HOST,
RuleInfo: ruleName,
Ip: fmt.Sprintf("%s (%s)", weblogbean.SRC_IP, utils.GetCountry(weblogbean.SRC_IP)),
})
}()
resBytes := []byte("<html><head><title>您的访问被阻止</title></head><body><center><h1>" + blockInfo + "</h1> <br> 访问识别码:<h3>" + weblogbean.REQ_UUID + "</h3></center></body> </html>")
w.WriteHeader(403)
@@ -404,11 +408,11 @@ func EchoErrorInfo(w http.ResponseWriter, r *http.Request, weblogbean innerbean.
weblogbean.STATUS_CODE = 403
weblogbean.TASK_FLAG = 1
weblogbean.GUEST_IDENTIFICATION = "可疑用户"
global.GQEQUE_LOG_DB.PushBack(weblogbean)
global.GQEQUE_LOG_DB.Enqueue(weblogbean)
}
func (waf *WafEngine) errorResponse() func(http.ResponseWriter, *http.Request, error) {
return func(w http.ResponseWriter, req *http.Request, err error) {
zlog.Error("服务不可用 response:", zap.Any("err", err))
zlog.Error("服务不可用 response:", zap.Any("err", err.Error()))
resBytes := []byte("<html><head><title>服务不可用</title></head><body><center><h1>服务不可用</h1> <br><h3></h3></center></body> </html>")
w.WriteHeader(http.StatusServiceUnavailable)
@@ -551,7 +555,7 @@ func (waf *WafEngine) modifyResponse() func(*http.Response) error {
weblogfrist.TASK_FLAG = 1
if global.GWAF_RUNTIME_RECORD_LOG_TYPE == "all" {
if waf.HostTarget[host].Host.EXCLUDE_URL_LOG == "" {
global.GQEQUE_LOG_DB.PushBack(weblogfrist)
global.GQEQUE_LOG_DB.Enqueue(weblogfrist)
} else {
lines := strings.Split(waf.HostTarget[host].Host.EXCLUDE_URL_LOG, "\n")
isRecordLog := true
@@ -562,11 +566,11 @@ func (waf *WafEngine) modifyResponse() func(*http.Response) error {
}
}
if isRecordLog {
global.GQEQUE_LOG_DB.PushBack(weblogfrist)
global.GQEQUE_LOG_DB.Enqueue(weblogfrist)
}
}
} else if global.GWAF_RUNTIME_RECORD_LOG_TYPE == "abnormal" && weblogfrist.ACTION != "放行" {
global.GQEQUE_LOG_DB.PushBack(weblogfrist)
global.GQEQUE_LOG_DB.Enqueue(weblogfrist)
}
}
@@ -673,7 +677,7 @@ func (waf *WafEngine) StartWaf() {
OpType: "信息",
OpContent: "WAF启动",
}
global.GQEQUE_LOG_DB.PushBack(wafSysLog)
global.GQEQUE_LOG_DB.Enqueue(wafSysLog)
waf.StartAllProxyServer()
}
@@ -697,7 +701,7 @@ func (waf *WafEngine) CloseWaf() {
OpType: "信息",
OpContent: "WAF关闭",
}
global.GQEQUE_LOG_DB.PushBack(wafSysLog)
global.GQEQUE_LOG_DB.Enqueue(wafSysLog)
waf.EngineCurrentStatus = 0
waf.StopAllProxyServer()
@@ -985,7 +989,7 @@ func (waf *WafEngine) StartProxyServer(innruntime innerbean.ServerRunTime) {
OpType: "系统运行错误",
OpContent: "HTTPS端口被占用: " + strconv.Itoa(innruntime.Port) + ",请检查",
}
global.GQEQUE_LOG_DB.PushBack(wafSysLog)
global.GQEQUE_LOG_DB.Enqueue(wafSysLog)
zlog.Error("[HTTPServer] https server start fail, cause:[%v]", err)
}
zlog.Info("server https shutdown")
@@ -1023,7 +1027,7 @@ func (waf *WafEngine) StartProxyServer(innruntime innerbean.ServerRunTime) {
OpType: "系统运行错误",
OpContent: "HTTP端口被占用: " + strconv.Itoa(innruntime.Port) + ",请检查",
}
global.GQEQUE_LOG_DB.PushBack(wafSysLog)
global.GQEQUE_LOG_DB.Enqueue(wafSysLog)
zlog.Error("[HTTPServer] http server start fail, cause:[%v]", err)
}
zlog.Info("server http shutdown")
+8
View File
@@ -0,0 +1,8 @@
package wafinterface
import "SamWaf/innerbean"
type WafNotify interface {
NotifyBatch(logs []innerbean.WebLog) error // 处理多个日志
NotifySingle(log innerbean.WebLog) error // 处理单个日志
}
+3 -3
View File
@@ -1,10 +1,10 @@
package wafmangeweb
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/middleware"
"SamWaf/router"
"SamWaf/utils/zlog"
"SamWaf/wafmangeweb/static"
"context"
"errors"
@@ -29,7 +29,7 @@ func (web *WafWebManager) initRouter(r *gin.Engine) {
router.PublicApiGroupApp.InitCenterRouter(PublicRouterGroup) //注册中心接收接口
RouterGroup := r.Group("")
RouterGroup.Use(middleware.Auth(), middleware.CenterApi(), middleware.SecApi()) //TODO 中心管控 特定
RouterGroup.Use(middleware.Auth(), middleware.CenterApi(), middleware.SecApi(), middleware.GinGlobalExceptionMiddleWare()) //TODO 中心管控 特定
{
router.ApiGroupApp.InitHostRouter(RouterGroup)
router.ApiGroupApp.InitLogRouter(RouterGroup)
@@ -57,7 +57,7 @@ func (web *WafWebManager) initRouter(r *gin.Engine) {
router.ApiGroupApp.InitLoadBalanceRouter(RouterGroup)
router.ApiGroupApp.InitSslConfigRouter(RouterGroup)
}
r.Use(middleware.GinGlobalExceptionMiddleWare())
//r.Use(middleware.GinGlobalExceptionMiddleWare())
if global.GWAF_RELEASE == "true" {
static.Static(r, func(handlers ...gin.HandlerFunc) {
r.NoRoute(handlers...)
+1 -1
View File
@@ -1,9 +1,9 @@
package static
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/public"
"SamWaf/utils/zlog"
"errors"
"fmt"
"github.com/gin-gonic/gin"
+152
View File
@@ -0,0 +1,152 @@
package kafka
import (
"SamWaf/common/zlog"
"SamWaf/innerbean"
"context"
"encoding/json"
"fmt"
"github.com/twmb/franz-go/pkg/kgo"
"sync"
"time"
)
// KafkaNotifier 实现 WafNotify 接口
type KafkaNotifier struct {
client *kgo.Client
topic string
brokers []string
lastErrTime int64 //上次错误时间
}
// 初始化 Kafka 客户端
func NewKafkaNotifier(brokers []string, topic string, enable int64) (*KafkaNotifier, error) {
if enable == 0 {
return &KafkaNotifier{
client: nil,
topic: topic,
brokers: brokers,
}, nil
}
client, err := connect(brokers)
return &KafkaNotifier{
client: client,
topic: topic,
brokers: brokers,
}, err
}
func connect(brokers []string) (*kgo.Client, error) {
opts := []kgo.Opt{
kgo.SeedBrokers(brokers...),
kgo.RecordPartitioner(kgo.StickyKeyPartitioner(nil)),
}
client, err := kgo.NewClient(opts...)
if err != nil {
return nil, fmt.Errorf("failed to create Kafka client: %v", err)
}
ctx, _ := context.WithTimeout(context.Background(), 10*time.Second)
err = client.Ping(ctx)
if err != nil {
zlog.Error("failed to ping Kafka server: %v", err)
return nil, err
}
return client, nil
}
// 关闭 Kafka 客户端
func (kn *KafkaNotifier) Close() {
if kn.client != nil {
kn.client.Close()
}
}
// 关闭 Kafka 客户端
func (kn *KafkaNotifier) ReConnect() error {
//TODO 如果错误时间间隔在1,分钟之内,不进行重连操作
if kn.client != nil {
kn.client.Close()
}
client, err := connect(kn.brokers)
if err != nil {
zlog.Error("failed to reconnect create Kafka notifier: %v", err)
return err
}
kn.client = client
kn.lastErrTime = 0
return nil
}
// 实现 WafNotify 接口中的 NotifySingle 方法
func (kn *KafkaNotifier) NotifySingle(log innerbean.WebLog) error {
err := kn.sendMessage(log)
if err != nil {
return err
}
zlog.Debug("已发送")
return nil
}
// 实现 WafNotify 接口中的 NotifyBatch 方法
func (kn *KafkaNotifier) NotifyBatch(logs []innerbean.WebLog) error {
go func() {
for _, log := range logs {
if err := kn.sendMessage(log); err != nil {
zlog.Debug("发送批量日志时出错:", err.Error())
return
}
}
}()
return nil
}
// 发送消息到 Kafka
func (kn *KafkaNotifier) sendMessage(log innerbean.WebLog) error {
if kn.client == nil {
err := kn.ReConnect()
if err != nil {
return err
}
}
// 将日志对象转换为 JSON
jsonData, err := json.Marshal(log)
if err != nil {
zlog.Error("failed to marshal log: %v", err)
//精简重新继续发
logSimple := innerbean.WebLog{ /* 初始化日志 */ }
logSimple.REQ_UUID = log.REQ_UUID
logSimple.CREATE_TIME = log.CREATE_TIME
logSimple.ACTION = log.ACTION
logSimple.SRC_IP = log.SRC_IP
logSimple.SRC_PORT = log.SRC_PORT
logSimple.HOST_CODE = log.HOST_CODE
logSimple.GUEST_IDENTIFICATION = log.GUEST_IDENTIFICATION
logSimple.RISK_LEVEL = log.RISK_LEVEL
jsonData, err = json.Marshal(logSimple)
if err != nil {
zlog.Error("failed to marshal simple log: %v", err)
return err
}
}
record := &kgo.Record{
Topic: kn.topic,
Value: jsonData,
}
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second) // 为每次操作创建子 context
defer cancel()
var wg sync.WaitGroup
wg.Add(1)
kn.client.Produce(ctx, record, func(r *kgo.Record, err error) {
defer wg.Done()
if err != nil {
zlog.Error("failed to produce message: %v\n", err.Error())
} else {
zlog.Debug("message produced successfully")
}
})
wg.Wait()
return nil
}
+69
View File
@@ -0,0 +1,69 @@
package kafka
import (
"SamWaf/innerbean"
"fmt"
"testing"
"time"
)
func TestKafka(t *testing.T) {
// 初始化 Kafka
brokers := []string{"localhost:9092"}
topic := "samwaf_logs_topic"
kafkaNotifier, err := NewKafkaNotifier(brokers, topic)
if err != nil {
fmt.Println("Failed to initialize Kafka:", err)
return
}
defer kafkaNotifier.Close()
done := make(chan error)
log := innerbean.WebLog{ /* 初始化日志 */ }
log.ACTION = "通过"
log.BODY = "测试"
log.SRC_IP = "127.0.0.3"
//logs := []innerbean.WebLog{log, /* 其他日志 */ }
kafkaNotifier.NotifySingle(log)
go func() {
err := <-done
if err != nil {
fmt.Println("Failed to send single log to Kafka:", err)
} else {
fmt.Println("Single log sent to Kafka successfully.")
}
}()
// 保持主进程继续运行,直到所有异步任务完成
time.Sleep(10 * time.Second)
}
func BenchmarkKafka(b *testing.B) {
b.ReportAllocs()
// 初始化 Kafka
brokers := []string{"localhost:9092"}
topic := "samwaf_logs_topic"
kafkaNotifier, err := NewKafkaNotifier(brokers, topic)
if err != nil {
fmt.Println("Failed to initialize Kafka:", err)
return
}
defer kafkaNotifier.Close()
total := b.N
fmt.Println(total)
for i := 0; i < b.N; i++ {
datetimeNow := time.Now()
log := innerbean.WebLog{ /* 初始化日志 */ }
log.ACTION = "通过"
log.BODY = "测试"
log.SRC_IP = "127.0.0.3"
log.CREATE_TIME = datetimeNow.Format("2006-01-02 15:04:05")
log.UNIX_ADD_TIME = datetimeNow.UnixNano() / 1e6
//logs := []innerbean.WebLog{log, /* 其他日志 */ }
kafkaNotifier.NotifySingle(log)
}
// 保持主进程继续运行,直到所有异步任务完成
}
+19
View File
@@ -0,0 +1,19 @@
package wafnotify
import (
"SamWaf/wafnotify/kafka"
"fmt"
"strings"
)
func InitNotifyKafkaEngine(enable int64, url string, topic string) *WafNotifyService {
brokers := strings.Split(url, ",") // Kafka brokers
notifier, err := kafka.NewKafkaNotifier(brokers, topic, enable)
if err != nil {
fmt.Printf("Failed to create notifier: %v\n", err)
return NewWafNotifyService(notifier, enable)
}
// 创建日志服务并注入 notifier
return NewWafNotifyService(notifier, enable)
}
View File
+62
View File
@@ -1 +1,63 @@
package wafnotify
import (
"SamWaf/common/zlog"
"SamWaf/innerbean"
"SamWaf/wafinterface"
"fmt"
)
type WafNotifyService struct {
notifier wafinterface.WafNotify
enable int64
isStart int64
}
// 初始化 NotifyService
func NewWafNotifyService(notifier wafinterface.WafNotify, enable int64) *WafNotifyService {
return &WafNotifyService{notifier: notifier, enable: enable}
}
// 改激活状态
func (ls *WafNotifyService) ChangeEnable(enable int64) {
if ls == nil {
return
}
ls.enable = enable
}
// 处理并发送单条日志
func (ls *WafNotifyService) ProcessSingleLog(log innerbean.WebLog) error {
if ls.enable == 0 {
return nil
}
if ls.notifier == nil {
zlog.Debug("not init kafka")
return nil
}
// 日志处理逻辑
zlog.Debug("Processing single log...")
// 通知处理
if err := ls.notifier.NotifySingle(log); err != nil {
return fmt.Errorf("failed to notify single log: %v", err)
}
return nil
}
// 处理并发送多条日志
func (ls *WafNotifyService) ProcessBatchLogs(logs []innerbean.WebLog) error {
if ls.enable == 0 {
zlog.Debug("kafka没有开启")
return nil
}
// 日志处理逻辑
zlog.Debug("Processing batch logs...")
// 通知处理
if err := ls.notifier.NotifyBatch(logs); err != nil {
return fmt.Errorf("failed to notify batch logs: %v", err)
}
return nil
}
View File
+1 -1
View File
@@ -55,7 +55,7 @@ func OneKeyModifyBt(btSavePath string) (error, string) {
// 读取文件内容
afterContent, _ := ioutil.ReadFile(filePath)
//插入记录
global.GQEQUE_LOG_DB.PushBack(model.OneKeyMod{
global.GQEQUE_LOG_DB.Enqueue(model.OneKeyMod{
BaseOrm: baseorm.BaseOrm{
Id: uuid.NewV4().String(),
USER_CODE: global.GWAF_USER_CODE,
+25
View File
@@ -0,0 +1,25 @@
package wafqueue
import (
"SamWaf/global"
"time"
)
/*
*
处理核心队列信息
*/
func ProcessCoreDequeEngine() {
for {
for !global.GQEQUE_DB.Empty() {
bean, ok := global.GQEQUE_DB.Dequeue()
if ok {
if bean != nil {
global.GWAF_LOCAL_DB.Create(bean)
}
}
}
time.Sleep(100 * time.Millisecond)
}
}
+18
View File
@@ -0,0 +1,18 @@
package wafqueue
import (
"SamWaf/common/queue"
"SamWaf/global"
)
/*
*
初始化队列
*/
func InitDequeEngine() {
global.GQEQUE_DB = queue.NewQueue()
global.GQEQUE_LOG_DB = queue.NewQueue()
global.GQEQUE_STATS_DB = queue.NewQueue()
global.GQEQUE_STATS_UPDATE_DB = queue.NewQueue()
global.GQEQUE_MESSAGE_DB = queue.NewQueue()
}
+48
View File
@@ -0,0 +1,48 @@
package wafqueue
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/innerbean"
"sync/atomic"
"time"
)
/*
*
处理Log队列信息
*/
func ProcessLogDequeEngine() {
for {
global.GWAF_MEASURE_PROCESS_DEQUEENGINE.WriteData(time.Now().UnixNano() / 1e6)
if global.GDATA_CURRENT_CHANGE {
//如果正在切换库 跳过
zlog.Debug("正在切换数据库等待中队列")
} else {
var webLogArray []innerbean.WebLog
for !global.GQEQUE_LOG_DB.Empty() {
atomic.AddUint64(&global.GWAF_RUNTIME_LOG_PROCESS, 1) // 原子增加计数器
weblogbean, ok := global.GQEQUE_LOG_DB.Dequeue()
if !ok {
continue
}
if weblogbean != nil {
// 进行类型断言将其转为具体的结构
if logValue, ok := weblogbean.(innerbean.WebLog); ok {
webLogArray = append(webLogArray, logValue)
} else {
//插入其他类型内容
global.GWAF_LOCAL_LOG_DB.Create(weblogbean)
}
}
}
if len(webLogArray) > 0 {
global.GWAF_LOCAL_LOG_DB.CreateInBatches(webLogArray, global.GDATA_BATCH_INSERT)
global.GNOTIFY_KAKFA_SERVICE.ProcessBatchLogs(webLogArray)
}
}
time.Sleep(100 * time.Millisecond)
global.GWAF_MEASURE_PROCESS_DEQUEENGINE.WriteData(time.Now().UnixNano() / 1e6)
}
}
@@ -1,83 +1,30 @@
package wafenginecore
package wafqueue
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/innerbean"
"SamWaf/model"
"SamWaf/utils"
"SamWaf/utils/zlog"
"SamWaf/wafsec"
"encoding/json"
"github.com/edwingeng/deque"
uuid "github.com/satori/go.uuid"
"sync/atomic"
"time"
)
/*
*
初始化队列
处理消息队列信息
*/
func InitDequeEngine() {
global.GQEQUE_DB = deque.NewDeque()
global.GQEQUE_LOG_DB = deque.NewDeque()
global.GQEQUE_STATS_DB = deque.NewDeque()
global.GQEQUE_STATS_UPDATE_DB = deque.NewDeque()
global.GQEQUE_MESSAGE_DB = deque.NewDeque()
}
/*
*
处理队列信息
*/
func ProcessDequeEngine() {
zlog.Info("ProcessDequeEngine start")
func ProcessMessageDequeEngine() {
for {
global.GWAF_MEASURE_PROCESS_DEQUEENGINE.WriteData(time.Now().UnixNano() / 1e6)
for !global.GQEQUE_DB.Empty() {
weblogbean := global.GQEQUE_DB.PopFront()
if weblogbean != nil {
global.GWAF_LOCAL_DB.Create(weblogbean)
}
}
if global.GDATA_CURRENT_CHANGE {
//如果正在切换库 跳过
zlog.Debug("正在切换数据库等待中队列")
} else {
var webLogArray []innerbean.WebLog
for !global.GQEQUE_LOG_DB.Empty() {
atomic.AddUint64(&global.GWAF_RUNTIME_LOG_PROCESS, 1) // 原子增加计数器
weblogbean := global.GQEQUE_LOG_DB.PopFront()
if weblogbean != nil {
// 进行类型断言将其转为具体的结构
if logValue, ok := weblogbean.(innerbean.WebLog); ok {
webLogArray = append(webLogArray, logValue)
} else {
//插入其他类型内容
global.GWAF_LOCAL_LOG_DB.Create(weblogbean)
}
}
}
if len(webLogArray) > 0 {
global.GWAF_LOCAL_LOG_DB.CreateInBatches(webLogArray, global.GDATA_BATCH_INSERT)
}
}
for !global.GQEQUE_STATS_DB.Empty() {
bean := global.GQEQUE_STATS_DB.PopFront()
global.GWAF_LOCAL_STATS_DB.Create(bean)
}
for !global.GQEQUE_STATS_UPDATE_DB.Empty() {
bean := global.GQEQUE_STATS_UPDATE_DB.PopFront()
// 进行类型断言将其转为具体的结构
if UpdateValue, ok := bean.(innerbean.UpdateModel); ok {
global.GWAF_LOCAL_STATS_DB.Model(UpdateValue.Model).Where(UpdateValue.Query,
UpdateValue.Args...).Updates(UpdateValue.Update)
}
}
for !global.GQEQUE_MESSAGE_DB.Empty() {
messageinfo := global.GQEQUE_MESSAGE_DB.PopFront().(interface{})
popFront, ok := global.GQEQUE_MESSAGE_DB.Dequeue()
if !ok {
zlog.Error("来得信息未空")
continue
}
messageinfo := popFront.(interface{})
isCanSend := false
switch messageinfo.(type) {
case innerbean.RuleMessageInfo:
@@ -101,44 +48,35 @@ func ProcessDequeEngine() {
} else {
utils.NotifyHelperApp.SendRuleInfo(rulemessage)
}
if rulemessage.BaseMessageInfo.OperaType == "命中保护规则" {
//发送websocket
for _, ws := range global.GWebSocket.GetAllWebSocket() {
}
if rulemessage.BaseMessageInfo.OperaType == "命中保护规则" {
//发送websocket
for _, ws := range global.GWebSocket.SocketMap {
if ws != nil {
//信息包体进行单独处理
/*msgBody:=model.MsgDataPacket{
MessageId: uuid.NewV4().String(),
MessageType: "命中保护规则",
MessageData: rulemessage.RuleInfo + rulemessage.Ip,
MessageAttach: nil,
MessageDateTime: time.Now().Format("2006-01-02 15:04:05"),
MessageUnReadStatus: true,
}*/
msgBody, _ := json.Marshal(model.MsgDataPacket{
MessageId: uuid.NewV4().String(),
MessageType: "命中保护规则",
MessageData: rulemessage.RuleInfo + rulemessage.Ip,
MessageAttach: nil,
MessageDateTime: time.Now().Format("2006-01-02 15:04:05"),
MessageUnReadStatus: true,
})
encryptStr, _ := wafsec.AesEncrypt(msgBody, global.GWAF_COMMUNICATION_KEY)
//写入ws数据
msgBytes, err := json.Marshal(model.MsgPacket{
MsgCode: "200",
MsgDataPacket: encryptStr,
MsgCmdType: "Info",
})
err = ws.WriteMessage(1, msgBytes)
if err != nil {
continue
if ws != nil {
msgBody, _ := json.Marshal(model.MsgDataPacket{
MessageId: uuid.NewV4().String(),
MessageType: "命中保护规则",
MessageData: rulemessage.RuleInfo + rulemessage.Ip,
MessageAttach: nil,
MessageDateTime: time.Now().Format("2006-01-02 15:04:05"),
MessageUnReadStatus: true,
})
encryptStr, _ := wafsec.AesEncrypt(msgBody, global.GWAF_COMMUNICATION_KEY)
//写入ws数据
msgBytes, err := json.Marshal(model.MsgPacket{
MsgCode: "200",
MsgDataPacket: encryptStr,
MsgCmdType: "Info",
})
err = ws.WriteMessage(1, msgBytes)
if err != nil {
continue
}
}
}
}
}
break
case innerbean.OperatorMessageInfo:
operatorMessage := messageinfo.(innerbean.OperatorMessageInfo)
@@ -147,7 +85,7 @@ func ProcessDequeEngine() {
case innerbean.ExportResultMessageInfo:
exportResult := messageinfo.(innerbean.ExportResultMessageInfo)
//发送websocket
for _, ws := range global.GWebSocket.SocketMap {
for _, ws := range global.GWebSocket.GetAllWebSocket() {
if ws != nil {
//信息包体进行单独处理
msgBody, _ := json.Marshal(model.MsgDataPacket{
@@ -177,7 +115,7 @@ func ProcessDequeEngine() {
//升级结果
updatemessage := messageinfo.(innerbean.UpdateResultMessageInfo)
//发送websocket
for _, ws := range global.GWebSocket.SocketMap {
for _, ws := range global.GWebSocket.GetAllWebSocket() {
if ws != nil {
//信息包体进行单独处理
msgBody, _ := json.Marshal(model.MsgDataPacket{
@@ -207,7 +145,7 @@ func ProcessDequeEngine() {
//操作实时结果
updatemessage := messageinfo.(innerbean.OpResultMessageInfo)
//发送websocket
for _, ws := range global.GWebSocket.SocketMap {
for _, ws := range global.GWebSocket.GetAllWebSocket() {
if ws != nil {
//信息包体进行单独处理
msgBody, _ := json.Marshal(model.MsgDataPacket{
@@ -237,6 +175,5 @@ func ProcessDequeEngine() {
}
time.Sleep(100 * time.Millisecond)
global.GWAF_MEASURE_PROCESS_DEQUEENGINE.WriteData(time.Now().UnixNano() / 1e6)
}
}
+34
View File
@@ -0,0 +1,34 @@
package wafqueue
import (
"SamWaf/global"
"SamWaf/innerbean"
"time"
)
/*
*
处理Stat队列信息
*/
func ProcessStatDequeEngine() {
for {
for !global.GQEQUE_STATS_DB.Empty() {
dequeue, ok := global.GQEQUE_STATS_DB.Dequeue()
if ok {
global.GWAF_LOCAL_STATS_DB.Create(dequeue)
}
}
for !global.GQEQUE_STATS_UPDATE_DB.Empty() {
bean, ok := global.GQEQUE_STATS_UPDATE_DB.Dequeue()
if ok {
// 进行类型断言将其转为具体的结构
if UpdateValue, ok := bean.(innerbean.UpdateModel); ok {
global.GWAF_LOCAL_STATS_DB.Model(UpdateValue.Model).Where(UpdateValue.Query,
UpdateValue.Args...).Updates(UpdateValue.Update)
}
}
}
time.Sleep(100 * time.Millisecond)
}
}
+1 -1
View File
@@ -1,8 +1,8 @@
package wafsafeclear
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/utils/zlog"
)
func SafeClear() {
+14 -11
View File
@@ -1,6 +1,7 @@
package waftask
import (
"SamWaf/common/zlog"
"SamWaf/customtype"
"SamWaf/global"
"SamWaf/innerbean"
@@ -8,9 +9,8 @@ import (
"SamWaf/model/baseorm"
"SamWaf/service/waf_service"
"SamWaf/utils"
"SamWaf/utils/zlog"
"SamWaf/utils/wechat"
"SamWaf/wafsec"
"SamWaf/wechat"
"encoding/json"
"fmt"
uuid "github.com/satori/go.uuid"
@@ -90,6 +90,9 @@ func TaskCounter() {
//一、 主机聚合统计
{
var resultHosts []CountHostResult
explain := global.GWAF_LOCAL_LOG_DB.Debug().Explain("SELECT host_code, user_code,tenant_id ,action,count(req_uuid) as count,day,host FROM \"web_logs\" where task_flag = ? and unix_add_time > ? GROUP BY host_code, user_code,action,tenant_id,day,host",
1, currenyDayMillisecondsBak)
zlog.Debug(explain)
global.GWAF_LOCAL_LOG_DB.Raw("SELECT host_code, user_code,tenant_id ,action,count(req_uuid) as count,day,host FROM \"web_logs\" where task_flag = ? and unix_add_time > ? GROUP BY host_code, user_code,action,tenant_id,day,host",
1, currenyDayMillisecondsBak).Scan(&resultHosts)
/****
@@ -116,7 +119,7 @@ func TaskCounter() {
Type: value.ACTION,
Count: value.Count,
}
global.GQEQUE_STATS_DB.PushBack(statDay2)
global.GQEQUE_STATS_DB.Enqueue(statDay2)
} else {
statDayMap := map[string]interface{}{
"Count": value.Count + statDay.Count,
@@ -128,7 +131,7 @@ func TaskCounter() {
Update: statDayMap,
}
updateBean.Args = append(updateBean.Args, value.TenantId, value.UserCode, value.HostCode, value.ACTION, value.Day)
global.GQEQUE_STATS_UPDATE_DB.PushBack(updateBean)
global.GQEQUE_STATS_UPDATE_DB.Enqueue(updateBean)
}
}
}
@@ -163,7 +166,7 @@ func TaskCounter() {
Count: value.Count,
IP: value.Ip,
}
global.GQEQUE_STATS_DB.PushBack(statDay2)
global.GQEQUE_STATS_DB.Enqueue(statDay2)
} else {
statDayMap := map[string]interface{}{
"Count": value.Count + statDay.Count,
@@ -176,7 +179,7 @@ func TaskCounter() {
Update: statDayMap,
}
updateBean.Args = append(updateBean.Args, value.TenantId, value.UserCode, value.HostCode, value.Ip, value.ACTION, value.Day)
global.GQEQUE_STATS_UPDATE_DB.PushBack(updateBean)
global.GQEQUE_STATS_UPDATE_DB.Enqueue(updateBean)
}
}
@@ -214,7 +217,7 @@ func TaskCounter() {
Province: value.Province,
City: value.City,
}
global.GQEQUE_STATS_DB.PushBack(statDay2)
global.GQEQUE_STATS_DB.Enqueue(statDay2)
} else {
statDayMap := map[string]interface{}{
"Count": value.Count + statDay.Count,
@@ -227,7 +230,7 @@ func TaskCounter() {
Update: statDayMap,
}
updateBean.Args = append(updateBean.Args, value.TenantId, value.UserCode, value.HostCode, value.Country, value.Province, value.City, value.ACTION, value.Day)
global.GQEQUE_STATS_UPDATE_DB.PushBack(updateBean)
global.GQEQUE_STATS_UPDATE_DB.Enqueue(updateBean)
}
}
@@ -256,7 +259,7 @@ func TaskStatusNotify() {
if err == nil {
noticeStr := fmt.Sprintf("今日访问量:%d 今天恶意访问量:%d 昨日恶意访问量:%d", statHomeInfo.VisitCountOfToday, statHomeInfo.AttackCountOfToday, statHomeInfo.AttackCountOfYesterday)
global.GQEQUE_MESSAGE_DB.PushBack(innerbean.OperatorMessageInfo{
global.GQEQUE_MESSAGE_DB.Enqueue(innerbean.OperatorMessageInfo{
BaseMessageInfo: innerbean.BaseMessageInfo{OperaType: "汇总通知"},
OperaCnt: noticeStr,
})
@@ -272,7 +275,7 @@ func TaskStatusNotify() {
*/
func TaskDeleteHistoryInfo() {
zlog.Debug("TaskDeleteHistoryInfo")
deleteBeforeDay := time.Now().AddDate(0, 0, -global.GDATA_DELETE_INTERVAL).Format("2006-01-02 15:04")
deleteBeforeDay := time.Now().AddDate(0, 0, -int(global.GDATA_DELETE_INTERVAL)).Format("2006-01-02 15:04")
waf_service.WafLogServiceApp.DeleteHistory(deleteBeforeDay)
}
@@ -289,7 +292,7 @@ func TaskDelayInfo() {
msg := models[i]
sendSuccess := 0
//发送websocket
for _, ws := range global.GWebSocket.SocketMap {
for _, ws := range global.GWebSocket.GetAllWebSocket() {
if ws != nil {
cmdType := "Info"
+1 -1
View File
@@ -1,8 +1,8 @@
package waftask
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/utils/zlog"
"SamWaf/wafdb"
"github.com/spf13/viper"
"testing"
+102 -178
View File
@@ -1,195 +1,119 @@
package waftask
import (
"SamWaf/common/zlog"
"SamWaf/global"
"SamWaf/model/request"
"SamWaf/utils/zlog"
"strconv"
)
// 加载配置数据
func TaskLoadSetting() {
zlog.Debug("TaskLoadSetting")
configItem := wafSystemConfigService.GetDetailByItem("record_max_req_body_length")
func setConfigIntValue(name string, value int64, change int) {
// 更新全局配置值
switch name {
case "record_max_req_body_length":
global.GCONFIG_RECORD_MAX_BODY_LENGTH = value
case "record_max_res_body_length":
global.GCONFIG_RECORD_MAX_RES_BODY_LENGTH = value
case "record_resp":
global.GCONFIG_RECORD_RESP = value
case "delete_history_log_day":
global.GDATA_DELETE_INTERVAL = value
case "log_db_size":
global.GDATA_SHARE_DB_SIZE = value
case "auto_load_ssl_file":
global.GCONFIG_RECORD_AUTO_LOAD_SSL = value
case "kafka_enable":
if global.GCONFIG_RECORD_KAFKA_ENABLE != value && global.GNOTIFY_KAKFA_SERVICE != nil {
global.GNOTIFY_KAKFA_SERVICE.ChangeEnable(value)
}
global.GCONFIG_RECORD_KAFKA_ENABLE = value
default:
zlog.Warn("Unknown config item:", name)
}
}
func setConfigStringValue(name string, value string, change int) {
// 更新全局配置值
switch name {
case "dns_server":
global.GWAF_RUNTIME_DNS_SERVER = value
case "record_log_type":
global.GWAF_RUNTIME_RECORD_LOG_TYPE = value
case "gwaf_center_enable":
global.GWAF_CENTER_ENABLE = value
case "gwaf_center_url":
global.GWAF_CENTER_URL = value
case "gwaf_proxy_header":
global.GCONFIG_RECORD_PROXY_HEADER = value
case "kafka_url":
global.GCONFIG_RECORD_KAFKA_URL = value
case "kafka_topic":
global.GCONFIG_RECORD_KAFKA_TOPIC = value
default:
zlog.Warn("Unknown config item:", name)
}
}
func updateConfigIntItem(initLoad bool, itemClass string, itemName string, defaultValue int64, remarks string, itemType string, options string) {
configItem := wafSystemConfigService.GetDetailByItem(itemName)
if configItem.Id != "" {
value, err := strconv.ParseInt(configItem.Value, 10, 0)
if err == nil {
if global.GCONFIG_RECORD_MAX_BODY_LENGTH != value {
global.GCONFIG_RECORD_MAX_BODY_LENGTH = value
}
if err == nil && defaultValue != value {
setConfigIntValue(itemName, value, 1)
} else if err == nil && initLoad == true {
setConfigIntValue(itemName, value, 0)
}
} else {
wafSystemConfigAddReq := request.WafSystemConfigAddReq{
Item: "record_max_req_body_length",
Value: strconv.FormatInt(global.GCONFIG_RECORD_MAX_BODY_LENGTH, 10),
Remarks: "记录请求最大报文",
ItemType: "int",
}
wafSystemConfigService.AddApi(wafSystemConfigAddReq)
}
configItem = wafSystemConfigService.GetDetailByItem("record_max_rep_body_length")
if configItem.Id != "" {
value, err := strconv.ParseInt(configItem.Value, 10, 0)
if err == nil {
if global.GCONFIG_RECORD_MAX_RES_BODY_LENGTH != value {
global.GCONFIG_RECORD_MAX_RES_BODY_LENGTH = value
}
}
} else {
wafSystemConfigAddReq := request.WafSystemConfigAddReq{
Item: "record_max_rep_body_length",
Value: strconv.FormatInt(global.GCONFIG_RECORD_MAX_RES_BODY_LENGTH, 10),
Remarks: "如果可以记录,满足最大响应报文大小才记录",
ItemType: "int",
}
wafSystemConfigService.AddApi(wafSystemConfigAddReq)
}
configItem = wafSystemConfigService.GetDetailByItem("record_resp")
if configItem.Id != "" {
value, err := strconv.ParseInt(configItem.Value, 10, 0)
if err == nil {
if global.GCONFIG_RECORD_RESP != value {
global.GCONFIG_RECORD_RESP = value
}
}
} else {
wafSystemConfigAddReq := request.WafSystemConfigAddReq{
Item: "record_resp",
Value: strconv.FormatInt(global.GCONFIG_RECORD_RESP, 10),
Remarks: "是否记录响应报文",
ItemType: "int",
}
wafSystemConfigService.AddApi(wafSystemConfigAddReq)
}
configItem = wafSystemConfigService.GetDetailByItem("delete_history_log_day")
if configItem.Id != "" {
value, err := strconv.Atoi(configItem.Value)
if err == nil {
if global.GDATA_DELETE_INTERVAL != value {
global.GDATA_DELETE_INTERVAL = value
}
}
} else {
wafSystemConfigAddReq := request.WafSystemConfigAddReq{
Item: "delete_history_log_day",
Value: strconv.Itoa(global.GDATA_DELETE_INTERVAL),
Remarks: "删除多少天前的日志数据(单位:天)",
ItemType: "int",
}
wafSystemConfigService.AddApi(wafSystemConfigAddReq)
}
configItem = wafSystemConfigService.GetDetailByItem("log_db_size")
if configItem.Id != "" {
value, err := strconv.ParseInt(configItem.Value, 10, 64)
if err == nil {
if global.GDATA_SHARE_DB_SIZE != value {
global.GDATA_SHARE_DB_SIZE = value
}
}
} else {
wafSystemConfigAddReq := request.WafSystemConfigAddReq{
Item: "log_db_size",
Value: strconv.FormatInt(global.GDATA_SHARE_DB_SIZE, 10),
Remarks: "日志归档最大记录数量",
ItemType: "int",
}
wafSystemConfigService.AddApi(wafSystemConfigAddReq)
}
//dns查询
configItem = wafSystemConfigService.GetDetailByItem("dns_server")
if configItem.Id != "" {
if global.GWAF_RUNTIME_DNS_SERVER != configItem.Value {
global.GWAF_RUNTIME_DNS_SERVER = configItem.Value
}
} else {
wafSystemConfigService.AddApi(request.WafSystemConfigAddReq{
Item: "dns_server",
Value: global.GWAF_RUNTIME_DNS_SERVER,
Remarks: "DNS服务器",
ItemType: "options",
Options: "119.29.29.29|腾讯DNS,8.8.8.8|谷歌DNS",
})
}
//日志记录类型
configItem = wafSystemConfigService.GetDetailByItem("record_log_type")
if configItem.Id != "" {
if global.GWAF_RUNTIME_RECORD_LOG_TYPE != configItem.Value {
global.GWAF_RUNTIME_RECORD_LOG_TYPE = configItem.Value
}
} else {
wafSystemConfigService.AddApi(request.WafSystemConfigAddReq{
Item: "record_log_type",
Value: global.GWAF_RUNTIME_RECORD_LOG_TYPE,
Remarks: "日志记录类型",
ItemType: "options",
Options: "all|全部,abnormal|非正常",
})
}
//控制中心-开关
configItem = wafSystemConfigService.GetDetailByItem("gwaf_center_enable")
if configItem.Id != "" {
if global.GWAF_CENTER_ENABLE != configItem.Value {
global.GWAF_CENTER_ENABLE = configItem.Value
}
} else {
wafSystemConfigService.AddApi(request.WafSystemConfigAddReq{
Item: "gwaf_center_enable",
Value: global.GWAF_CENTER_ENABLE,
Remarks: "中心开关",
ItemType: "bool",
Options: "false|关闭,true|开启",
})
}
//控制中心-中心url
configItem = wafSystemConfigService.GetDetailByItem("gwaf_center_url")
if configItem.Id != "" {
if global.GWAF_CENTER_URL != configItem.Value {
global.GWAF_CENTER_URL = configItem.Value
}
} else {
wafSystemConfigService.AddApi(request.WafSystemConfigAddReq{
Item: "gwaf_center_url",
Value: global.GWAF_CENTER_URL,
Remarks: "中心URL",
ItemType: "string",
Options: "",
})
}
//获取用户IP头方式
configItem = wafSystemConfigService.GetDetailByItem("gwaf_proxy_header")
if configItem.Id != "" {
if global.GCONFIG_RECORD_PROXY_HEADER != configItem.Value {
global.GCONFIG_RECORD_PROXY_HEADER = configItem.Value
}
} else {
wafSystemConfigService.AddApi(request.WafSystemConfigAddReq{
Item: "gwaf_proxy_header",
Value: global.GCONFIG_RECORD_PROXY_HEADER,
Remarks: "获取访客IP头信息(按照顺序)",
ItemType: "string",
Options: "",
})
}
//是否每天凌晨3点自动加载ssl证书
configItem = wafSystemConfigService.GetDetailByItem("auto_load_ssl_file")
if configItem.Id != "" {
value, err := strconv.ParseInt(configItem.Value, 10, 0)
if err == nil {
if global.GCONFIG_RECORD_AUTO_LOAD_SSL != value {
global.GCONFIG_RECORD_AUTO_LOAD_SSL = value
}
}
} else {
wafSystemConfigAddReq := request.WafSystemConfigAddReq{
Item: "auto_load_ssl_file",
Value: strconv.FormatInt(global.GCONFIG_RECORD_AUTO_LOAD_SSL, 10),
Remarks: "是否每天凌晨3点自动加载ssl证书",
ItemType: "int",
ItemClass: itemClass,
Item: itemName,
Value: strconv.FormatInt(defaultValue, 10),
Remarks: remarks,
ItemType: itemType,
Options: options,
}
wafSystemConfigService.AddApi(wafSystemConfigAddReq)
}
}
func updateConfigStringItem(initLoad bool, itemClass string, itemName string, defaultValue string, remarks string, itemType string, options string) {
configItem := wafSystemConfigService.GetDetailByItem(itemName)
if configItem.Id != "" {
if defaultValue != configItem.Value {
setConfigStringValue(itemName, configItem.Value, 1)
} else if initLoad == true {
setConfigStringValue(itemName, configItem.Value, 0)
}
} else {
wafSystemConfigAddReq := request.WafSystemConfigAddReq{
ItemClass: itemClass,
Item: itemName,
Value: defaultValue,
Remarks: remarks,
ItemType: itemType,
Options: options,
}
wafSystemConfigService.AddApi(wafSystemConfigAddReq)
}
}
// 加载配置数据
func TaskLoadSetting(initLoad bool) {
zlog.Debug("TaskLoadSetting")
updateConfigIntItem(initLoad, "system", "record_max_req_body_length", global.GCONFIG_RECORD_MAX_BODY_LENGTH, "记录请求最大报文", "int", "")
updateConfigIntItem(initLoad, "system", "record_max_res_body_length", global.GCONFIG_RECORD_MAX_RES_BODY_LENGTH, "如果可以记录,满足最大响应报文大小才记录", "int", "")
updateConfigIntItem(initLoad, "system", "record_resp", global.GCONFIG_RECORD_RESP, "是否记录响应报文", "int", "")
updateConfigIntItem(initLoad, "system", "delete_history_log_day", global.GDATA_DELETE_INTERVAL, "删除多少天前的日志数据(单位:天)", "int", "")
updateConfigIntItem(initLoad, "system", "log_db_size", global.GDATA_SHARE_DB_SIZE, "日志归档最大记录数量", "int", "")
updateConfigIntItem(initLoad, "system", "auto_load_ssl_file", global.GCONFIG_RECORD_AUTO_LOAD_SSL, "是否每天凌晨3点自动加载ssl证书", "int", "")
updateConfigStringItem(initLoad, "system", "dns_server", global.GWAF_RUNTIME_DNS_SERVER, "DNS服务器", "options", "119.29.29.29|腾讯DNS,8.8.8.8|谷歌DNS")
updateConfigStringItem(initLoad, "system", "record_log_type", global.GWAF_RUNTIME_RECORD_LOG_TYPE, "日志记录类型", "options", "all|全部,abnormal|非正常")
updateConfigStringItem(initLoad, "system", "gwaf_center_enable", global.GWAF_CENTER_ENABLE, "中心开关", "bool", "false|关闭,true|开启")
updateConfigStringItem(initLoad, "system", "gwaf_center_url", global.GWAF_CENTER_URL, "中心URL", "string", "")
updateConfigStringItem(initLoad, "system", "gwaf_proxy_header", global.GCONFIG_RECORD_PROXY_HEADER, "获取访客IP头信息(按照顺序)", "string", "")
updateConfigIntItem(initLoad, "kafka", "kafka_enable", global.GCONFIG_RECORD_KAFKA_ENABLE, "kafka 是否激活", "int", "")
updateConfigStringItem(initLoad, "kafka", "kafka_url", global.GCONFIG_RECORD_KAFKA_URL, "kafka url地址", "string", "")
updateConfigStringItem(initLoad, "kafka", "kafka_topic", global.GCONFIG_RECORD_KAFKA_TOPIC, "kafka topic", "string", "")
}
+1 -1
View File
@@ -1,13 +1,13 @@
package waftask
import (
"SamWaf/common/zlog"
"SamWaf/customtype"
"SamWaf/global"
"SamWaf/innerbean"
"SamWaf/model"
"SamWaf/model/baseorm"
"SamWaf/utils"
"SamWaf/utils/zlog"
"SamWaf/wafdb"
"fmt"
uuid "github.com/satori/go.uuid"
+1 -1
View File
@@ -1,12 +1,12 @@
package waftask
import (
"SamWaf/common/zlog"
"SamWaf/enums"
"SamWaf/global"
"SamWaf/model/spec"
"SamWaf/service/waf_service"
"SamWaf/utils"
"SamWaf/utils/zlog"
)
var (
+1 -1
View File
@@ -1,11 +1,11 @@
package waftask
import (
"SamWaf/common/zlog"
"SamWaf/customtype"
"SamWaf/global"
"SamWaf/model/request"
"SamWaf/service/waf_service"
"SamWaf/utils/zlog"
"SamWaf/wafsec"
"encoding/json"
"io/ioutil"