From 558edc864b02e4659e5d050cb06351c37ef3bdab Mon Sep 17 00:00:00 2001 From: Your Name Date: Tue, 6 Jan 2026 16:20:02 +0800 Subject: [PATCH] feat: add send_stats table --- migrate/migrate.go | 1 + models/send_stats.go | 199 ++++++++++++++++++ models/send_tasks_logs.go | 77 ++++--- routers/api/v1/statistic.go | 33 +++ routers/router.go | 1 + service/send_message_service/send_message.go | 46 ++++ service/statistic_service/statistic.go | 27 +++ .../components/pages/dashboard/Dashboard.vue | 30 +-- 8 files changed, 372 insertions(+), 42 deletions(-) create mode 100644 models/send_stats.go diff --git a/migrate/migrate.go b/migrate/migrate.go index 5215ed6..0b23adb 100644 --- a/migrate/migrate.go +++ b/migrate/migrate.go @@ -69,6 +69,7 @@ func Setup() { &models.HostedMessage{}, &models.LoginLog{}, &models.Template{}, + &models.SendStats{}, } for _, table := range tables { diff --git a/models/send_stats.go b/models/send_stats.go new file mode 100644 index 0000000..651b338 --- /dev/null +++ b/models/send_stats.go @@ -0,0 +1,199 @@ +package models + +import ( + "fmt" + "message-nest/pkg/util" +) + +// SendStats 发送统计表 +type SendStats struct { + ID uint `gorm:"primaryKey;autoIncrement" json:"id"` + TaskID *uint `json:"task_id" gorm:"type:bigint unsigned;index:idx_task_type_day_status"` + TaskType string `json:"task_type" gorm:"type:varchar(20);default:'task';index:idx_task_type_day_status;comment:'任务类型:task-发信任务,template-模板任务'"` + Day string `json:"day" gorm:"type:varchar(10);index:idx_task_type_day_status;index:idx_day_status"` + Status string `json:"status" gorm:"type:varchar(20);index:idx_task_type_day_status;index:idx_day_status"` + Num int64 `json:"num" gorm:"type:bigint;default:0"` +} + +// SendStatsData 首页统计数据结构 +type SendStatsData struct { + TodaySuccNum int64 `json:"today_succ_num"` + TodayFailedNum int64 `json:"today_failed_num"` + TodayTotalNum int64 `json:"today_total_num"` + TotalSuccNum int64 `json:"total_succ_num"` + TotalFailedNum int64 `json:"total_failed_num"` + TotalNum int64 `json:"total_num"` + DailyStats []DailyStatsData `json:"daily_stats"` + StatusStats []StatusStatsData `json:"status_stats"` + TaskTypeStats []TaskTypeStatsData `json:"task_type_stats"` +} + +// DailyStatsData 每日统计数据 +type DailyStatsData struct { + Day string `json:"day"` + SuccNum int64 `json:"succ_num"` + FailedNum int64 `json:"failed_num"` + TotalNum int64 `json:"total_num"` +} + +// StatusStatsData 状态统计数据 +type StatusStatsData struct { + Status string `json:"status"` + Num int64 `json:"num"` +} + +// TaskTypeStatsData 任务类型统计数据 +type TaskTypeStatsData struct { + TaskType string `json:"task_type"` + Num int64 `json:"num"` +} + +// GetSendStatsData 获取发送统计数据 +func GetSendStatsData(days int) (SendStatsData, error) { + var result SendStatsData + statsTable := GetSchema(SendStats{}) + currDay := util.GetNowTimeStr()[:10] + + // 今日统计 + todayQuery := db.Table(statsTable). + Select(` + SUM(CASE WHEN status = 'success' THEN num ELSE 0 END) AS today_succ_num, + SUM(CASE WHEN status = 'failed' THEN num ELSE 0 END) AS today_failed_num, + SUM(num) AS today_total_num + `). + Where("day = ?", currDay) + + todayQuery.Scan(&result) + + // 总计统计 + totalQuery := db.Table(statsTable). + Select(` + SUM(CASE WHEN status = 'success' THEN num ELSE 0 END) AS total_succ_num, + SUM(CASE WHEN status = 'failed' THEN num ELSE 0 END) AS total_failed_num, + SUM(num) AS total_num + `) + + totalQuery.Scan(&result) + + // 最近N天每日统计 + if days > 0 { + now := util.GetNowTime() + past := now.AddDate(0, 0, -days) + pastDate := past.Format("2006-01-02") + + var dailyStats []DailyStatsData + dailyQuery := db.Table(statsTable). + Select(` + day, + SUM(CASE WHEN status = 'success' THEN num ELSE 0 END) AS succ_num, + SUM(CASE WHEN status = 'failed' THEN num ELSE 0 END) AS failed_num, + SUM(num) AS total_num + `). + Where("day >= ?", pastDate). + Group("day"). + Order("day DESC") + + dailyQuery.Scan(&dailyStats) + result.DailyStats = dailyStats + } + + // 状态分布统计 + var statusStats []StatusStatsData + statusQuery := db.Table(statsTable). + Select("status, SUM(num) AS num"). + Group("status"). + Order("num DESC") + + statusQuery.Scan(&statusStats) + result.StatusStats = statusStats + + // 任务类型分布统计 + var taskTypeStats []TaskTypeStatsData + taskTypeQuery := db.Table(statsTable). + Select("task_type, SUM(num) AS num"). + Group("task_type"). + Order("num DESC") + + taskTypeQuery.Scan(&taskTypeStats) + result.TaskTypeStats = taskTypeStats + + return result, nil +} + +// GetSendStatsByTask 获取指定任务的统计数据 +func GetSendStatsByTask(taskID uint, days int) (SendStatsData, error) { + var result SendStatsData + statsTable := GetSchema(SendStats{}) + currDay := util.GetNowTimeStr()[:10] + + // 今日统计 + todayQuery := db.Table(statsTable). + Select(` + SUM(CASE WHEN status = 'success' THEN num ELSE 0 END) AS today_succ_num, + SUM(CASE WHEN status = 'failed' THEN num ELSE 0 END) AS today_failed_num, + SUM(num) AS today_total_num + `). + Where("day = ? AND task_id = ?", currDay, taskID) + + todayQuery.Scan(&result) + + // 总计统计 + totalQuery := db.Table(statsTable). + Select(` + SUM(CASE WHEN status = 'success' THEN num ELSE 0 END) AS total_succ_num, + SUM(CASE WHEN status = 'failed' THEN num ELSE 0 END) AS total_failed_num, + SUM(num) AS total_num + `). + Where("task_id = ?", taskID) + + totalQuery.Scan(&result) + + // 最近N天每日统计 + if days > 0 { + now := util.GetNowTime() + past := now.AddDate(0, 0, -days) + pastDate := past.Format("2006-01-02") + + var dailyStats []DailyStatsData + dailyQuery := db.Table(statsTable). + Select(` + day, + SUM(CASE WHEN status = 'success' THEN num ELSE 0 END) AS succ_num, + SUM(CASE WHEN status = 'failed' THEN num ELSE 0 END) AS failed_num, + SUM(num) AS total_num + `). + Where("day >= ? AND task_id = ?", pastDate, taskID). + Group("day"). + Order("day DESC") + + dailyQuery.Scan(&dailyStats) + result.DailyStats = dailyStats + } + + // 状态分布统计 + var statusStats []StatusStatsData + statusQuery := db.Table(statsTable). + Select("status, SUM(num) AS num"). + Where("task_id = ?", taskID). + Group("status"). + Order("num DESC") + + statusQuery.Scan(&statusStats) + result.StatusStats = statusStats + + return result, nil +} + +// IncrementSendStats 增加发送统计(用于实时更新统计数据) +func IncrementSendStats(taskID *uint, taskType string, day string, status string, num int64) error { + statsTable := GetSchema(SendStats{}) + + // 使用 ON DUPLICATE KEY UPDATE 或 UPSERT 逻辑 + query := fmt.Sprintf(` + INSERT INTO %s (task_id, task_type, day, status, num) + VALUES (?, ?, ?, ?, ?) + ON DUPLICATE KEY UPDATE num = num + ? + `, statsTable) + + return db.Exec(query, taskID, taskType, day, status, num, num).Error +} diff --git a/models/send_tasks_logs.go b/models/send_tasks_logs.go index f60b72a..108182b 100644 --- a/models/send_tasks_logs.go +++ b/models/send_tasks_logs.go @@ -55,7 +55,7 @@ func GetSendLogs(pageNum int, pageSize int, name string, taskId string, maps map } query = query.Where(maps) - + // 按名称搜索(搜索日志表的 name 字段) if name != "" { query = query.Where(fmt.Sprintf("%s.name like ?", logt), fmt.Sprintf("%%%s%%", name)) @@ -68,7 +68,7 @@ func GetSendLogs(pageNum int, pageSize int, name string, taskId string, maps map query = query.Offset(pageNum).Limit(pageSize) } query.Scan(&logs) - + //v1 接口的历史日志数据兼容处理 // 应用层处理:为历史数据(type=task 且 name 为空)补充任务名称 fillTaskNamesForLogs(&logs) @@ -130,7 +130,7 @@ func fillTaskNamesForLogs(logs *[]LogsResult) { func GetSendLogsTotal(name string, taskId string, maps map[string]interface{}) (int64, error) { var total int64 logt := GetSchema(SendTasksLogs{}) - + // 简化查询,只查询日志表 query := db.Table(logt) @@ -141,7 +141,7 @@ func GetSendLogsTotal(name string, taskId string, maps map[string]interface{}) ( } query = query.Where(maps) - + // 按名称搜索(搜索日志表的 name 字段) if name != "" { query = query.Where(fmt.Sprintf("%s.name like ?", logt), fmt.Sprintf("%%%s%%", name)) @@ -156,7 +156,7 @@ func GetSendLogsTotal(name string, taskId string, maps map[string]interface{}) ( // GetSendLogsTotal 获取所有日志总数 func DeleteOutDateLogs(keepNum int) (int, error) { var affectedRows int - + // 优化方案:使用GORM的Offset和Limit找到临界ID,兼容多种数据库 // 1. 获取第 keepNum 条记录的ID作为临界值 var threshold SendTasksLogs @@ -166,18 +166,18 @@ func DeleteOutDateLogs(keepNum int) (int, error) { Offset(keepNum - 1). Limit(1). First(&threshold) - + // 如果记录总数不足keepNum条,则不需要删除 if result.Error != nil { return 0, nil } - + // 2. 删除ID小于临界值的记录 deleteResult := db.Where("id < ?", threshold.ID).Delete(&SendTasksLogs{}) if deleteResult.Error != nil { return affectedRows, deleteResult.Error } - + affectedRows = int(deleteResult.RowsAffected) return affectedRows, nil } @@ -317,27 +317,27 @@ func GetBasicStatisticData() (BasicStatisticData, error) { return statistic, nil } -// GetTrendStatisticData 获取趋势统计数据 +// GetTrendStatisticData 获取趋势统计数据(使用 send_stats 表) func GetTrendStatisticData() (TrendStatisticData, error) { var statistic TrendStatisticData var latestData []LatestSendData - logt := GetSchema(SendTasksLogs{}) + statsTable := GetSchema(SendStats{}) // 最近30天数据 days := 30 now := util.GetNowTime() past := now.AddDate(0, 0, -days) pastDate := past.Format("2006-01-02") - next := now.AddDate(0, 0, 1) - nextDate := next.Format("2006-01-02") + queryData := db. - Table(logt). + Table(statsTable). Select(` - CAST(DATE(created_on) AS CHAR) AS day, - SUM(CASE WHEN status = 1 THEN 1 ELSE 0 END) AS day_succ_num, - SUM(CASE WHEN status != 1 or status is null THEN 1 ELSE 0 END) AS day_failed_num, - COUNT(*) AS num`). - Where(fmt.Sprintf(" created_on >= '%s' and created_on <= '%s' ", pastDate, nextDate)). + day, + SUM(CASE WHEN status = 'success' THEN num ELSE 0 END) AS day_succ_num, + SUM(CASE WHEN status = 'failed' THEN num ELSE 0 END) AS day_failed_num, + SUM(num) AS num + `). + Where("day >= ?", pastDate). Group("day"). Order("day") @@ -346,21 +346,42 @@ func GetTrendStatisticData() (TrendStatisticData, error) { return statistic, nil } -// GetChannelStatisticData 获取渠道统计数据(包含任务实例和模板实例) +// GetChannelStatisticData 获取渠道统计数据(使用 send_stats 表的任务类型统计) func GetChannelStatisticData() (ChannelStatisticData, error) { var statistic ChannelStatisticData var wayCateData []WayCateData - inst := GetSchema(SendTasksIns{}) - wayst := GetSchema(SendWays{}) + statsTable := GetSchema(SendStats{}) + + // 任务类型映射 + taskTypeMap := map[string]string{ + "task": "发信任务", + "template": "模板任务", + } + + // 统计任务类型分布 + var taskTypeStats []struct { + TaskType string + Num int64 + } - // 统计所有实例的渠道分布(包含任务实例和模板实例) - // 不区分 task_id 和 template_id,统计所有关联到渠道的实例 db. - Table(inst). - Select(fmt.Sprintf("%s.name as way_name, count(%s.id) as count_num", wayst, inst)). - Joins(fmt.Sprintf("JOIN %s ON %s.way_id = %s.id", wayst, inst, wayst)). - Group(fmt.Sprintf("%s.id", wayst)). - Scan(&wayCateData) + Table(statsTable). + Select("task_type, SUM(num) AS num"). + Group("task_type"). + Order("num DESC"). + Scan(&taskTypeStats) + + // 转换为 WayCateData 格式(复用前端已有的数据结构) + for _, stat := range taskTypeStats { + wayName := taskTypeMap[stat.TaskType] + if wayName == "" { + wayName = stat.TaskType + } + wayCateData = append(wayCateData, WayCateData{ + WayName: wayName, + CountNum: int(stat.Num), + }) + } statistic.WayCateData = wayCateData return statistic, nil diff --git a/routers/api/v1/statistic.go b/routers/api/v1/statistic.go index 0d56585..4fd1461 100644 --- a/routers/api/v1/statistic.go +++ b/routers/api/v1/statistic.go @@ -41,6 +41,14 @@ func GetStatisticData(c *gin.Context) { return } appG.CResponse(http.StatusOK, "获取渠道统计成功", data) + case "send_stats": + // 新增:基于 send_stats 表的统计数据 + data, err := msgService.GetSendStatsData() + if err != nil { + appG.CResponse(http.StatusInternalServerError, fmt.Sprintf("获取发送统计失败!原因:%s", err), nil) + return + } + appG.CResponse(http.StatusOK, "获取发送统计成功", data) default: // 默认返回完整统计数据(保持向后兼容) data, err := msgService.GetStatisticData() @@ -52,3 +60,28 @@ func GetStatisticData(c *gin.Context) { } } +// GetSendStatsByTask 获取指定任务的发送统计数据 +func GetSendStatsByTask(c *gin.Context) { + var ( + appG = app.Gin{C: c} + ) + + taskID := c.Query("task_id") + if taskID == "" { + appG.CResponse(http.StatusBadRequest, "任务ID不能为空", nil) + return + } + + msgService := statistic_service.StatisticService{ + TaskID: taskID, + } + + data, err := msgService.GetSendStatsByTask() + if err != nil { + appG.CResponse(http.StatusInternalServerError, fmt.Sprintf("获取任务统计失败!原因:%s", err), nil) + return + } + + appG.CResponse(http.StatusOK, "获取任务统计成功", data) +} + diff --git a/routers/router.go b/routers/router.go index 08fe9fe..c647aac 100644 --- a/routers/router.go +++ b/routers/router.go @@ -88,6 +88,7 @@ func InitRouter(f embed.FS) *gin.Engine { // statistic apiV1.GET("/statistic", v1.GetStatisticData) + apiV1.GET("/statistic/task", v1.GetSendStatsByTask) // cronMessage apiV1.POST("/cronmessages/addone", v1.AddCronMsgTask) diff --git a/service/send_message_service/send_message.go b/service/send_message_service/send_message.go index 1b9dc11..53d2bcb 100644 --- a/service/send_message_service/send_message.go +++ b/service/send_message_service/send_message.go @@ -6,6 +6,7 @@ import ( "fmt" "message-nest/models" "message-nest/pkg/constant" + "message-nest/pkg/util" "message-nest/service/send_message_service/unified" "message-nest/service/send_task_service" "message-nest/service/send_way_service" @@ -294,6 +295,8 @@ func (sm *SendMessageService) Send(task models.TaskIns) (string, error) { sm.AppendSendContent() // 日志写到数据库 sm.RecordSendLog() + // 更新统计数据(任务级别:一次任务算一次) + sm.UpdateSendStats() totalOutputLog := strings.Join(sm.LogOutput, "\n") if sm.Status == SendSuccess { @@ -339,6 +342,49 @@ func (sm *SendMessageService) RecordSendLog() { } } +// UpdateSendStats 更新发送统计数据(任务级别) +func (sm *SendMessageService) UpdateSendStats() { + // 获取当前日期 + currentDay := sm.getCurrentDay() + + // 解析任务ID(如果是数字ID) + var taskID *uint + if sm.TaskID != "" { + // 尝试将字符串ID转换为uint(如果是数字ID) + var id uint64 + _, err := fmt.Sscanf(sm.TaskID, "%d", &id) + if err == nil { + taskIDValue := uint(id) + taskID = &taskIDValue + } + } + + // 确定任务类型 + taskType := "task" + if sm.SendMode == SendModeTemplate { + taskType = "template" + } + + // 根据任务的最终状态更新统计(一次任务只记录一次) + var status string + if sm.Status == SendSuccess { + status = "success" + } else { + status = "failed" + } + + // 更新统计:每次任务执行记录为1次 + err := models.IncrementSendStats(taskID, taskType, currentDay, status, 1) + if err != nil { + logrus.Errorf("更新发送统计失败:%s", err) + } +} + +// getCurrentDay 获取当前日期(YYYY-MM-DD格式) +func (sm *SendMessageService) getCurrentDay() string { + return util.GetNowTimeStr()[:10] +} + // TransError 转化错误 func (sm *SendMessageService) TransError(err string) string { if err == "" { diff --git a/service/statistic_service/statistic.go b/service/statistic_service/statistic.go index 3b6cdb0..2687f8f 100644 --- a/service/statistic_service/statistic.go +++ b/service/statistic_service/statistic.go @@ -2,9 +2,12 @@ package statistic_service import ( "message-nest/models" + "strconv" ) type StatisticService struct { + TaskID string + Days int } func (sw *StatisticService) GetStatisticData() (models.StatisticData, error) { @@ -25,3 +28,27 @@ func (sw *StatisticService) GetTrendStatisticData() (models.TrendStatisticData, func (sw *StatisticService) GetChannelStatisticData() (models.ChannelStatisticData, error) { return models.GetChannelStatisticData() } + +// GetSendStatsData 获取发送统计数据(基于 send_stats 表) +func (sw *StatisticService) GetSendStatsData() (models.SendStatsData, error) { + days := sw.Days + if days <= 0 { + days = 30 // 默认30天 + } + return models.GetSendStatsData(days) +} + +// GetSendStatsByTask 获取指定任务的发送统计数据 +func (sw *StatisticService) GetSendStatsByTask() (models.SendStatsData, error) { + days := sw.Days + if days <= 0 { + days = 30 // 默认30天 + } + + taskID, err := strconv.ParseUint(sw.TaskID, 10, 64) + if err != nil { + return models.SendStatsData{}, err + } + + return models.GetSendStatsByTask(uint(taskID), days) +} diff --git a/web/src/components/pages/dashboard/Dashboard.vue b/web/src/components/pages/dashboard/Dashboard.vue index 89bdf9d..4eea437 100644 --- a/web/src/components/pages/dashboard/Dashboard.vue +++ b/web/src/components/pages/dashboard/Dashboard.vue @@ -30,10 +30,10 @@ let state = reactive({ today_failed_num: 0, }, trendData: { - latest_send_data: [] as SendData[], + latest_send_data: [] as SendData[] | null, }, channelData: { - way_cate_data: [] as CateData[], + way_cate_data: [] as CateData[] | null, }, loading: { basic: false, @@ -111,24 +111,25 @@ const loadAllStatisticData = async () => { } const renderLineChart = () => { + const latestSendData = state.trendData.latest_send_data || []; const options = { series: [ { name: '发送总数', - data: state.trendData.latest_send_data.length > 0 - ? state.trendData.latest_send_data.map(item => item.num || 0) + data: latestSendData.length > 0 + ? latestSendData.map(item => item.num || 0) : [] }, { name: '发送成功数', - data: state.trendData.latest_send_data.length > 0 - ? state.trendData.latest_send_data.map(item => item.day_succ_num || 0) + data: latestSendData.length > 0 + ? latestSendData.map(item => item.day_succ_num || 0) : [] }, { name: '发送失败数', - data: state.trendData.latest_send_data.length > 0 - ? state.trendData.latest_send_data.map(item => item.day_failed_num || 0) + data: latestSendData.length > 0 + ? latestSendData.map(item => item.day_failed_num || 0) : [] }, ], @@ -175,8 +176,8 @@ const renderLineChart = () => { } }, xaxis: { - categories: state.trendData.latest_send_data.length > 0 - ? state.trendData.latest_send_data.map(item => item.day) + categories: latestSendData.length > 0 + ? latestSendData.map(item => item.day) : [], axisBorder: { show: false @@ -332,9 +333,10 @@ const renderLineChart = () => { } const renderPieChart = () => { + const wayCateData = state.channelData.way_cate_data || []; const options = { - series: state.channelData.way_cate_data.length > 0 - ? state.channelData.way_cate_data.map(item => item.count_num) + series: wayCateData.length > 0 + ? wayCateData.map(item => item.count_num) : [], chart: { type: 'pie', @@ -355,8 +357,8 @@ const renderPieChart = () => { } } }, - labels: state.channelData.way_cate_data.length > 0 - ? state.channelData.way_cate_data.map(item => item.way_name) + labels: wayCateData.length > 0 + ? wayCateData.map(item => item.way_name) : [], colors: ['#3b82f6', '#10b981', '#f59e0b', '#ef4444', '#8b5cf6'], legend: {