- 新增访问统计:visitstats 服务与查询、visit_capture 中间件、爬虫/来源识别、管理端 analytics 页面 - 新增侧边栏:sidebar 服务与 handler、管理端 sidebar 配置页、前端侧边栏组件 - 新增 widget 嵌入代码生成(widgetCode)与 widgetApi - 遥测逻辑迁移至 handler 层,删除 service/telemetry - env 示例补充 Docker 部署可信代理 IP 说明
296 lines
8.3 KiB
Go
296 lines
8.3 KiB
Go
package service
|
||
|
||
import (
|
||
"context"
|
||
"crypto/sha256"
|
||
"encoding/hex"
|
||
"log"
|
||
"regexp"
|
||
"strings"
|
||
"sync"
|
||
"sync/atomic"
|
||
"time"
|
||
|
||
"github.com/freefire/jiang13-bbs/model"
|
||
"gorm.io/gorm"
|
||
"gorm.io/gorm/clause"
|
||
)
|
||
|
||
// 访问统计管道参数
|
||
const (
|
||
visitChannelSize = 8192 // 事件缓冲:满则丢弃(统计场景可接受)
|
||
visitFlushBatch = 500 // 攒批上限:满批立即落库
|
||
visitFlushInterval = 3 * time.Second // 或每 3 秒落库一次
|
||
visitDeleteBatch = 5000 // 保留期清理分批大小,避免长事务
|
||
)
|
||
|
||
// VisitorCookieName 浏览器访客 ID cookie 名(handler 签发、捕获中间件读取)
|
||
const VisitorCookieName = "j13_vid"
|
||
|
||
var vidUUIDRe = regexp.MustCompile(`(?i)^[0-9a-f]{8}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{4}-[0-9a-f]{12}$`)
|
||
|
||
// todayDate 当天零点(本地时区)
|
||
func todayDate() time.Time {
|
||
now := time.Now()
|
||
return time.Date(now.Year(), now.Month(), now.Day(), 0, 0, 0, 0, now.Location())
|
||
}
|
||
|
||
// HashVisitorID 访客 ID 仅存 SHA-256 哈希(明细不含原始 vid)
|
||
func HashVisitorID(vid string) string {
|
||
sum := sha256.Sum256([]byte(strings.TrimSpace(strings.ToLower(vid))))
|
||
return hex.EncodeToString(sum[:])
|
||
}
|
||
|
||
// ValidVisitorID 校验浏览器访客 ID(UUID 形态)
|
||
func ValidVisitorID(vid string) bool {
|
||
return vidUUIDRe.MatchString(strings.TrimSpace(vid))
|
||
}
|
||
|
||
// VisitStatsService 访问统计:内存攒批落库 + 保留期滚动清理 + 管理端聚合查询
|
||
type VisitStatsService struct {
|
||
db *gorm.DB
|
||
setting *SettingService
|
||
ch chan model.VisitEvent
|
||
dropped atomic.Uint64
|
||
flushOnce sync.Once
|
||
}
|
||
|
||
func NewVisitStatsService(db *gorm.DB, settingSvc *SettingService) *VisitStatsService {
|
||
return &VisitStatsService{
|
||
db: db,
|
||
setting: settingSvc,
|
||
ch: make(chan model.VisitEvent, visitChannelSize),
|
||
}
|
||
}
|
||
|
||
// Enqueue 非阻塞投递事件;缓冲满时丢弃并计数
|
||
func (s *VisitStatsService) Enqueue(ev model.VisitEvent) {
|
||
if ev.CreatedAt.IsZero() {
|
||
ev.CreatedAt = time.Now()
|
||
}
|
||
select {
|
||
case s.ch <- ev:
|
||
default:
|
||
s.dropped.Add(1)
|
||
}
|
||
}
|
||
|
||
// DroppedCount 已因缓冲满而丢弃的事件数(诊断用)
|
||
func (s *VisitStatsService) DroppedCount() uint64 { return s.dropped.Load() }
|
||
|
||
// pvDelta 一天内某批 pageview 事件的 PV/UV 增量
|
||
type pvDelta struct {
|
||
PV int64
|
||
LoggedInPV int64
|
||
VidHashes map[string]struct{}
|
||
}
|
||
|
||
// summarizePageviews 批内按 (日期, vid) 归并 PV/UV 增量;bot/probe 不计。纯函数,便于单测。
|
||
func summarizePageviews(events []model.VisitEvent) map[time.Time]*pvDelta {
|
||
out := make(map[time.Time]*pvDelta, 2)
|
||
for _, ev := range events {
|
||
if ev.Kind != model.VisitKindPageview {
|
||
continue
|
||
}
|
||
// 取事件本地日期(不能用 Truncate:它按 UTC 绝对时间截断,非零时区会错位)
|
||
t := ev.CreatedAt.In(time.Local)
|
||
day := time.Date(t.Year(), t.Month(), t.Day(), 0, 0, 0, 0, time.Local)
|
||
d, ok := out[day]
|
||
if !ok {
|
||
d = &pvDelta{VidHashes: make(map[string]struct{})}
|
||
out[day] = d
|
||
}
|
||
d.PV++
|
||
if ev.UserID > 0 {
|
||
d.LoggedInPV++
|
||
}
|
||
if ev.VidHash != "" {
|
||
d.VidHashes[ev.VidHash] = struct{}{}
|
||
}
|
||
}
|
||
return out
|
||
}
|
||
|
||
// flush 落库一批事件:明细批量 INSERT + 同一事务合并日统计(PV 必增、当日 vid 首次 UV+1)。
|
||
// bot/probe 只写明细。失败记日志不重试(统计允许少量丢失)。
|
||
func (s *VisitStatsService) flush(events []model.VisitEvent) {
|
||
if len(events) == 0 {
|
||
return
|
||
}
|
||
if err := s.db.CreateInBatches(&events, visitFlushBatch).Error; err != nil {
|
||
log.Printf("[visitstats] 明细落库失败(%d 条丢弃): %v", len(events), err)
|
||
return
|
||
}
|
||
|
||
deltas := summarizePageviews(events)
|
||
if len(deltas) == 0 {
|
||
return
|
||
}
|
||
err := s.db.Transaction(func(tx *gorm.DB) error {
|
||
for day, d := range deltas {
|
||
// 当日访客去重:插入成功(此前未出现)才计 UV
|
||
uv := int64(0)
|
||
for vid := range d.VidHashes {
|
||
vis := model.SiteDailyVisitor{Date: day, VidHash: vid}
|
||
res := tx.Clauses(clause.OnConflict{DoNothing: true}).Create(&vis)
|
||
if res.Error != nil {
|
||
return res.Error
|
||
}
|
||
uv += res.RowsAffected
|
||
}
|
||
// 日统计 upsert:行不存在自动创建,存在则累加
|
||
if err := tx.Exec(`
|
||
INSERT INTO site_daily_stats (date, pv, uv, logged_in_pv, updated_at)
|
||
VALUES (?, ?, ?, ?, ?)
|
||
ON CONFLICT (date) DO UPDATE SET
|
||
pv = site_daily_stats.pv + EXCLUDED.pv,
|
||
uv = site_daily_stats.uv + EXCLUDED.uv,
|
||
logged_in_pv = site_daily_stats.logged_in_pv + EXCLUDED.logged_in_pv,
|
||
updated_at = EXCLUDED.updated_at
|
||
`, day, d.PV, uv, d.LoggedInPV, time.Now()).Error; err != nil {
|
||
return err
|
||
}
|
||
}
|
||
return nil
|
||
})
|
||
if err != nil {
|
||
log.Printf("[visitstats] 日统计合并失败: %v", err)
|
||
}
|
||
}
|
||
|
||
// StartFlusher 攒批落库协程:满 500 条或每 3 秒触发;ctx 取消时冲刷余量后退出
|
||
func (s *VisitStatsService) StartFlusher(ctx context.Context) {
|
||
s.flushOnce.Do(func() {
|
||
ticker := time.NewTicker(visitFlushInterval)
|
||
defer ticker.Stop()
|
||
buf := make([]model.VisitEvent, 0, visitFlushBatch)
|
||
for {
|
||
select {
|
||
case <-ctx.Done():
|
||
s.drain(buf)
|
||
return
|
||
case ev := <-s.ch:
|
||
buf = append(buf, ev)
|
||
if len(buf) >= visitFlushBatch {
|
||
s.flush(buf)
|
||
buf = make([]model.VisitEvent, 0, visitFlushBatch)
|
||
}
|
||
case <-ticker.C:
|
||
if len(buf) > 0 {
|
||
s.flush(buf)
|
||
buf = make([]model.VisitEvent, 0, visitFlushBatch)
|
||
}
|
||
}
|
||
}
|
||
})
|
||
}
|
||
|
||
// drain 退出前尽量冲刷缓冲中的余量(非阻塞)
|
||
func (s *VisitStatsService) drain(buf []model.VisitEvent) {
|
||
for {
|
||
select {
|
||
case ev := <-s.ch:
|
||
buf = append(buf, ev)
|
||
if len(buf) >= visitFlushBatch {
|
||
s.flush(buf)
|
||
buf = make([]model.VisitEvent, 0, visitFlushBatch)
|
||
}
|
||
default:
|
||
s.flush(buf)
|
||
return
|
||
}
|
||
}
|
||
}
|
||
|
||
// retentionDays 当前保留天数(缺行/非法回退默认 90)
|
||
func (s *VisitStatsService) retentionDays() int {
|
||
days, err := s.setting.AnalyticsRetentionDays()
|
||
if err != nil || days <= 0 {
|
||
return DefaultAnalyticsRetentionDays
|
||
}
|
||
return days
|
||
}
|
||
|
||
// deleteOlderThan 分批删除截止时间前的明细,返回删除总数
|
||
func (s *VisitStatsService) deleteOlderThan(cutoff time.Time) int64 {
|
||
var total int64
|
||
for {
|
||
res := s.db.Exec(
|
||
`DELETE FROM visit_events WHERE id IN (
|
||
SELECT id FROM visit_events WHERE created_at < ? LIMIT ?
|
||
)`, cutoff, visitDeleteBatch)
|
||
if res.Error != nil {
|
||
log.Printf("[visitstats] 保留期清理失败: %v", res.Error)
|
||
return total
|
||
}
|
||
total += res.RowsAffected
|
||
if res.RowsAffected < visitDeleteBatch {
|
||
return total
|
||
}
|
||
}
|
||
}
|
||
|
||
// RunRetention 保留期滚动清理:每天 04:30 低峰触发
|
||
func (s *VisitStatsService) RunRetention(ctx context.Context) {
|
||
for {
|
||
now := time.Now()
|
||
next := time.Date(now.Year(), now.Month(), now.Day(), 4, 30, 0, 0, now.Location())
|
||
if !next.After(now) {
|
||
next = next.AddDate(0, 0, 1)
|
||
}
|
||
select {
|
||
case <-ctx.Done():
|
||
return
|
||
case <-time.After(next.Sub(now)):
|
||
}
|
||
cutoff := time.Now().AddDate(0, 0, -s.retentionDays())
|
||
if n := s.deleteOlderThan(cutoff); n > 0 {
|
||
log.Printf("[visitstats] 保留期清理 %d 条(> %s)", n, cutoff.Format("2006-01-02"))
|
||
}
|
||
}
|
||
}
|
||
|
||
// PurgeBefore 手动清理 N 天前明细:goroutine 异步分批执行,受理即返回(前端稍后刷新容量)。
|
||
func (s *VisitStatsService) PurgeBefore(days int) bool {
|
||
if days < MinAnalyticsRetentionDays {
|
||
return false
|
||
}
|
||
if days > MaxAnalyticsRetentionDays {
|
||
days = MaxAnalyticsRetentionDays
|
||
}
|
||
go func() {
|
||
cutoff := time.Now().AddDate(0, 0, -days)
|
||
if n := s.deleteOlderThan(cutoff); n > 0 {
|
||
log.Printf("[visitstats] 手动清理 %d 条(> %s)", n, cutoff.Format("2006-01-02"))
|
||
}
|
||
}()
|
||
return true
|
||
}
|
||
|
||
// VisitUsage 明细表容量(行数 + 磁盘占用字节)
|
||
type VisitUsage struct {
|
||
Rows int64 `json:"rows"`
|
||
Bytes int64 `json:"bytes"`
|
||
}
|
||
|
||
// Usage 统计明细行数与磁盘占用(含索引/TOAST)
|
||
func (s *VisitStatsService) Usage() (VisitUsage, error) {
|
||
var u VisitUsage
|
||
if err := s.db.Model(&model.VisitEvent{}).Count(&u.Rows).Error; err != nil {
|
||
return u, err
|
||
}
|
||
if err := s.db.Raw(`SELECT pg_total_relation_size('visit_events')`).Scan(&u.Bytes).Error; err != nil {
|
||
return u, err
|
||
}
|
||
return u, nil
|
||
}
|
||
|
||
// Enabled 采集开关(缺行视为开启)
|
||
func (s *VisitStatsService) Enabled() bool {
|
||
on, err := s.setting.AnalyticsEnabled()
|
||
return err == nil && on
|
||
}
|
||
|
||
// CanonicalIP 归一化客户端 IP(导出供 handler/middleware 使用)
|
||
func CanonicalIP(ip string) string { return canonicalIP(ip) }
|