diff --git a/.release_log b/.release_log index 7412cb5..a497020 100644 --- a/.release_log +++ b/.release_log @@ -17,3 +17,4 @@ 17. 支持数据统计展示 18. 支持docker部署,从环境变量启动 19. 支持微信测试公众号模板消息发送 +20. 支持设置自定义的定时消息推送 diff --git a/main.go b/main.go index 3d69cae..64eb3c6 100644 --- a/main.go +++ b/main.go @@ -11,6 +11,7 @@ import ( "message-nest/pkg/logging" "message-nest/pkg/setting" "message-nest/routers" + "message-nest/service/cron_msg_service" "message-nest/service/cron_service" "net/http" "os" @@ -34,6 +35,7 @@ func init() { func main() { cron_service.StartLogsCronRun() + cron_msg_service.StartUpMsgCronTask() gin.SetMode(setting.ServerSetting.RunMode) routersInit := routers.InitRouter(f) diff --git a/migrate/migrate.go b/migrate/migrate.go index b473063..0056369 100644 --- a/migrate/migrate.go +++ b/migrate/migrate.go @@ -65,6 +65,7 @@ func Setup() { &models.SendTasksLogs{}, &models.SendTasksIns{}, &models.Settings{}, + &models.CronMessages{}, } for _, table := range tables { diff --git a/models/cron_messages.go b/models/cron_messages.go new file mode 100644 index 0000000..ea0c7d4 --- /dev/null +++ b/models/cron_messages.go @@ -0,0 +1,118 @@ +package models + +import ( + "errors" + "fmt" + "github.com/jinzhu/gorm" + "message-nest/pkg/util" +) + +type CronMessages struct { + UUIDModel + + Name string `json:"name" gorm:"type:varchar(200) comment '关联的消息名称';default:'';"` + TaskID string `json:"task_id" gorm:"type:varchar(36) comment '关联的消息ID';default:'';"` + Cron string `json:"cron" gorm:"type:varchar(4096) comment '定时表达式';default:'';"` + Title string `json:"title" gorm:"type:varchar(1000) comment '消息名称';default:'';"` + Content string `json:"content" gorm:"type:varchar(4096) comment '消息内容';default:'';"` + //MarkDown string `json:"markdown" gorm:"type:varchar(4096) comment 'markdown内容';default:'';"` + Url string `json:"url" gorm:"type:varchar(4096) comment 'url地址';default:'';"` + Enable int `json:"enable" gorm:"type:int comment '开启、暂停状态';default:1;"` +} + +func GenerateMsgUniqueID() string { + newUUID := util.GenerateUniqueID() + return fmt.Sprintf("C-%s", newUUID) +} + +func AddSendCronMsg( + name string, + task_id string, + cron string, + title string, + content string, + url string, + createdBy string, +) (string, error) { + newUUID := GenerateMsgUniqueID() + msg := CronMessages{ + UUIDModel: UUIDModel{ + ID: newUUID, + CreatedBy: createdBy, + ModifiedBy: createdBy, + }, + Name: name, + TaskID: task_id, + Cron: cron, + Title: title, + Content: content, + Url: url, + Enable: 1, + } + if err := db.Create(&msg).Error; err != nil { + return newUUID, err + } + return newUUID, nil +} + +// GetCronMessages 获取所有任务 +func GetCronMessages(pageNum int, pageSize int, name string, maps interface{}) ([]CronMessages, error) { + var ( + msgs []CronMessages + err error + ) + query := db.Where(maps) + if name != "" { + query = query.Where("name like ?", fmt.Sprintf("%%%s%%", name)) + } + query = query.Order("created_on DESC") + if pageSize > 0 || pageNum > 0 { + query = query.Offset(pageNum).Limit(pageSize) + } + err = query.Find(&msgs).Error + if err != nil && err != gorm.ErrRecordNotFound { + return nil, err + } + return msgs, nil +} + +// GetCronMessagesTotal 获取所有任务总数 +func GetCronMessagesTotal(name string, maps interface{}) (int, error) { + var ( + err error + total int + ) + query := db.Model(&CronMessages{}).Where(maps) + if name != "" { + query = query.Where("name like ?", fmt.Sprintf("%%%s%%", name)) + } + + err = query.Count(&total).Error + if err != nil { + return 0, err + } + return total, nil +} + +func DeleteCronMsg(id string) error { + if err := db.Where("id = ?", id).Delete(&CronMessages{}).Error; err != nil { + return err + } + return nil +} + +func EditCronMsg(id string, data interface{}) error { + if err := db.Model(&CronMessages{}).Where("id = ? ", id).Updates(data).Error; err != nil { + return err + } + return nil +} + +func GetCronMsgByID(id string) (CronMessages, error) { + var msg CronMessages + err := db.Where("id = ? ", id).Find(&msg).Error + if err != nil && errors.Is(err, gorm.ErrRecordNotFound) { + return msg, err + } + return msg, nil +} diff --git a/models/send_tasks.go b/models/send_tasks.go index b2b87ae..852461c 100644 --- a/models/send_tasks.go +++ b/models/send_tasks.go @@ -174,3 +174,12 @@ func EditSendTask(id string, data interface{}) error { return nil } + +func GetTaskByID(id string) (SendTasks, error) { + var task SendTasks + err := db.Where("id = ? ", id).Find(&task).Error + if err != nil && errors.Is(err, gorm.ErrRecordNotFound) { + return task, err + } + return task, nil +} diff --git a/pkg/constant/constant.go b/pkg/constant/constant.go index 6890784..bc1e5c1 100644 --- a/pkg/constant/constant.go +++ b/pkg/constant/constant.go @@ -1,5 +1,7 @@ package constant +import "github.com/robfig/cron/v3" + const CleanLogsTaskId = "T-IM1GBswSRY" const SiteSettingSectionName = "site_config" @@ -31,4 +33,5 @@ const AboutSectionName = "about" // 限制goroutine的最大数量 var MaxSendTaskSemaphoreChan = make(chan string, 2048) -// 内存缓存存放所有的微信公共号的缓存实例 +// cron消息任务内存缓存map +var CronMsgIdMapMemoryCache = make(map[string]cron.EntryID) diff --git a/routers/api/v1/cron_msg.go b/routers/api/v1/cron_msg.go new file mode 100644 index 0000000..b702aa8 --- /dev/null +++ b/routers/api/v1/cron_msg.go @@ -0,0 +1,163 @@ +package v1 + +import ( + "message-nest/pkg/e" + "message-nest/pkg/util" + "net/http" + + "github.com/gin-gonic/gin" + "message-nest/pkg/app" + "message-nest/service/cron_msg_service" +) + +type DeleteCronMsgTaskReq struct { + ID string `json:"id" validate:"required,len=12" label:"任务id"` +} + +// DeleteCronMsgTask 删除定时消息 +func DeleteCronMsgTask(c *gin.Context) { + var ( + appG = app.Gin{C: c} + req DeleteCronMsgTaskReq + ) + + errCode, errMsg := app.BindJsonAndPlayValid(c, &req) + if errCode != e.SUCCESS { + appG.CResponse(errCode, errMsg, nil) + return + } + + CronMsgService := cron_msg_service.CronMsgService{ + ID: req.ID, + } + msg, _ := CronMsgService.GetByID() + + err := CronMsgService.Delete() + if err != nil { + appG.CResponse(http.StatusBadRequest, "删除定时消息失败!", nil) + return + } + cron_msg_service.RemoveCronMsgToCronServer(msg) + appG.CResponse(http.StatusOK, "删除定时消息成功!", nil) +} + +// GetCronMsgList 获取定时消息列表 +func GetCronMsgList(c *gin.Context) { + appG := app.Gin{C: c} + name := c.Query("name") + + offset, limit := util.GetPageSize(c) + CronMsgService := cron_msg_service.CronMsgService{ + Name: name, + PageNum: offset, + PageSize: limit, + } + tasks, err := CronMsgService.GetAll() + if err != nil { + appG.CResponse(http.StatusInternalServerError, "获取定时消息失败!", nil) + return + } + + count, err := CronMsgService.Count() + if err != nil { + appG.CResponse(http.StatusInternalServerError, "获取定时消息总数失败!", nil) + return + } + + appG.CResponse(http.StatusOK, "获取定时消息成功", map[string]interface{}{ + "lists": tasks, + "total": count, + }) +} + +type AddCronMsgTaskReq struct { + Name string `json:"name" validate:"required,max=100,min=1" label:"消息名称"` + TaskID string `json:"task_id" validate:"required,max=100,min=1" label:"关联的任务id"` + Cron string `json:"cron" validate:"required,cron" label:"cron表达式"` + Title string `json:"title" validate:"required,max=100,min=1" label:"消息标题"` + Content string `json:"content" validate:"required,max=10000,min=1" label:"消息内容"` + Url string `json:"url" validate:"" label:"消息详情地址"` +} + +// AddCronMsgTask 添加定时消息 +func AddCronMsgTask(c *gin.Context) { + var ( + appG = app.Gin{C: c} + req AddCronMsgTaskReq + ) + + currentUser := app.GetCurrentUserName(c) + errCode, errStr := app.BindJsonAndPlayValid(c, &req) + if errCode != e.SUCCESS { + appG.CResponse(errCode, errStr, nil) + return + } + + CronMsgService := cron_msg_service.CronMsgService{ + Name: req.Name, + TaskID: req.TaskID, + Cron: req.Cron, + Title: req.Title, + Content: req.Content, + Url: req.Url, + CreatedBy: currentUser, + ModifiedBy: currentUser, + } + + uuidstr, err := CronMsgService.Add() + if err != nil { + appG.CResponse(http.StatusBadRequest, "添加定时消息失败!", nil) + return + } + CronMsgService.ID = uuidstr + msg, _ := CronMsgService.GetByID() + cron_msg_service.UpdateCronMsgToCronServer(msg) + appG.CResponse(http.StatusOK, "添加定时消息成功!", nil) + +} + +type EditCronMsgTaskReq struct { + ID string `json:"id" validate:"required,len=12" label:"定时消息id"` + Name string `json:"name" validate:"required,max=100,min=1" label:"消息名称"` + TaskID string `json:"task_id" validate:"required,max=100,min=1" label:"关联的任务id"` + Cron string `json:"cron" validate:"required,cron" label:"cron表达式"` + Title string `json:"title" validate:"required,max=100,min=1" label:"消息标题"` + Content string `json:"content" validate:"required,max=10000,min=1" label:"消息内容"` + Url string `json:"url" validate:"" label:"消息详情地址"` + Enable int `json:"enable" validate:"oneof=0 1" label:"是否开启"` +} + +// EditCronMsgTask 编辑定时消息任务 +func EditCronMsgTask(c *gin.Context) { + var ( + appG = app.Gin{C: c} + req EditCronMsgTaskReq + ) + + errCode, errMsg := app.BindJsonAndPlayValid(c, &req) + if errCode != e.SUCCESS { + appG.CResponse(errCode, errMsg, nil) + return + } + + CronMsgService := cron_msg_service.CronMsgService{ + ID: req.ID, + } + + data := make(map[string]interface{}) + data["name"] = req.Name + data["task_id"] = req.TaskID + data["cron"] = req.Cron + data["title"] = req.Title + data["content"] = req.Content + data["url"] = req.Url + data["enable"] = req.Enable + err := CronMsgService.Edit(data) + if err != nil { + appG.CResponse(http.StatusBadRequest, "编辑定时消息失败!", nil) + return + } + msg, _ := CronMsgService.GetByID() + cron_msg_service.UpdateCronMsgToCronServer(msg) + appG.CResponse(http.StatusOK, "编辑定时消息成功!", nil) +} diff --git a/routers/api/v1/send_task.go b/routers/api/v1/send_task.go index 287e4d3..cbfedf8 100644 --- a/routers/api/v1/send_task.go +++ b/routers/api/v1/send_task.go @@ -134,3 +134,26 @@ func EditMsgSendTask(c *gin.Context) { appG.CResponse(http.StatusOK, "编辑发信任务成功!", nil) } + +// GetMsgSendTask 获取消息任务 +func GetMsgSendTask(c *gin.Context) { + appG := app.Gin{C: c} + id := c.Query("id") + + if id == "" { + appG.CResponse(http.StatusBadRequest, "任务id为空!", nil) + return + } + + sendTaskService := send_task_service.SendTaskService{ + ID: id, + } + + task, err := sendTaskService.GetByID() + if err != nil { + appG.CResponse(http.StatusBadRequest, "获取到的任务信息为空!", nil) + return + } + + appG.CResponse(http.StatusOK, "获取任务信息成功", task) +} diff --git a/routers/router.go b/routers/router.go index 52c858e..1254f93 100644 --- a/routers/router.go +++ b/routers/router.go @@ -62,6 +62,7 @@ func InitRouter(f embed.FS) *gin.Engine { apiV1.POST("/sendtasks/add", v1.AddMsgSendTask) apiV1.POST("/sendtasks/delete", v1.DeleteMsgSendTask) apiV1.POST("/sendtasks/edit", v1.EditMsgSendTask) + apiV1.GET("/sendtasks/get", v1.GetMsgSendTask) // sendtasks/ins apiV1.POST("/sendtasks/ins/addmany", v1.AddManyTasksIns) @@ -84,6 +85,12 @@ func InitRouter(f embed.FS) *gin.Engine { // statistic apiV1.GET("/statistic", v1.GetStatisticData) + // cronMessage + apiV1.POST("/sendmessages/addone", v1.AddCronMsgTask) + apiV1.GET("/sendmessages/list", v1.GetCronMsgList) + apiV1.POST("/sendmessages/delete", v1.DeleteCronMsgTask) + apiV1.POST("/sendmessages/edit", v1.EditCronMsgTask) + } return app diff --git a/service/cron_msg_service/cron_msg.go b/service/cron_msg_service/cron_msg.go new file mode 100644 index 0000000..db84c35 --- /dev/null +++ b/service/cron_msg_service/cron_msg.go @@ -0,0 +1,56 @@ +package cron_msg_service + +import ( + "message-nest/models" +) + +type CronMsgService struct { + ID string + + Name string + TaskID string + Cron string + Title string + Content string + Url string + + CreatedBy string + ModifiedBy string + CreatedOn string + + PageNum int + PageSize int +} + +func (st *CronMsgService) Add() (string, error) { + return models.AddSendCronMsg(st.Name, st.TaskID, st.Cron, st.Title, st.Content, st.Url, st.CreatedBy) +} + +func (st *CronMsgService) Edit(data map[string]interface{}) error { + return models.EditCronMsg(st.ID, data) +} + +func (st *CronMsgService) GetByID() (models.CronMessages, error) { + return models.GetCronMsgByID(st.ID) +} + +func (st *CronMsgService) Count() (int, error) { + return models.GetCronMessagesTotal(st.Name, st.getMaps()) +} + +func (st *CronMsgService) GetAll() ([]models.CronMessages, error) { + msgs, err := models.GetCronMessages(st.PageNum, st.PageSize, st.Name, st.getMaps()) + if err != nil { + return nil, err + } + return msgs, nil +} + +func (st *CronMsgService) getMaps() map[string]interface{} { + maps := make(map[string]interface{}) + return maps +} + +func (st *CronMsgService) Delete() error { + return models.DeleteCronMsg(st.ID) +} diff --git a/service/cron_msg_service/startup_cron_msg.go b/service/cron_msg_service/startup_cron_msg.go new file mode 100644 index 0000000..d9a3b1a --- /dev/null +++ b/service/cron_msg_service/startup_cron_msg.go @@ -0,0 +1,102 @@ +package cron_msg_service + +import ( + "github.com/sirupsen/logrus" + "message-nest/models" + "message-nest/pkg/constant" + "message-nest/service/cron_service" + "message-nest/service/send_message_service" +) + +type MsgCronTask struct { +} + +func (s MsgCronTask) Register() { + // 获取所有的定时消息任务 + limit := 10000 + filter := make(map[string]interface{}) + filter["enable"] = 1 + data, err := models.GetCronMessages(0, limit, "", filter) + if err != nil { + logrus.Errorf("获取定时消息任务失败!原因:%s", err.Error()) + return + } + if len(data) == 0 { + logrus.Infof("没有定时消息任务需要注册") + return + } + //注册定时任务 + for _, msg := range data { + AddCronMsgToCronServer(msg) + } + length := len(data) + if length > 0 { + logrus.Infof("完成用户自定义的定时消息注册,个数:%d", length) + } +} + +// AddCronMsgToCronServer 注册定时任务到定时服务 +func AddCronMsgToCronServer(msg models.CronMessages) { + if msg.Enable != 1 { + return + } + taskId := cron_service.AddTask(cron_service.ScheduledTask{ + Schedule: msg.Cron, + Job: func() { + CronMsgSendF(msg) + }, + }) + constant.CronMsgIdMapMemoryCache[msg.ID] = taskId +} + +// 执行任务的构造函数 +func CronMsgSendF(msg models.CronMessages) { + logrus.Infof("开始只能执行定时消息发送任务: %s , 消息名: %s", msg.ID, msg.Name) + task, err := models.GetTaskByID(msg.TaskID) + if err != nil { + logrus.Infof("消息任务不存在: %s ", msg.TaskID) + return + } + sender := send_message_service.SendMessageService{ + TaskID: task.ID, + Title: msg.Title, + Text: msg.Content, + URL: msg.Url, + CallerIp: "corn", + DefaultLogger: logrus.WithFields(logrus.Fields{ + "prefix": "[Cron Message]", + }), + } + taskData, _ := sender.SendPreCheck() + sender.Send(taskData) +} + +// UpdateCronMsgToCronServer 更新定时服务的任务 +func UpdateCronMsgToCronServer(msg models.CronMessages) { + if entryId, ok := constant.CronMsgIdMapMemoryCache[msg.ID]; ok { + // 先删除之前的定时任务 + delete(constant.CronMsgIdMapMemoryCache, msg.ID) + cron_service.RemoveTask(entryId) + // 再注册新的定时任务 + AddCronMsgToCronServer(msg) + } else { + // 注册新的定时任务 + AddCronMsgToCronServer(msg) + } + logrus.Infof("完成定时消息的定时更新,消息id: %s, 总数:%d", msg.ID, len(constant.CronMsgIdMapMemoryCache)) +} + +// RemoveCronMsgToCronServer 删除定时任务中心的任务 +func RemoveCronMsgToCronServer(msg models.CronMessages) { + if entryId, ok := constant.CronMsgIdMapMemoryCache[msg.ID]; ok { + // 先删除之前的定时任务 + delete(constant.CronMsgIdMapMemoryCache, msg.ID) + cron_service.RemoveTask(entryId) + } + logrus.Infof("完成定时消息的定时删除,消息id: %s, 总数:%d", msg.ID, len(constant.CronMsgIdMapMemoryCache)) +} + +// StartUpMsgCronTask 启动注册定时任务 +func StartUpMsgCronTask() { + MsgCronTask{}.Register() +} diff --git a/service/cron_service/tasks.go b/service/cron_service/tasks.go index ae74315..2d4cebd 100644 --- a/service/cron_service/tasks.go +++ b/service/cron_service/tasks.go @@ -46,7 +46,7 @@ func AddTask(task ScheduledTask) cron.EntryID { logrus.Errorf("注册定时任务失败, job: %s, 原因:%s", jobName, err) } else { TaskList[taskId] = &task - logrus.Infof("注册定时任务成功, job: %s, entryID: %d, cron: %s", jobName, taskId, task.Schedule) + //logrus.Infof("注册定时任务成功, job: %s, entryID: %d, cron: %s", jobName, taskId, task.Schedule) } return taskId } diff --git a/service/send_task_service/send_task.go b/service/send_task_service/send_task.go index f917933..42e4d6a 100644 --- a/service/send_task_service/send_task.go +++ b/service/send_task_service/send_task.go @@ -51,3 +51,7 @@ func (st *SendTaskService) getMaps() map[string]interface{} { maps := make(map[string]interface{}) return maps } + +func (st *SendTaskService) GetByID() (interface{}, error) { + return models.GetTaskByID(st.ID) +} diff --git a/web/src/router/index.js b/web/src/router/index.js index 63effbd..88bfd91 100644 --- a/web/src/router/index.js +++ b/web/src/router/index.js @@ -27,6 +27,11 @@ const router = createRouter({ name: 'sendtasks', component: () => import('../views/tabsTools/sendTasks/sendTasks.vue') }, + { + path: '/sendmessages', + name: 'sendmessages', + component: () => import('../views/tabsTools/sendMessage/sensMessage.vue') + }, { path: '/sendlogs', name: 'sendlogs', diff --git a/web/src/views/home/index.vue b/web/src/views/home/index.vue index 3017582..aec180b 100644 --- a/web/src/views/home/index.vue +++ b/web/src/views/home/index.vue @@ -35,9 +35,10 @@ export default { const menuData = reactive([ { id: '0', title: '数据统计', path: '/statistic' }, { id: '1', title: '发信日志', path: '/sendlogs' }, - { id: '2', title: '发信任务', path: '/sendtasks' }, - { id: '3', title: '发信渠道', path: '/sendways' }, - { id: '4', title: '设置', path: '/settings' }, + { id: '2', title: '定时发信', path: '/sendmessages' }, + { id: '3', title: '发信任务', path: '/sendtasks' }, + { id: '4', title: '发信渠道', path: '/sendways' }, + { id: '5', title: '设置', path: '/settings' }, ]); const checkIsLogin = () => { diff --git a/web/src/views/tabsTools/sendMessage/sensMessage.vue b/web/src/views/tabsTools/sendMessage/sensMessage.vue new file mode 100644 index 0000000..ecbf3d7 --- /dev/null +++ b/web/src/views/tabsTools/sendMessage/sensMessage.vue @@ -0,0 +1,204 @@ + + + + + \ No newline at end of file diff --git a/web/src/views/tabsTools/sendMessage/view/addCronMsgPopUp.vue b/web/src/views/tabsTools/sendMessage/view/addCronMsgPopUp.vue new file mode 100644 index 0000000..5463587 --- /dev/null +++ b/web/src/views/tabsTools/sendMessage/view/addCronMsgPopUp.vue @@ -0,0 +1,211 @@ + + + + + diff --git a/web/src/views/tabsTools/sendMessage/view/editCronMsgPopUp.vue b/web/src/views/tabsTools/sendMessage/view/editCronMsgPopUp.vue new file mode 100644 index 0000000..006063a --- /dev/null +++ b/web/src/views/tabsTools/sendMessage/view/editCronMsgPopUp.vue @@ -0,0 +1,229 @@ + + + + +