/* Copyright 2026 Arctel.net Licensed under the Apache License, Version 2.0 (the "License"); you may not use this file except in compliance with the License. You may obtain a copy of the License at http://www.apache.org/licenses/LICENSE-2.0 Unless required by applicable law or agreed to in writing, software distributed under the License is distributed on an "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the License for the specific language governing permissions and limitations under the License. */ package risk_control import ( "context" "time" "github.com/Rain-kl/Wavelet/internal/config" "github.com/Rain-kl/Wavelet/internal/db" "github.com/Rain-kl/Wavelet/internal/logger" ) // UserAccessLog 用户访问记录 type UserAccessLog struct { ID uint64 `json:"id,string"` UserID uint64 `json:"user_id,string"` Path string `json:"path"` Method string `json:"method"` IP string `json:"ip"` UserAgent string `json:"user_agent"` Headers string `json:"headers"` Status int32 `json:"status"` Latency int64 `json:"latency"` // 耗时毫秒 CreatedAt time.Time `json:"created_at"` } var ( logChan chan *UserAccessLog ) const ( defaultQueueSize = 10000 maxBatchSize = 1000 flushInterval = 1 * time.Second ) // InitLogWriter 初始化日志写入通道和后台写入协程 func InitLogWriter() { if !config.Config.ClickHouse.Enabled { return } logChan = make(chan *UserAccessLog, defaultQueueSize) go startBatchWorker() } // IsBufferFull 检查当前本地缓冲队列是否已满 // 如果没有启用 ClickHouse,默认返回 false,不触发限流 func IsBufferFull() bool { if !config.Config.ClickHouse.Enabled || logChan == nil { return false } return len(logChan) >= cap(logChan) } // QueueAccessLog 异步非阻塞地将日志推入缓冲队列 func QueueAccessLog(logItem *UserAccessLog) { if !config.Config.ClickHouse.Enabled || logChan == nil { return } select { case logChan <- logItem: default: // 如果在极端并发下仍然写满了,这里做非阻塞丢弃,防止卡死 logger.WarnF(context.Background(), "[RiskControl] Log queue full, dropping log item for path: %s", logItem.Path) } } // startBatchWorker 后台批量写入 ClickHouse 的工作协程 func startBatchWorker() { ticker := time.NewTicker(flushInterval) defer ticker.Stop() var batch []*UserAccessLog flush := func() { if len(batch) == 0 { return } if db.ChConn == nil { batch = nil return } ctx := context.Background() b, err := db.ChConn.PrepareBatch(ctx, "INSERT INTO user_access_logs (id, user_id, path, method, ip, user_agent, headers, status, latency, created_at)") if err != nil { logger.ErrorF(ctx, "[RiskControl] Prepare ClickHouse batch failed: %v", err) batch = nil return } for _, item := range batch { err = b.Append( item.ID, item.UserID, item.Path, item.Method, item.IP, item.UserAgent, item.Headers, item.Status, item.Latency, item.CreatedAt, ) if err != nil { logger.ErrorF(ctx, "[RiskControl] Append item to ClickHouse batch failed: %v", err) } } if err := b.Send(); err != nil { logger.ErrorF(ctx, "[RiskControl] Send ClickHouse batch failed: %v", err) } batch = nil } for { select { case item, ok := <-logChan: if !ok { flush() return } batch = append(batch, item) if len(batch) >= maxBatchSize { flush() } case <-ticker.C: flush() } } }