Files
ai-agent/internal/services/channel_message_outbox_service.go
T
t 18c9354095 refactor: 将客服后端重构为宿主可嵌入模块
- 注入数据库、运行时配置、统一响应、文件存储和平台 AI 能力,补充业务读写工具与客户快捷操作契约。

- 移除模块内重复的组织、客户、工单、标签、技能、旧工作流、MCP 和迁移实现,将身份权限与业务主体交由宿主管理。

- 使用 libSQL 重构向量存储,并完善图片消息、访客身份、排队调度、企业微信和支持聊天页面。

- 统一 HTTP、DTO 与 WebSocket 的 snake_case 协议,补齐模块初始化、业务动作和公共载荷等回归测试。
2026-08-28 22:23:13 +08:00

185 lines
6.4 KiB
Go

package services
import (
"code.tczkiot.com/wlw/ai-agent/internal/models"
"code.tczkiot.com/wlw/ai-agent/internal/pkg/dto"
"code.tczkiot.com/wlw/ai-agent/internal/pkg/enums"
"code.tczkiot.com/wlw/ai-agent/internal/pkg/errorsx"
"code.tczkiot.com/wlw/ai-agent/internal/repositories"
"encoding/json"
"strings"
"time"
"code.tczkiot.com/wlw/ai-agent/internal/pkg/httpx/params"
"github.com/mlogclub/simple/sqls"
)
var ChannelMessageOutboxService = newChannelMessageOutboxService()
func newChannelMessageOutboxService() *channelMessageOutboxService {
return &channelMessageOutboxService{}
}
type channelMessageOutboxService struct {
}
func (s *channelMessageOutboxService) Get(id int64) *models.ChannelMessageOutbox {
return repositories.ChannelMessageOutboxRepository.Get(sqls.DB(), id)
}
func (s *channelMessageOutboxService) Take(where ...interface{}) *models.ChannelMessageOutbox {
return repositories.ChannelMessageOutboxRepository.Take(sqls.DB(), where...)
}
func (s *channelMessageOutboxService) Find(cnd *sqls.Cnd) []models.ChannelMessageOutbox {
return repositories.ChannelMessageOutboxRepository.Find(sqls.DB(), cnd)
}
func (s *channelMessageOutboxService) FindOne(cnd *sqls.Cnd) *models.ChannelMessageOutbox {
return repositories.ChannelMessageOutboxRepository.FindOne(sqls.DB(), cnd)
}
func (s *channelMessageOutboxService) FindPageByParams(params *params.QueryParams) (list []models.ChannelMessageOutbox, paging *sqls.Paging) {
return repositories.ChannelMessageOutboxRepository.FindPageByParams(sqls.DB(), params)
}
func (s *channelMessageOutboxService) FindPageByCnd(cnd *sqls.Cnd) (list []models.ChannelMessageOutbox, paging *sqls.Paging) {
return repositories.ChannelMessageOutboxRepository.FindPageByCnd(sqls.DB(), cnd)
}
func (s *channelMessageOutboxService) Count(cnd *sqls.Cnd) int64 {
return repositories.ChannelMessageOutboxRepository.Count(sqls.DB(), cnd)
}
func (s *channelMessageOutboxService) Create(t *models.ChannelMessageOutbox) error {
return repositories.ChannelMessageOutboxRepository.Create(sqls.DB(), t)
}
func (s *channelMessageOutboxService) Update(t *models.ChannelMessageOutbox) error {
return repositories.ChannelMessageOutboxRepository.Update(sqls.DB(), t)
}
func (s *channelMessageOutboxService) Updates(id int64, columns map[string]interface{}) error {
return repositories.ChannelMessageOutboxRepository.Updates(sqls.DB(), id, columns)
}
func (s *channelMessageOutboxService) UpdateColumn(id int64, name string, value interface{}) error {
return repositories.ChannelMessageOutboxRepository.UpdateColumn(sqls.DB(), id, name, value)
}
func (s *channelMessageOutboxService) Delete(id int64) {
repositories.ChannelMessageOutboxRepository.Delete(sqls.DB(), id)
}
// GetByMessageID retrieves the outbox entry by message ID and channel type.
func (s *channelMessageOutboxService) GetByMessageID(channelType string, messageID int64) *models.ChannelMessageOutbox {
return repositories.ChannelMessageOutboxRepository.Take(sqls.DB(), "channel_type = ? AND message_id = ?", channelType, messageID)
}
func (s *channelMessageOutboxService) EnqueueWxWorkKFMessage(conversation *models.Conversation, message *models.Message) error {
if conversation == nil || message == nil {
return nil
}
channel := ChannelService.Get(conversation.ChannelID)
if channel == nil || channel.ChannelType != enums.ChannelTypeWxWorkKF {
return nil
}
if message.SenderType != enums.IMSenderTypeAgent && message.SenderType != enums.IMSenderTypeAI {
return nil
}
if message.MessageType != enums.IMMessageTypeText && message.MessageType != enums.IMMessageTypeHTML {
return nil
}
if existing := s.GetByMessageID(enums.ChannelTypeWxWorkKF, message.ID); existing != nil {
return nil
}
payload, err := json.Marshal(map[string]any{
"conversation_id": conversation.ID,
"message_id": message.ID,
"message_type": message.MessageType,
"content": strings.TrimSpace(message.Content),
"payload": strings.TrimSpace(message.Payload),
"sender_id": message.SenderID,
})
if err != nil {
return err
}
now := time.Now()
return s.Create(&models.ChannelMessageOutbox{
ChannelType: enums.ChannelTypeWxWorkKF,
ConversationID: conversation.ID,
MessageID: message.ID,
Payload: string(payload),
SendStatus: string(enums.ChannelMessageOutboxStatusPending),
AuditFields: models.AuditFields{
CreatedAt: now,
CreateUserID: message.UpdateUserID,
CreateUserName: message.UpdateUserName,
UpdatedAt: now,
UpdateUserID: message.UpdateUserID,
UpdateUserName: message.UpdateUserName,
},
})
}
func (s *channelMessageOutboxService) ListPending(channelType string, limit int) []models.ChannelMessageOutbox {
if limit <= 0 {
limit = 20
}
cnd := sqls.NewCnd().
Eq("channel_type", strings.TrimSpace(channelType)).
In("send_status", []string{
string(enums.ChannelMessageOutboxStatusPending),
string(enums.ChannelMessageOutboxStatusFailed),
}).
Asc("id").
Limit(limit)
return s.Find(cnd)
}
func (s *channelMessageOutboxService) RetryWxWorkFailure(id int64, operator *dto.AuthPrincipal) error {
item := s.Get(id)
if item == nil || item.ChannelType != enums.ChannelTypeWxWorkKF {
return errorsx.InvalidParam("outbox record does not exist")
}
if item.SendStatus != string(enums.ChannelMessageOutboxStatusFailed) &&
item.SendStatus != string(enums.ChannelMessageOutboxStatusIgnored) {
return errorsx.InvalidParam("only failed or ignored outbox records can be retried")
}
now := time.Now()
columns := map[string]interface{}{
"send_status": string(enums.ChannelMessageOutboxStatusPending),
"next_retry_at": nil,
"updated_at": now,
}
if operator != nil {
columns["update_user_id"] = operator.UserID
columns["update_user_name"] = operator.Username
}
return s.Updates(id, columns)
}
func (s *channelMessageOutboxService) IgnoreWxWorkFailure(id int64, operator *dto.AuthPrincipal) error {
item := s.Get(id)
if item == nil || item.ChannelType != enums.ChannelTypeWxWorkKF {
return errorsx.InvalidParam("outbox record does not exist")
}
if item.SendStatus != string(enums.ChannelMessageOutboxStatusFailed) {
return errorsx.InvalidParam("only failed outbox records can be ignored")
}
now := time.Now()
columns := map[string]interface{}{
"send_status": string(enums.ChannelMessageOutboxStatusIgnored),
"next_retry_at": nil,
"updated_at": now,
}
if operator != nil {
columns["update_user_id"] = operator.UserID
columns["update_user_name"] = operator.Username
}
return s.Updates(id, columns)
}