diff --git a/main.go b/main.go index 88aeee1..d4cd2e0 100644 --- a/main.go +++ b/main.go @@ -14,16 +14,12 @@ import ( "message-nest/routers" "message-nest/service/cron_service" "message-nest/service/env_service" - "message-nest/service/send_message_service" "net/http" - "sync" ) var ( //go:embed web/dist/* f embed.FS - - wg sync.WaitGroup ) func init() { @@ -36,7 +32,7 @@ func init() { cron_service.Setup() } -func GinServerUp() { +func main() { gin.SetMode(setting.ServerSetting.RunMode) routersInit := routers.InitRouter(f) @@ -60,13 +56,3 @@ func GinServerUp() { logrus.Error("Server err: ", err) } } - -func main() { - wg.Add(1) - - go GinServerUp() - go send_message_service.MessageConsumer(&wg) - - wg.Wait() - fmt.Println("Server exit...") -} diff --git a/pkg/constant/constant.go b/pkg/constant/constant.go index 07aa3d6..9b9e0ff 100644 --- a/pkg/constant/constant.go +++ b/pkg/constant/constant.go @@ -28,4 +28,4 @@ var LogsCleanDefaultValueMap = map[string]string{ const AboutSectionName = "about" // 限制goroutine的最大数量 -var MaxSemaphore = make(chan struct{}, 2048) +var MaxSendTaskSemaphoreChan = make(chan string, 2048) diff --git a/routers/api/v1/send_message.go b/routers/api/v1/send_message.go index 0acdd93..63f0cd8 100644 --- a/routers/api/v1/send_message.go +++ b/routers/api/v1/send_message.go @@ -47,7 +47,7 @@ func DoSendMassage(c *gin.Context) { appG.CResponse(http.StatusOK, "发送成功!", nil) } else { // 异步发送 - send_message_service.Buffer <- msgService + msgService.AsyncSend() appG.CResponse(http.StatusOK, "提交成功!", nil) } } diff --git a/service/send_message_service/consumer.go b/service/send_message_service/consumer.go deleted file mode 100644 index a1e8430..0000000 --- a/service/send_message_service/consumer.go +++ /dev/null @@ -1,43 +0,0 @@ -package send_message_service - -import ( - "github.com/sirupsen/logrus" - "message-nest/pkg/constant" - - "sync" -) - -var maxBufferSize = 10 -var Buffer = make(chan SendMessageService, maxBufferSize) - -func DoSendTask(task SendMessageService, wg *sync.WaitGroup) { - defer wg.Done() - defer func() { - if r := recover(); r != nil { - logrus.Error("DoSendTask: Recovered from panic:", r) - } - }() - - constant.MaxSemaphore <- struct{}{} - defer func() { - <-constant.MaxSemaphore - }() - - go task.Send() -} - -func MessageConsumer(wg *sync.WaitGroup) { - defer wg.Done() - - for { - task, ok := <-Buffer - if !ok { - logrus.Error("MessageConsumer: Channel closed. Exiting.") - return - } - - wg.Add(1) - go DoSendTask(task, wg) - } - -} diff --git a/service/send_message_service/send_message.go b/service/send_message_service/send_message.go index ae59ee7..63ea8b6 100644 --- a/service/send_message_service/send_message.go +++ b/service/send_message_service/send_message.go @@ -4,6 +4,7 @@ import ( "fmt" "github.com/sirupsen/logrus" "message-nest/models" + "message-nest/pkg/constant" "message-nest/service/send_task_service" "message-nest/service/send_way_service" "strings" @@ -14,6 +15,10 @@ const ( SendFail = 0 ) +//var logger = logrus.WithFields(logrus.Fields{ +// "prefix": "[Message]", +//}) + func errStrIsSuccess(errStr string) int { if errStr == "" { return SendSuccess @@ -41,6 +46,23 @@ func (sm *SendMessageService) LogsAndStatusMark(errStr string, status int) { logrus.Infof("%s, 状态:%d", strings.Trim(errStr, "\n"), status) } +// AsyncSend 异步发送一个消息任务的所有实例 +func (sm *SendMessageService) AsyncSend() { + defer func() { + if r := recover(); r != nil { + logrus.Error("AsyncSend: Recovered from panic:", r) + } + }() + + // 限制并发异步发送数量 + constant.MaxSendTaskSemaphoreChan <- "" + defer func() { + <-constant.MaxSendTaskSemaphoreChan + }() + + go sm.Send() +} + // Send 发送一个消息任务的所有实例 func (sm *SendMessageService) Send() string { sm.Status = SendSuccess