Go语言定时消息调度循环中CreateSmsMessage函数执行异常
Go定时消息调度器:CreateSmsMessage执行异常问题
我正在Go应用中实现一个无限循环的消息调度器,用于遍历数据库获取定时消息并发送。代码可正常执行,但CreateSmsMessage函数存在未完全执行甚至不执行的情况。以下是相关代码:
调度循环函数代码
func SchedularLoop() { ctx := context.Background() stop := make(chan struct{}) schedularChan := make(chan []models.MessageScheduler) errorChan := make(chan error) createMessageChan := make(chan *models.Message) var wg sync.WaitGroup var createMessageWG sync.WaitGroup db := databaseapi.NewApiConfig{} messageService := messages.NewMessage(db.Postgres_DB_Api(), ctx) go func() { defer close(schedularChan) defer close(errorChan) for { data, err := GetAllSchedules(ctx) if err != nil { errorChan <- fmt.Errorf("error: %v", err) } schedularChan <- data time.Sleep(5 * time.Second) } }() go func() { for msgs := range schedularChan { currentTime := time.Now() TodayDate := utilities.FormatDate(currentTime) fmt.Println("Today's Date", TodayDate) for _, data := range msgs { schDate := utilities.FormatDate(data.ScheduledDate) if schDate == TodayDate { messageData, err := messageService.GetMessage(data.MessageID, ctx) if err != nil { errorChan <- err } if !messageData.IsSms { fmt.Println("I'm not a voice message") } smsMessage := &models.Message{ UserID: data.UserID, GroupID: messageData.GroupID, ID: uuid.New(), MessageTitle: messageData.MessageTitle, IsScheduled: false, Recipients: messageData.Recipients, Content: messageData.Content, Sender: messageData.Sender, ScheduledDate: time.Now().String(), CreatedAt: messageData.CreatedAt, UpdatedAt: time.Now(), } createMessageWG.Add(1) go func() { defer createMessageWG.Done() fmt.Println("Createmessage routine") CrmsgData, err := messageService.CreateSmsMessage(smsMessage, ctx) fmt.Println("Sent Data", CrmsgData) if err != nil { fmt.Println("Error creating SMS message:", err.Error()) errorChan <- err } else { createMessageChan <- CrmsgData fmt.Println("Testing Message Created", createMessageChan) } }() } } } }() go func() { for msg := range createMessageChan { fmt.Println(msg) fmt.Println("Create Message working oooo") } wg.Done() }() wg.Add(1) wg.Wait() defer close(stop) }
CreateSmsMessage函数代码
func (msgService *MessagesServiceImplementation) CreateSmsMessage(msg *models.Message, ctx context.Context) (*models.Message, error) { fmt.Println(msg.ScheduledDate) userService := usersimpl.NewUserService(msgService.pg, ctx) _, err := userService.GetUserById(msg.UserID, ctx) if err != nil { return nil, fmt.Errorf("error: %v", err) } errorChannel := make(chan error) dbChanel := make(chan database_utils.Message) msgChannel := make(chan *mnotify.SMSData) if msg.GroupID.UUID != uuid.Nil { // Search/Check for the existence of groupID groupService := group.NewGroupService(msgService.pg, ctx) // Check for group existence groupId := msg.GroupID.UUID _, err := groupService.GetGroupById(groupId, ctx) if err != nil { return nil, fmt.Errorf("group does not exist Reason: %v", err) } // If the group exist, lets get all the subscribers that belong to the user phoneNumbers var sub models.Subscribers sub.UserID = msg.UserID sub.GroupID = msg.GroupID subscriberService := subscribers.NewUserService(msgService.pg, ctx) phoneNumbers, err := subscriberService.GetAllUserSubscribersPhoneByGroupID(&sub, ctx) if err != nil { return nil, fmt.Errorf("error: %v", err) } // Lets Pass the users in the group phone numbers msg.Recipients = phoneNumbers params := &database_utils.CreateSMSMessageParams{ ID: uuid.New(), MessageTitle: msg.MessageTitle, IsSms: true, Sender: msg.Sender, GroupID: msg.GroupID, Recipients: msg.Recipients, Content: msg.Content, UserID: msg.UserID, MsgLanguage: msg.Title, CreatedAt: msg.CreatedAt, ScheduledDate: msg.ScheduledDate, IsScheduled: msg.IsScheduled, } go func() { data, err := msgService.pg.Postgres_DB_Api().DB.CreateSMSMessage(ctx, *params) if err != nil { errorChannel <- fmt.Errorf("error: %v", err) return } dbChanel <- data }() // Lets get the sender by ID var sender models.SmsSender sender.UserID = msg.UserID sender.ID = msg.Sender.UUID senderService := smssendergo.NewSmsSenderService(msgService.pg, ctx) senderData, err := senderService.GetSmsSenderAndUserID(&sender, ctx) if err != nil { return nil, fmt.Errorf("error: %v", err) } var mnotifyMsg models.QuickBulkSms mnotifyMsg.Is_Schedule = false mnotifyMsg.Message = msg.Content mnotifyMsg.Recipient = msg.Recipients mnotifyMsg.Sender = senderData.Name // Lets send the message go func() { mnotifyService := mnotifycampaign.NewCampaign(msgService.pg, ctx) data, sendErr := mnotifyService.CreateQuickBulkSMS(&mnotifyMsg) if sendErr != nil { errorChannel <- fmt.Errorf("error: %v", sendErr) } msgChannel <- data }() select { case data := <-dbChanel: return models.DbMsgtoDbMsg(data), nil case errorChnx := <-errorChannel: return nil, errorChnx case <-time.After(time.Second * 5): return nil, errors.New("timeout error") } } if msg.GroupID.UUID == uuid.Nil { params := &database_utils.CreateSMSMessageParams{ ID: uuid.New(), MessageTitle: msg.MessageTitle, IsSms: true, Sender: msg.Sender, GroupID: msg.GroupID, Recipients: msg.Recipients, Content: msg.Content, UserID: msg.UserID, MsgLanguage: msg.Title, CreatedAt: msg.CreatedAt, UpdatedAt: msg.UpdatedAt, IsScheduled: msg.IsScheduled, } go func() { data, err := msgService.pg.Postgres_DB_Api().DB.CreateSMSMessage(ctx, *params) if err != nil { errorChannel <- fmt.Errorf("error: %v", err.Error()) return } dbChanel <- data }() // // Lets get the sender by ID var sender models.SmsSender sender.UserID = msg.UserID sender.ID = msg.Sender.UUID senderService := smssendergo.NewSmsSenderService(msgService.pg, ctx) senderData, err := senderService.GetSmsSenderAndUserID(&sender, ctx) if err != nil { return nil, fmt.Errorf("error: %v", err.Error()) } var mnotifyMsg models.QuickBulkSms mnotifyMsg.Is_Schedule = false mnotifyMsg.Message = msg.Content mnotifyMsg.Recipient = msg.Recipients mnotifyMsg.Sender = senderData.Name // // Lets send the message go func() { mnotifyService := mnotifycampaign.NewCampaign(msgService.pg, ctx) response, sendErr := mnotifyService.CreateQuickBulkSMS(&mnotifyMsg) if sendErr != nil { errorChannel <- fmt.Errorf("error: %v", sendErr.Error()) } msgChannel <- response }() select { case data := <-dbChanel: return models.DbMsgtoDbMsg(data), nil case errorChnx := <-errorChannel: return nil, errorChnx case <-time.After(time.Second * 5): return nil, errors.New("timeout error") } } return nil, errors.New("invalid request") }
内容的提问来源于stack exchange,提问作者Daniel Dexoangle
相关产品推荐
相关产品推荐

