Files
Message-Push-Nest/service/send_message_service/consumer.go
T
2024-01-03 19:55:38 +08:00

43 lines
725 B
Go

package send_message_service
import (
"message-nest/pkg/constant"
"message-nest/pkg/logging"
"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 {
logging.Logger.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 {
logging.Logger.Error("MessageConsumer: Channel closed. Exiting.")
return
}
wg.Add(1)
go DoSendTask(task, wg)
}
}