You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.17 16:27:02