diff --git a/internal/bootstrap/routes.go b/internal/bootstrap/routes.go index 9843f9ec..374c1479 100644 --- a/internal/bootstrap/routes.go +++ b/internal/bootstrap/routes.go @@ -501,3 +501,20 @@ func registerThirdTikTokRoutes(group *gin.RouterGroup) { group.POST("/webhook", third.TikTokPostWebhook) group.POST("/webhook/:channel_id", third.TikTokPostWebhook) } + +func registerThirdLineRoutes(group *gin.RouterGroup) { + group.POST("/webhook", third.LinePostWebhook) + group.POST("/webhook/:channel_id", third.LinePostWebhook) +} + +func registerThirdViberRoutes(group *gin.RouterGroup) { + group.POST("/webhook", third.ViberPostWebhook) + group.POST("/webhook/:channel_id", third.ViberPostWebhook) +} + +func registerThirdThreadsRoutes(group *gin.RouterGroup) { + group.GET("/webhook", third.ThreadsGetWebhook) + group.GET("/webhook/:channel_id", third.ThreadsGetWebhook) + group.POST("/webhook", third.ThreadsPostWebhook) + group.POST("/webhook/:channel_id", third.ThreadsPostWebhook) +} diff --git a/internal/bootstrap/server.go b/internal/bootstrap/server.go index 76b632bc..7d2f798a 100644 --- a/internal/bootstrap/server.go +++ b/internal/bootstrap/server.go @@ -205,6 +205,9 @@ func addRouter(app *gin.Engine) { registerThirdSlackRoutes(thirdGroup.Group("/slack")) registerThirdXRoutes(thirdGroup.Group("/x")) registerThirdTikTokRoutes(thirdGroup.Group("/tiktok")) + registerThirdLineRoutes(thirdGroup.Group("/line")) + registerThirdViberRoutes(thirdGroup.Group("/viber")) + registerThirdThreadsRoutes(thirdGroup.Group("/threads")) } type spaShellRewrite struct { diff --git a/internal/handlers/third/line_handler.go b/internal/handlers/third/line_handler.go new file mode 100644 index 00000000..431914f5 --- /dev/null +++ b/internal/handlers/third/line_handler.go @@ -0,0 +1,36 @@ +package third + +import ( + "bytes" + "io" + "net/http" + "strings" + + "agent-desk/internal/services" + + "github.com/gin-gonic/gin" +) + +// LinePostWebhook receives incoming webhook events from the LINE Platform. +func LinePostWebhook(ctx *gin.Context) { + channelID := strings.TrimSpace(ctx.Param("channel_id")) + if channelID == "" { + channelID = strings.TrimSpace(ctx.Query("channel_id")) + } + + signature := ctx.GetHeader("X-Line-Signature") + + bodyBytes, err := io.ReadAll(ctx.Request.Body) + if err != nil { + ctx.JSON(http.StatusBadRequest, gin.H{"ok": false, "error": "failed to read body"}) + return + } + ctx.Request.Body = io.NopCloser(bytes.NewBuffer(bodyBytes)) + + if err := services.LineInboundService.HandleWebhook(ctx.Request.Context(), channelID, signature, bodyBytes); err != nil { + ctx.JSON(http.StatusOK, gin.H{"ok": false, "error": err.Error()}) + return + } + + ctx.JSON(http.StatusOK, gin.H{"ok": true}) +} diff --git a/internal/handlers/third/threads_handler.go b/internal/handlers/third/threads_handler.go new file mode 100644 index 00000000..7f011ac2 --- /dev/null +++ b/internal/handlers/third/threads_handler.go @@ -0,0 +1,71 @@ +package third + +import ( + "bytes" + "io" + "net/http" + "strings" + + "agent-desk/internal/pkg/enums" + "agent-desk/internal/services" + + "github.com/gin-gonic/gin" +) + +// ThreadsGetWebhook handles Meta webhook verification (hub.challenge). +func ThreadsGetWebhook(ctx *gin.Context) { + mode := strings.TrimSpace(ctx.Query("hub.mode")) + token := strings.TrimSpace(ctx.Query("hub.verify_token")) + challenge := strings.TrimSpace(ctx.Query("hub.challenge")) + + channelID := strings.TrimSpace(ctx.Param("channel_id")) + if channelID == "" { + channelID = strings.TrimSpace(ctx.Query("channel_id")) + } + + if mode == "subscribe" { + if channelID != "" { + channel := services.ChannelService.Take("channel_id = ? AND channel_type = ? AND status = ?", channelID, enums.ChannelTypeThreads, enums.StatusOk) + if channel != nil { + if cfg, err := services.ChannelService.ParseThreadsChannelConfig(channel.ConfigJSON); err == nil && cfg != nil { + if cfg.WebhookVerifyToken != "" && cfg.WebhookVerifyToken != token { + ctx.String(http.StatusForbidden, "Verification token mismatch") + return + } + } + } + } + + ctx.String(http.StatusOK, challenge) + return + } + + ctx.String(http.StatusBadRequest, "Invalid verification request") +} + +// ThreadsPostWebhook receives incoming webhook events from Meta Threads. +func ThreadsPostWebhook(ctx *gin.Context) { + channelID := strings.TrimSpace(ctx.Param("channel_id")) + if channelID == "" { + channelID = strings.TrimSpace(ctx.Query("channel_id")) + } + + sigHeader := ctx.GetHeader("X-Hub-Signature-256") + if sigHeader == "" { + sigHeader = ctx.GetHeader("X-Hub-Signature") + } + + bodyBytes, err := io.ReadAll(ctx.Request.Body) + if err != nil { + ctx.JSON(http.StatusBadRequest, gin.H{"ok": false, "error": "failed to read body"}) + return + } + ctx.Request.Body = io.NopCloser(bytes.NewBuffer(bodyBytes)) + + if err := services.ThreadsInboundService.HandleWebhook(ctx.Request.Context(), channelID, sigHeader, bodyBytes); err != nil { + ctx.JSON(http.StatusOK, gin.H{"ok": false, "error": err.Error()}) + return + } + + ctx.JSON(http.StatusOK, gin.H{"ok": true, "message": "EVENT_RECEIVED"}) +} diff --git a/internal/handlers/third/viber_handler.go b/internal/handlers/third/viber_handler.go new file mode 100644 index 00000000..850c9dd5 --- /dev/null +++ b/internal/handlers/third/viber_handler.go @@ -0,0 +1,46 @@ +package third + +import ( + "bytes" + "io" + "net/http" + "strings" + + "agent-desk/internal/services" + + "github.com/gin-gonic/gin" +) + +// ViberPostWebhook receives incoming callbacks from Viber. +// +// For a conversation_started callback with a welcome message configured, +// the welcome message JSON is written to the response body as required +// by the Viber API. +func ViberPostWebhook(ctx *gin.Context) { + channelID := strings.TrimSpace(ctx.Param("channel_id")) + if channelID == "" { + channelID = strings.TrimSpace(ctx.Query("channel_id")) + } + + signature := ctx.GetHeader("X-Viber-Content-Signature") + + bodyBytes, err := io.ReadAll(ctx.Request.Body) + if err != nil { + ctx.JSON(http.StatusBadRequest, gin.H{"ok": false, "error": "failed to read body"}) + return + } + ctx.Request.Body = io.NopCloser(bytes.NewBuffer(bodyBytes)) + + responseBody, err := services.ViberInboundService.HandleWebhook(ctx.Request.Context(), channelID, signature, bodyBytes) + if err != nil { + ctx.JSON(http.StatusOK, gin.H{"ok": false, "error": err.Error()}) + return + } + + if responseBody != "" { + ctx.Data(http.StatusOK, "application/json", []byte(responseBody)) + return + } + + ctx.JSON(http.StatusOK, gin.H{"ok": true}) +} diff --git a/internal/line/client.go b/internal/line/client.go new file mode 100644 index 00000000..7c7bddde --- /dev/null +++ b/internal/line/client.go @@ -0,0 +1,116 @@ +package line + +import ( + "bytes" + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/base64" + "encoding/json" + "fmt" + "io" + "net/http" + "strings" + "time" +) + +const defaultBaseURL = "https://api.line.me" + +type Client struct { + channelAccessToken string + baseURL string + httpClient *http.Client +} + +func NewClient(channelAccessToken string) *Client { + return &Client{ + channelAccessToken: strings.TrimSpace(channelAccessToken), + baseURL: defaultBaseURL, + httpClient: &http.Client{Timeout: 15 * time.Second}, + } +} + +func (c *Client) SetBaseURL(url string) { + if strings.TrimSpace(url) != "" { + c.baseURL = strings.TrimRight(strings.TrimSpace(url), "/") + } +} + +// VerifyWebhookSignature validates the x-line-signature header value. +// The signature is HMAC-SHA256 of the raw body keyed by the channel secret, +// encoded as base64. +func VerifyWebhookSignature(channelSecret string, signature string, payload []byte) bool { + secret := strings.TrimSpace(channelSecret) + sig := strings.TrimSpace(signature) + if secret == "" || sig == "" { + return false + } + mac := hmac.New(sha256.New, []byte(secret)) + mac.Write(payload) + expected := base64.StdEncoding.EncodeToString(mac.Sum(nil)) + return hmac.Equal([]byte(expected), []byte(sig)) +} + +// PushMessage sends a push message to a user via the LINE Messaging API. +func (c *Client) PushMessage(ctx context.Context, req PushMessageRequest) (*PushMessageResponse, error) { + if strings.TrimSpace(req.To) == "" { + return nil, fmt.Errorf("line recipient (to) is required") + } + if len(req.Messages) == 0 { + return nil, fmt.Errorf("line message list is required") + } + + var resp PushMessageResponse + if err := c.doRequest(ctx, "/v2/bot/message/push", req, &resp); err != nil { + return nil, err + } + return &resp, nil +} + +func (c *Client) doRequest(ctx context.Context, path string, payload any, result any) error { + if c.channelAccessToken == "" { + return fmt.Errorf("line channel access token is required") + } + + endpoint := c.baseURL + path + + var bodyReader io.Reader + if payload != nil { + bodyBytes, err := json.Marshal(payload) + if err != nil { + return fmt.Errorf("marshal line request failed: %w", err) + } + bodyReader = bytes.NewBuffer(bodyBytes) + } + + req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, bodyReader) + if err != nil { + return fmt.Errorf("create line request failed: %w", err) + } + req.Header.Set("Authorization", "Bearer "+c.channelAccessToken) + if payload != nil { + req.Header.Set("Content-Type", "application/json") + } + + res, err := c.httpClient.Do(req) + if err != nil { + return fmt.Errorf("line http request failed: %w", err) + } + defer res.Body.Close() + + bodyBytes, err := io.ReadAll(res.Body) + if err != nil { + return fmt.Errorf("read line response failed: %w", err) + } + + if res.StatusCode < 200 || res.StatusCode >= 300 { + return fmt.Errorf("line api error (%d): %s", res.StatusCode, string(bodyBytes)) + } + + if result != nil && len(bodyBytes) > 0 { + if err := json.Unmarshal(bodyBytes, result); err != nil { + return fmt.Errorf("unmarshal line response failed: %w (body: %s)", err, string(bodyBytes)) + } + } + return nil +} diff --git a/internal/line/client_test.go b/internal/line/client_test.go new file mode 100644 index 00000000..05f950a2 --- /dev/null +++ b/internal/line/client_test.go @@ -0,0 +1,67 @@ +package line + +import ( + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/base64" + "encoding/json" + "net/http" + "net/http/httptest" + "testing" +) + +func TestLinePushMessage(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/v2/bot/message/push" { + t.Errorf("expected path /v2/bot/message/push, got %s", r.URL.Path) + } + if r.Header.Get("Authorization") != "Bearer test_token" { + t.Errorf("expected Bearer test_token, got %s", r.Header.Get("Authorization")) + } + var req PushMessageRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + t.Errorf("decode request failed: %v", err) + } + if req.To != "U4af4980629" { + t.Errorf("expected to U4af4980629, got %s", req.To) + } + if len(req.Messages) != 1 || req.Messages[0].Text != "hello" { + t.Errorf("unexpected messages: %+v", req.Messages) + } + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + w.Write([]byte(`{"sentMessages":[{"id":"4612309"}]}`)) + })) + defer server.Close() + + client := NewClient("test_token") + client.SetBaseURL(server.URL) + + resp, err := client.PushMessage(context.Background(), PushMessageRequest{ + To: "U4af4980629", + Messages: []MessageObject{{Type: "text", Text: "hello"}}, + }) + if err != nil { + t.Fatalf("PushMessage failed: %v", err) + } + if len(resp.SentMessages) != 1 || resp.SentMessages[0].ID != "4612309" { + t.Errorf("unexpected response: %+v", resp) + } +} + +func TestLineVerifyWebhookSignature(t *testing.T) { + const secret = "8c570fa6dd201bb328f1c1eac23a96d8" + body := []byte(`{"destination":"U8e742f61d673b39c7fff3cecb7536ef0","events":[]}`) + + mac := hmac.New(sha256.New, []byte(secret)) + mac.Write(body) + valid := base64.StdEncoding.EncodeToString(mac.Sum(nil)) + + if !VerifyWebhookSignature(secret, valid, body) { + t.Errorf("expected valid signature to verify") + } + if VerifyWebhookSignature(secret, "bad-signature", body) { + t.Errorf("expected invalid signature to fail") + } +} diff --git a/internal/line/types.go b/internal/line/types.go new file mode 100644 index 00000000..c38ee0ce --- /dev/null +++ b/internal/line/types.go @@ -0,0 +1,51 @@ +package line + +// WebhookEvent is the top-level webhook payload sent by the LINE Platform. +type WebhookEvent struct { + Destination string `json:"destination,omitempty"` + Events []Event `json:"events,omitempty"` +} + +// Event is a single webhook event object. +type Event struct { + Type string `json:"type,omitempty"` // message | follow | unfollow | join | leave | postback ... + ReplyToken string `json:"replyToken,omitempty"` + Source *Source `json:"source,omitempty"` + Message *Msg `json:"message,omitempty"` + Timestamp int64 `json:"timestamp,omitempty"` +} + +// Source describes where the event came from. +type Source struct { + Type string `json:"type,omitempty"` // user | group | room + UserID string `json:"userId,omitempty"` +} + +// Msg is the message object carried by a message event. +type Msg struct { + ID string `json:"id,omitempty"` + Type string `json:"type,omitempty"` // text | image | video | audio | file | sticker ... + Text string `json:"text,omitempty"` +} + +// PushMessageRequest is the request body of the send push message endpoint. +type PushMessageRequest struct { + To string `json:"to"` + Messages []MessageObject `json:"messages"` +} + +// MessageObject is a message to be sent to a user. +type MessageObject struct { + Type string `json:"type"` // text + Text string `json:"text"` +} + +// PushMessageResponse is the response of the push message endpoint. +type PushMessageResponse struct { + SentMessages []SentMessage `json:"sentMessages,omitempty"` +} + +// SentMessage describes a message accepted by the LINE Platform. +type SentMessage struct { + ID string `json:"id,omitempty"` +} diff --git a/internal/pkg/dto/dto.go b/internal/pkg/dto/dto.go index 537fc198..e8d6d608 100644 --- a/internal/pkg/dto/dto.go +++ b/internal/pkg/dto/dto.go @@ -113,16 +113,16 @@ type SlackChannelConfig struct { } type XChannelConfig struct { - BearerToken string `json:"bearerToken,omitempty"` // X API v2 Bearer Token - APIKey string `json:"apiKey,omitempty"` // Consumer Key - APISecretKey string `json:"apiSecretKey,omitempty"` // Consumer Secret - AccessToken string `json:"accessToken,omitempty"` // Access Token - AccessTokenSecret string `json:"accessTokenSecret,omitempty"` // Access Token Secret - AccountID string `json:"accountId,omitempty"` // X Numeric User/Account ID - Username string `json:"username,omitempty"` // @handle - WebhookEnv string `json:"webhookEnv,omitempty"` // Webhook environment name - WebhookCRCSecret string `json:"webhookCRCSecret,omitempty"` // CRC response secret - WelcomeMessage string `json:"welcomeMessage,omitempty"` + BearerToken string `json:"bearerToken,omitempty"` // X API v2 Bearer Token + APIKey string `json:"apiKey,omitempty"` // Consumer Key + APISecretKey string `json:"apiSecretKey,omitempty"` // Consumer Secret + AccessToken string `json:"accessToken,omitempty"` // Access Token + AccessTokenSecret string `json:"accessTokenSecret,omitempty"` // Access Token Secret + AccountID string `json:"accountId,omitempty"` // X Numeric User/Account ID + Username string `json:"username,omitempty"` // @handle + WebhookEnv string `json:"webhookEnv,omitempty"` // Webhook environment name + WebhookCRCSecret string `json:"webhookCRCSecret,omitempty"` // CRC response secret + WelcomeMessage string `json:"welcomeMessage,omitempty"` } type TikTokChannelConfig struct { @@ -143,9 +143,18 @@ type LineChannelConfig struct { } type ViberChannelConfig struct { - AuthToken string `json:"authToken,omitempty"` // Viber Bot Authentication Token - BotName string `json:"botName,omitempty"` // Sender Name - AvatarURL string `json:"avatarUrl,omitempty"` // Sender Avatar URL - WebhookSecret string `json:"webhookSecret,omitempty"` // Secret string in webhook event + AuthToken string `json:"authToken,omitempty"` // Viber Bot Authentication Token + BotName string `json:"botName,omitempty"` // Sender Name + AvatarURL string `json:"avatarUrl,omitempty"` // Sender Avatar URL + WebhookSecret string `json:"webhookSecret,omitempty"` // Secret string in webhook event WelcomeMessage string `json:"welcomeMessage,omitempty"` } + +type ThreadsChannelConfig struct { + ThreadsUserID string `json:"threadsUserId,omitempty"` // Threads App-Scoped User ID of the business account + Username string `json:"username,omitempty"` // @username of the Threads account + AccessToken string `json:"accessToken,omitempty"` // Long-lived Threads User Access Token + WebhookVerifyToken string `json:"webhookVerifyToken,omitempty"` // Meta webhook verification token + AppSecret string `json:"appSecret,omitempty"` // Meta App Secret for X-Hub-Signature-256 + WelcomeMessage string `json:"welcomeMessage,omitempty"` +} diff --git a/internal/pkg/enums/external_identity.go b/internal/pkg/enums/external_identity.go index dcacd78d..e56a943b 100644 --- a/internal/pkg/enums/external_identity.go +++ b/internal/pkg/enums/external_identity.go @@ -22,6 +22,7 @@ const ( ExternalSourceTikTok ExternalSource = "tiktok" // TikTok Direct Messages ExternalSourceLine ExternalSource = "line" // LINE Official Account ExternalSourceViber ExternalSource = "viber" // Viber Business Bot + ExternalSourceThreads ExternalSource = "threads" // Meta Threads ) var externalSourceLabelMap = map[ExternalSource]string{ @@ -41,6 +42,7 @@ var externalSourceLabelMap = map[ExternalSource]string{ ExternalSourceTikTok: "TikTok", ExternalSourceLine: "LINE", ExternalSourceViber: "Viber", + ExternalSourceThreads: "Threads", } func GetExternalSourceLabel(v ExternalSource) string { diff --git a/internal/pkg/enums/wxwork_kf.go b/internal/pkg/enums/wxwork_kf.go index 6f8841d2..6f392abc 100644 --- a/internal/pkg/enums/wxwork_kf.go +++ b/internal/pkg/enums/wxwork_kf.go @@ -33,6 +33,7 @@ const ( ChannelTypeTikTok = "tiktok" ChannelTypeLine = "line" ChannelTypeViber = "viber" + ChannelTypeThreads = "threads" ) type WxWorkKFMessageSendStatus string diff --git a/internal/services/channel_message_outbox_service.go b/internal/services/channel_message_outbox_service.go index 8e990195..992e63dc 100644 --- a/internal/services/channel_message_outbox_service.go +++ b/internal/services/channel_message_outbox_service.go @@ -756,6 +756,195 @@ func (s *channelMessageOutboxService) EnqueueTikTokMessage(conversation *models. return nil } +func (s *channelMessageOutboxService) EnqueueLineMessage(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.ChannelTypeLine { + return nil + } + if message.SenderType != enums.IMSenderTypeAgent && message.SenderType != enums.IMSenderTypeAI { + return nil + } + if message.MessageType != enums.IMMessageTypeText && message.MessageType != enums.IMMessageTypeHTML && message.MessageType != enums.IMMessageTypeImage && message.MessageType != enums.IMMessageTypeAttachment { + return nil + } + if existing := s.GetByMessageID(enums.ChannelTypeLine, message.ID); existing != nil { + return nil + } + + payload, err := json.Marshal(map[string]any{ + "conversationId": conversation.ID, + "messageId": message.ID, + "messageType": message.MessageType, + "content": strings.TrimSpace(message.Content), + "payload": strings.TrimSpace(message.Payload), + "senderId": message.SenderID, + }) + if err != nil { + return err + } + + now := time.Now() + err = s.Create(&models.ChannelMessageOutbox{ + ChannelType: enums.ChannelTypeLine, + 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, + }, + }) + if err != nil { + return err + } + + // Trigger async dispatch immediately + go func() { + defer func() { + if r := recover(); r != nil { + slog.Error("recovered from panic in line outbound dispatch", "error", r) + } + }() + LineOutboundService.DispatchPendingOutbox() + }() + + return nil +} + +func (s *channelMessageOutboxService) EnqueueViberMessage(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.ChannelTypeViber { + return nil + } + if message.SenderType != enums.IMSenderTypeAgent && message.SenderType != enums.IMSenderTypeAI { + return nil + } + if message.MessageType != enums.IMMessageTypeText && message.MessageType != enums.IMMessageTypeHTML && message.MessageType != enums.IMMessageTypeImage && message.MessageType != enums.IMMessageTypeAttachment { + return nil + } + if existing := s.GetByMessageID(enums.ChannelTypeViber, message.ID); existing != nil { + return nil + } + + payload, err := json.Marshal(map[string]any{ + "conversationId": conversation.ID, + "messageId": message.ID, + "messageType": message.MessageType, + "content": strings.TrimSpace(message.Content), + "payload": strings.TrimSpace(message.Payload), + "senderId": message.SenderID, + }) + if err != nil { + return err + } + + now := time.Now() + err = s.Create(&models.ChannelMessageOutbox{ + ChannelType: enums.ChannelTypeViber, + 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, + }, + }) + if err != nil { + return err + } + + // Trigger async dispatch immediately + go func() { + defer func() { + if r := recover(); r != nil { + slog.Error("recovered from panic in viber outbound dispatch", "error", r) + } + }() + ViberOutboundService.DispatchPendingOutbox() + }() + + return nil +} + +func (s *channelMessageOutboxService) EnqueueThreadsMessage(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.ChannelTypeThreads { + return nil + } + if message.SenderType != enums.IMSenderTypeAgent && message.SenderType != enums.IMSenderTypeAI { + return nil + } + if message.MessageType != enums.IMMessageTypeText && message.MessageType != enums.IMMessageTypeHTML && message.MessageType != enums.IMMessageTypeImage && message.MessageType != enums.IMMessageTypeAttachment { + return nil + } + if existing := s.GetByMessageID(enums.ChannelTypeThreads, message.ID); existing != nil { + return nil + } + + payload, err := json.Marshal(map[string]any{ + "conversationId": conversation.ID, + "messageId": message.ID, + "messageType": message.MessageType, + "content": strings.TrimSpace(message.Content), + "payload": strings.TrimSpace(message.Payload), + "senderId": message.SenderID, + }) + if err != nil { + return err + } + + now := time.Now() + err = s.Create(&models.ChannelMessageOutbox{ + ChannelType: enums.ChannelTypeThreads, + 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, + }, + }) + if err != nil { + return err + } + + // Trigger async dispatch immediately + go func() { + defer func() { + if r := recover(); r != nil { + slog.Error("recovered from panic in threads outbound dispatch", "error", r) + } + }() + ThreadsOutboundService.DispatchPendingOutbox() + }() + + return nil +} + func (s *channelMessageOutboxService) ListPending(channelType string, limit int) []models.ChannelMessageOutbox { if limit <= 0 { limit = 20 diff --git a/internal/services/channel_service.go b/internal/services/channel_service.go index 6ddebc7b..e3b01966 100644 --- a/internal/services/channel_service.go +++ b/internal/services/channel_service.go @@ -1,6 +1,7 @@ package services import ( + "agent-desk/internal/messenger" "agent-desk/internal/models" "agent-desk/internal/pkg/config" "agent-desk/internal/pkg/dto" @@ -11,7 +12,6 @@ import ( "agent-desk/internal/pkg/httpx" "agent-desk/internal/pkg/utils" "agent-desk/internal/repositories" - "agent-desk/internal/messenger" "agent-desk/internal/telegram" "agent-desk/internal/wxwork" "context" @@ -583,6 +583,54 @@ func (s *channelService) ParseTikTokChannelConfig(raw string) (*dto.TikTokChanne return cfg, nil } +func (s *channelService) ParseLineChannelConfig(raw string) (*dto.LineChannelConfig, error) { + raw = strings.TrimSpace(raw) + cfg := &dto.LineChannelConfig{} + if raw != "" { + if err := json.Unmarshal([]byte(raw), cfg); err != nil { + return nil, err + } + } + cfg.ChannelID = strings.TrimSpace(cfg.ChannelID) + cfg.ChannelSecret = strings.TrimSpace(cfg.ChannelSecret) + cfg.ChannelAccessToken = strings.TrimSpace(cfg.ChannelAccessToken) + cfg.WelcomeMessage = strings.TrimSpace(cfg.WelcomeMessage) + return cfg, nil +} + +func (s *channelService) ParseViberChannelConfig(raw string) (*dto.ViberChannelConfig, error) { + raw = strings.TrimSpace(raw) + cfg := &dto.ViberChannelConfig{} + if raw != "" { + if err := json.Unmarshal([]byte(raw), cfg); err != nil { + return nil, err + } + } + cfg.AuthToken = strings.TrimSpace(cfg.AuthToken) + cfg.BotName = strings.TrimSpace(cfg.BotName) + cfg.AvatarURL = strings.TrimSpace(cfg.AvatarURL) + cfg.WebhookSecret = strings.TrimSpace(cfg.WebhookSecret) + cfg.WelcomeMessage = strings.TrimSpace(cfg.WelcomeMessage) + return cfg, nil +} + +func (s *channelService) ParseThreadsChannelConfig(raw string) (*dto.ThreadsChannelConfig, error) { + raw = strings.TrimSpace(raw) + cfg := &dto.ThreadsChannelConfig{} + if raw != "" { + if err := json.Unmarshal([]byte(raw), cfg); err != nil { + return nil, err + } + } + cfg.ThreadsUserID = strings.TrimSpace(cfg.ThreadsUserID) + cfg.Username = strings.TrimSpace(cfg.Username) + cfg.AccessToken = strings.TrimSpace(cfg.AccessToken) + cfg.WebhookVerifyToken = strings.TrimSpace(cfg.WebhookVerifyToken) + cfg.AppSecret = strings.TrimSpace(cfg.AppSecret) + cfg.WelcomeMessage = strings.TrimSpace(cfg.WelcomeMessage) + return cfg, nil +} + func (s *channelService) GetUserTokenSecret(channel *models.Channel) string { if channel == nil { return "" @@ -802,7 +850,7 @@ func (s *channelService) GetEnabledChannel(ctx *gin.Context) *models.Channel { func (s *channelService) buildChannelModel(id int64, req request.CreateChannelRequest) (*models.Channel, error) { channelType := strings.TrimSpace(req.ChannelType) - if channelType != enums.ChannelTypeWeb && channelType != enums.ChannelTypeWechatMP && channelType != enums.ChannelTypeWxWorkKF && channelType != enums.ChannelTypeTelegram && channelType != enums.ChannelTypeZaloOA && channelType != enums.ChannelTypeEmail && channelType != enums.ChannelTypeDiscord && channelType != enums.ChannelTypeMessenger && channelType != enums.ChannelTypeInstagram && channelType != enums.ChannelTypeWhatsApp && channelType != enums.ChannelTypeSlack && channelType != enums.ChannelTypeX && channelType != enums.ChannelTypeTikTok { + if channelType != enums.ChannelTypeWeb && channelType != enums.ChannelTypeWechatMP && channelType != enums.ChannelTypeWxWorkKF && channelType != enums.ChannelTypeTelegram && channelType != enums.ChannelTypeZaloOA && channelType != enums.ChannelTypeEmail && channelType != enums.ChannelTypeDiscord && channelType != enums.ChannelTypeMessenger && channelType != enums.ChannelTypeInstagram && channelType != enums.ChannelTypeWhatsApp && channelType != enums.ChannelTypeSlack && channelType != enums.ChannelTypeX && channelType != enums.ChannelTypeTikTok && channelType != enums.ChannelTypeLine && channelType != enums.ChannelTypeViber && channelType != enums.ChannelTypeThreads { return nil, errorsx.InvalidParamI18n("error.e0250") } name := strings.TrimSpace(req.Name) @@ -1110,6 +1158,68 @@ func (s *channelService) buildChannelModel(id int64, req request.CreateChannelRe return nil, err } configJSON = string(configBytes) + case enums.ChannelTypeLine: + if channelID == "" { + channelID = strs.UUID() + } + if exists := s.Take("channel_id = ? AND status <> ? AND id <> ?", channelID, enums.StatusDeleted, id); exists != nil { + return nil, errorsx.InvalidParamI18n("error.e0248") + } + cfg, err := s.ParseLineChannelConfig(configJSON) + if err != nil { + return nil, errorsx.InvalidParam("invalid line configuration") + } + if cfg == nil || cfg.ChannelAccessToken == "" { + return nil, errorsx.InvalidParam("line channelAccessToken is required") + } + configBytes, err := json.Marshal(cfg) + if err != nil { + return nil, err + } + configJSON = string(configBytes) + case enums.ChannelTypeViber: + if channelID == "" { + channelID = strs.UUID() + } + if exists := s.Take("channel_id = ? AND status <> ? AND id <> ?", channelID, enums.StatusDeleted, id); exists != nil { + return nil, errorsx.InvalidParamI18n("error.e0248") + } + cfg, err := s.ParseViberChannelConfig(configJSON) + if err != nil { + return nil, errorsx.InvalidParam("invalid viber configuration") + } + if cfg == nil || cfg.AuthToken == "" { + return nil, errorsx.InvalidParam("viber authToken is required") + } + configBytes, err := json.Marshal(cfg) + if err != nil { + return nil, err + } + configJSON = string(configBytes) + case enums.ChannelTypeThreads: + if channelID == "" { + channelID = strs.UUID() + } + if exists := s.Take("channel_id = ? AND status <> ? AND id <> ?", channelID, enums.StatusDeleted, id); exists != nil { + return nil, errorsx.InvalidParamI18n("error.e0248") + } + cfg, err := s.ParseThreadsChannelConfig(configJSON) + if err != nil { + return nil, errorsx.InvalidParam("invalid threads configuration") + } + if cfg == nil || cfg.AccessToken == "" || cfg.ThreadsUserID == "" { + return nil, errorsx.InvalidParam("threads accessToken and threadsUserId are required") + } + if cfg.WebhookVerifyToken == "" { + if secret, err := generateUserTokenSecret(); err == nil { + cfg.WebhookVerifyToken = secret + } + } + configBytes, err := json.Marshal(cfg) + if err != nil { + return nil, err + } + configJSON = string(configBytes) } return &models.Channel{ diff --git a/internal/services/cronx/cron.go b/internal/services/cronx/cron.go index 1b3dabec..51636412 100644 --- a/internal/services/cronx/cron.go +++ b/internal/services/cronx/cron.go @@ -58,6 +58,18 @@ func Init() { if slackCount > 0 { slog.Info("slack outbox dispatched", "count", slackCount) } + lineCount := services.LineOutboundService.DispatchPendingOutbox() + if lineCount > 0 { + slog.Info("line outbox dispatched", "count", lineCount) + } + viberCount := services.ViberOutboundService.DispatchPendingOutbox() + if viberCount > 0 { + slog.Info("viber outbox dispatched", "count", viberCount) + } + threadsCount := services.ThreadsOutboundService.DispatchPendingOutbox() + if threadsCount > 0 { + slog.Info("threads outbox dispatched", "count", threadsCount) + } }) c.Start() diff --git a/internal/services/line_inbound_service.go b/internal/services/line_inbound_service.go new file mode 100644 index 00000000..5f2afe80 --- /dev/null +++ b/internal/services/line_inbound_service.go @@ -0,0 +1,135 @@ +package services + +import ( + "context" + "encoding/json" + "fmt" + "strings" + "time" + + "agent-desk/internal/line" + "agent-desk/internal/models" + "agent-desk/internal/pkg/dto" + "agent-desk/internal/pkg/enums" + "agent-desk/internal/pkg/errorsx" + "agent-desk/internal/pkg/openidentity" +) + +var LineInboundService = newLineInboundService() + +func newLineInboundService() *lineInboundService { + return &lineInboundService{} +} + +type lineInboundService struct{} + +// HandleWebhook processes an incoming webhook from the LINE Platform. +func (s *lineInboundService) HandleWebhook(ctx context.Context, channelID string, signature string, rawPayload []byte) error { + channelID = strings.TrimSpace(channelID) + var channel *models.Channel + if channelID != "" { + channel = ChannelService.Take("channel_id = ? AND channel_type = ? AND status = ?", channelID, enums.ChannelTypeLine, enums.StatusOk) + } + if channel == nil { + channel = ChannelService.Take("channel_type = ? AND status = ?", enums.ChannelTypeLine, enums.StatusOk) + } + if channel == nil { + return errorsx.InvalidParam("line channel not found or disabled") + } + + cfg, err := ChannelService.ParseLineChannelConfig(channel.ConfigJSON) + if err != nil || cfg == nil || cfg.ChannelSecret == "" { + return errorsx.InvalidParam("line channel config invalid") + } + + // LINE requires webhook signature verification on every event. + if !line.VerifyWebhookSignature(cfg.ChannelSecret, signature, rawPayload) { + return errorsx.UnauthorizedI18n("error.auth.invalidSignature") + } + + var event line.WebhookEvent + if err := json.Unmarshal(rawPayload, &event); err != nil { + return fmt.Errorf("unmarshal line webhook failed: %w", err) + } + + for i := range event.Events { + if err := s.processEvent(channel, cfg, &event.Events[i]); err != nil { + return err + } + } + return nil +} + +func (s *lineInboundService) processEvent(channel *models.Channel, cfg *dto.LineChannelConfig, event *line.Event) error { + if event.Source == nil { + return nil + } + // Only handle 1:1 user events for now. + if event.Source.Type != "user" || strings.TrimSpace(event.Source.UserID) == "" { + return nil + } + + // Send the configured welcome message when a user follows the account. + if event.Type == "follow" { + welcome := strings.TrimSpace(cfg.WelcomeMessage) + if welcome == "" { + return nil + } + client := line.NewClient(cfg.ChannelAccessToken) + ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + if _, err := client.PushMessage(ctx, line.PushMessageRequest{ + To: strings.TrimSpace(event.Source.UserID), + Messages: []line.MessageObject{{Type: "text", Text: welcome}}, + }); err != nil { + return fmt.Errorf("send line welcome message failed: %w", err) + } + return nil + } + + if event.Type != "message" || event.Message == nil { + return nil + } + if event.Message.Type != "text" { + return nil // Ignore non-text messages for now + } + text := strings.TrimSpace(event.Message.Text) + if text == "" { + return nil + } + + externalID := strings.TrimSpace(event.Source.UserID) + name := fmt.Sprintf("LINE User %s", externalID) + + externalUser := openidentity.ExternalUser{ + ExternalSource: enums.ExternalSourceLine, + ExternalID: externalID, + ExternalName: name, + } + + conversation, err := ConversationService.Create(externalUser, channel.ID, channel.AIAgentID) + if err != nil { + return fmt.Errorf("create line conversation failed: %w", err) + } + + clientMsgID := fmt.Sprintf("line_%s", event.Message.ID) + payloadMap := map[string]any{ + "line_message_id": event.Message.ID, + "line_user_id": externalID, + "line_reply_token": event.ReplyToken, + "line_event_type": event.Type, + } + payloadBytes, _ := json.Marshal(payloadMap) + + if _, err := MessageService.SendCustomerMessage( + conversation.ID, + clientMsgID, + enums.IMMessageTypeText, + text, + string(payloadBytes), + externalUser, + ); err != nil { + return fmt.Errorf("send customer message failed: %w", err) + } + return nil +} diff --git a/internal/services/line_inbound_service_test.go b/internal/services/line_inbound_service_test.go new file mode 100644 index 00000000..4e286f96 --- /dev/null +++ b/internal/services/line_inbound_service_test.go @@ -0,0 +1,128 @@ +package services + +import ( + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/base64" + "encoding/json" + "fmt" + "testing" + "time" + + "agent-desk/internal/models" + "agent-desk/internal/pkg/dto" + "agent-desk/internal/pkg/enums" + "agent-desk/internal/repositories" + + "github.com/mlogclub/simple/sqls" +) + +const lineTestChannelSecret = "line_channel_secret_123" + +func signLinePayload(t *testing.T, secret string, payload []byte) string { + t.Helper() + mac := hmac.New(sha256.New, []byte(secret)) + mac.Write(payload) + return base64.StdEncoding.EncodeToString(mac.Sum(nil)) +} + +func TestLineInboundAndOutbound(t *testing.T) { + db := setupTikTokTestDB(t) + + now := time.Now() + aiAgent := &models.AIAgent{ + Name: "LINE AI Agent", + Status: enums.StatusOk, + PublishedRevisionID: 1, + AuditFields: models.AuditFields{CreatedAt: now, UpdatedAt: now}, + } + if err := db.Create(aiAgent).Error; err != nil { + t.Fatalf("create ai agent: %v", err) + } + + lineConfig := dto.LineChannelConfig{ + ChannelID: "2001234567", + ChannelSecret: lineTestChannelSecret, + ChannelAccessToken: "test_line_access_token", + } + cfgBytes, _ := json.Marshal(lineConfig) + + channel := &models.Channel{ + ChannelType: enums.ChannelTypeLine, + ChannelID: "line_channel_uuid_1", + AIAgentID: aiAgent.ID, + AIAgentRolloutPercent: 100, + Name: "LINE Support Channel", + ConfigJSON: string(cfgBytes), + Status: enums.StatusOk, + AuditFields: models.AuditFields{CreatedAt: now, UpdatedAt: now}, + } + if err := db.Create(channel).Error; err != nil { + t.Fatalf("create line channel: %v", err) + } + + payload := []byte(fmt.Sprintf( + `{"destination":"Udest123","events":[{"type":"message","replyToken":"reply_token_1","source":{"type":"user","userId":"Uline_cust_555"},"message":{"id":"msg_9001","type":"text","text":"Hello from LINE"},"timestamp":1725260000}]}`, + )) + signature := signLinePayload(t, lineTestChannelSecret, payload) + + ctx := context.Background() + if err := LineInboundService.HandleWebhook(ctx, "", signature, payload); err != nil { + t.Fatalf("HandleWebhook failed: %v", err) + } + + // Invalid signature must be rejected. + if err := LineInboundService.HandleWebhook(ctx, "", "invalid-signature", payload); err == nil { + t.Fatalf("expected invalid signature to be rejected") + } + + // Verify customer identity + identity := repositories.CustomerIdentityRepository.FindOne(sqls.DB(), sqls.NewCnd(). + Eq("external_source", enums.ExternalSourceLine). + Eq("external_id", "Uline_cust_555")) + if identity == nil { + t.Fatalf("expected customer identity for Uline_cust_555") + } + + // Verify conversation + conv := repositories.ConversationRepository.FindOne(sqls.DB(), sqls.NewCnd(). + Eq("customer_id", identity.CustomerID). + Eq("channel_id", channel.ID)) + if conv == nil { + t.Fatalf("expected conversation to be created") + } + + // Verify message + msg := repositories.MessageRepository.FindOne(sqls.DB(), sqls.NewCnd(). + Eq("conversation_id", conv.ID). + Eq("sender_type", enums.IMSenderTypeCustomer)) + if msg == nil { + t.Fatalf("expected customer message to be created") + } + if msg.Content != "Hello from LINE" { + t.Fatalf("expected message content 'Hello from LINE', got %s", msg.Content) + } + + // Verify outbox enqueue on agent reply + operator := &dto.AuthPrincipal{UserID: 1, Username: "tester"} + if _, err := MessageService.SendAIMessage(conv.ID, aiAgent.ID, "ai_line_reply_1", enums.IMMessageTypeText, "Hi, how can we help?", "", operator); err != nil { + t.Fatalf("SendAIMessage failed: %v", err) + } + + // Look up the outbox row created for the agent reply. + replyMsg := repositories.MessageRepository.FindOne(sqls.DB(), sqls.NewCnd(). + Eq("conversation_id", conv.ID). + Eq("sender_type", enums.IMSenderTypeAI). + Desc("id")) + if replyMsg == nil { + t.Fatalf("expected agent reply message to be created") + } + outbox := ChannelMessageOutboxService.GetByMessageID(enums.ChannelTypeLine, replyMsg.ID) + if outbox == nil { + t.Fatalf("expected line outbox row for agent reply") + } + if outbox.ChannelType != enums.ChannelTypeLine { + t.Fatalf("expected outbox channel type 'line', got %s", outbox.ChannelType) + } +} diff --git a/internal/services/line_outbound_service.go b/internal/services/line_outbound_service.go new file mode 100644 index 00000000..f545265e --- /dev/null +++ b/internal/services/line_outbound_service.go @@ -0,0 +1,169 @@ +package services + +import ( + "context" + "log/slog" + "strings" + "time" + + "agent-desk/internal/line" + "agent-desk/internal/models" + "agent-desk/internal/pkg/enums" + "agent-desk/internal/repositories" + "agent-desk/internal/services/storage" + + "github.com/mlogclub/simple/sqls" +) + +const ( + lineOutboxBatchSize = 20 + lineOutboxMaxRetry = 5 +) + +var LineOutboundService = newLineOutboundService() + +func newLineOutboundService() *lineOutboundService { + return &lineOutboundService{} +} + +type lineOutboundService struct{} + +func (s *lineOutboundService) DispatchPendingOutbox() int { + return s.doDispatchPendingOutbox(lineOutboxBatchSize) +} + +func (s *lineOutboundService) doDispatchPendingOutbox(limit int) int { + if limit <= 0 { + limit = lineOutboxBatchSize + } + items := ChannelMessageOutboxService.ListPending(enums.ChannelTypeLine, limit) + if len(items) == 0 { + return 0 + } + + successCount := 0 + for i := range items { + if err := s.processOutbox(items[i].ID); err != nil { + slog.Warn("process line outbox failed", + "outbox_id", items[i].ID, + "error", err, + ) + continue + } + successCount++ + } + return successCount +} + +func (s *lineOutboundService) processOutbox(outboxID int64) error { + outbox := ChannelMessageOutboxService.Get(outboxID) + if outbox == nil { + return nil + } + if outbox.ChannelType != enums.ChannelTypeLine { + return nil + } + if outbox.SendStatus == string(enums.ChannelMessageOutboxStatusSent) { + return nil + } + if outbox.NextRetryAt != nil && outbox.NextRetryAt.After(time.Now()) { + return nil + } + + if err := ChannelMessageOutboxService.Updates(outbox.ID, map[string]any{ + "send_status": string(enums.ChannelMessageOutboxStatusSending), + "updated_at": time.Now(), + }); err != nil { + return err + } + + message := MessageService.Get(outbox.MessageID) + if message == nil { + return s.markOutboxFailed(outbox, "message not found") + } + conversation := ConversationService.Get(outbox.ConversationID) + if conversation == nil { + return s.markOutboxFailed(outbox, "conversation not found") + } + + channel := ChannelService.Get(conversation.ChannelID) + if channel == nil || channel.Status != enums.StatusOk { + return s.markOutboxFailed(outbox, "line channel not found or disabled") + } + cfg, err := ChannelService.ParseLineChannelConfig(channel.ConfigJSON) + if err != nil || cfg == nil || cfg.ChannelAccessToken == "" { + return s.markOutboxFailed(outbox, "line credentials (channel access token) not configured") + } + + // Resolve target LINE user ID (ExternalID) + var recipientID string + customerIdentity := repositories.CustomerIdentityRepository.FindOne(sqls.DB(), sqls.NewCnd(). + Eq("customer_id", conversation.CustomerID). + Eq("external_source", enums.ExternalSourceLine)) + if customerIdentity != nil { + recipientID = strings.TrimSpace(customerIdentity.ExternalID) + } + if recipientID == "" { + return s.markOutboxFailed(outbox, "unable to resolve recipient line user id") + } + + text := strings.TrimSpace(message.Content) + if message.MessageType == enums.IMMessageTypeImage || message.MessageType == enums.IMMessageTypeAttachment { + assetPayload, err := parseIMMessageAssetPayload(message.Payload) + if err == nil && assetPayload != nil { + assetPayload = hydrateIMMessageAssetPayload(assetPayload) + if assetPayload.Provider != "" && assetPayload.StorageKey != "" { + if provider, err := storage.NewProvider(assetPayload.Provider); err == nil { + fileURL := provider.GetSignedURL(assetPayload.StorageKey) + if fileURL != "" { + if text != "" { + text += "\n" + fileURL + } else { + text = fileURL + } + } + } + } + } + } + if text == "" { + return s.markOutboxFailed(outbox, "line message has no text or resolvable media url") + } + + client := line.NewClient(cfg.ChannelAccessToken) + ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + + if _, err := client.PushMessage(ctx, line.PushMessageRequest{ + To: recipientID, + Messages: []line.MessageObject{{Type: "text", Text: text}}, + }); err != nil { + return s.markOutboxFailed(outbox, err.Error()) + } + + return ChannelMessageOutboxService.Updates(outbox.ID, map[string]any{ + "send_status": string(enums.ChannelMessageOutboxStatusSent), + "sent_at": time.Now(), + "updated_at": time.Now(), + }) +} + +func (s *lineOutboundService) markOutboxFailed(outbox *models.ChannelMessageOutbox, errMsg string) error { + if outbox == nil { + return nil + } + retryCount := outbox.RetryCount + 1 + status := string(enums.ChannelMessageOutboxStatusFailed) + if retryCount >= lineOutboxMaxRetry { + status = string(enums.ChannelMessageOutboxStatusIgnored) + } + nextRetryAt := time.Now().Add(time.Duration(retryCount*30) * time.Second) + + return ChannelMessageOutboxService.Updates(outbox.ID, map[string]any{ + "send_status": status, + "retry_count": retryCount, + "next_retry_at": &nextRetryAt, + "last_error": errMsg, + "updated_at": time.Now(), + }) +} diff --git a/internal/services/message_service.go b/internal/services/message_service.go index b306f63c..85ec8d2f 100644 --- a/internal/services/message_service.go +++ b/internal/services/message_service.go @@ -635,6 +635,33 @@ func (s *messageService) sendValidatedMessage(conversation *models.Conversation, "error", enqueueErr, ) } + + // LINE 渠道消息入队,异步发送 + if enqueueErr := ChannelMessageOutboxService.EnqueueLineMessage(conversation, message); enqueueErr != nil { + slog.Error("enqueue line outbox failed", + "conversation_id", conversation.ID, + "message_id", message.ID, + "error", enqueueErr, + ) + } + + // Viber 渠道消息入队,异步发送 + if enqueueErr := ChannelMessageOutboxService.EnqueueViberMessage(conversation, message); enqueueErr != nil { + slog.Error("enqueue viber outbox failed", + "conversation_id", conversation.ID, + "message_id", message.ID, + "error", enqueueErr, + ) + } + + // Threads 渠道消息入队,异步发送 + if enqueueErr := ChannelMessageOutboxService.EnqueueThreadsMessage(conversation, message); enqueueErr != nil { + slog.Error("enqueue threads outbox failed", + "conversation_id", conversation.ID, + "message_id", message.ID, + "error", enqueueErr, + ) + } // 客户发送消息,触发AI回复 if senderType == enums.IMSenderTypeCustomer { if TriggerAIReplyAsyncHook != nil { diff --git a/internal/services/threads_inbound_service.go b/internal/services/threads_inbound_service.go new file mode 100644 index 00000000..3f4e3146 --- /dev/null +++ b/internal/services/threads_inbound_service.go @@ -0,0 +1,149 @@ +package services + +import ( + "context" + "encoding/json" + "fmt" + "strings" + + "agent-desk/internal/models" + "agent-desk/internal/pkg/enums" + "agent-desk/internal/pkg/errorsx" + "agent-desk/internal/pkg/openidentity" + "agent-desk/internal/threads" +) + +var ThreadsInboundService = newThreadsInboundService() + +func newThreadsInboundService() *threadsInboundService { + return &threadsInboundService{} +} + +type threadsInboundService struct{} + +// HandleWebhook processes an incoming webhook from Meta Threads. +func (s *threadsInboundService) HandleWebhook(ctx context.Context, channelID string, signature string, rawPayload []byte) error { + channelID = strings.TrimSpace(channelID) + var channel *models.Channel + if channelID != "" { + channel = ChannelService.Take("channel_id = ? AND channel_type = ? AND status = ?", channelID, enums.ChannelTypeThreads, enums.StatusOk) + } + if channel == nil { + channel = ChannelService.Take("channel_type = ? AND status = ?", enums.ChannelTypeThreads, enums.StatusOk) + } + if channel == nil { + return errorsx.InvalidParam("threads channel not found or disabled") + } + + cfg, err := ChannelService.ParseThreadsChannelConfig(channel.ConfigJSON) + if err != nil || cfg == nil || cfg.AccessToken == "" { + return errorsx.InvalidParam("threads channel config invalid") + } + + // Verify X-Hub-Signature-256 when the app secret is configured. The + // check is fail-closed: once a secret is set, webhooks without a valid + // signature are rejected. + if cfg.AppSecret != "" { + if !threads.VerifyWebhookSignature(cfg.AppSecret, signature, rawPayload) { + return errorsx.UnauthorizedI18n("error.auth.invalidSignature") + } + } + + var payload threads.WebhookPayload + if err := json.Unmarshal(rawPayload, &payload); err != nil { + return fmt.Errorf("unmarshal threads webhook failed: %w", err) + } + + for _, value := range collectThreadsReplyValues(&payload) { + if err := s.processReply(channel, value); err != nil { + return err + } + } + return nil +} + +// collectThreadsReplyValues extracts reply objects from both documented +// webhook envelope shapes. +func collectThreadsReplyValues(payload *threads.WebhookPayload) []*threads.WebhookValue { + var values []*threads.WebhookValue + + appendValue := func(field string, value *threads.WebhookValue) { + if field == "replies" && value != nil && strings.TrimSpace(value.Text) != "" { + values = append(values, value) + } + } + + if payload == nil { + return values + } + for i := range payload.Entry { + for j := range payload.Entry[i].Changes { + appendValue(payload.Entry[i].Changes[j].Field, payload.Entry[i].Changes[j].Value) + } + } + if payload.Values != nil { + appendValue(payload.Values.Field, payload.Values.Value) + } + return values +} + +func (s *threadsInboundService) processReply(channel *models.Channel, value *threads.WebhookValue) error { + replyMediaID := strings.TrimSpace(value.ID) + if replyMediaID == "" { + replyMediaID = strings.TrimSpace(value.MediaID) + } + if replyMediaID == "" { + return nil + } + + // Threads reply webhooks do not carry a stable user id, but the + // @username is stable, so use it as the external identity to keep one + // customer per person. Fall back to the reply media id when missing. + externalID := strings.TrimSpace(value.Username) + if externalID == "" { + externalID = replyMediaID + } + + name := strings.TrimSpace(value.Username) + if name == "" { + name = fmt.Sprintf("Threads User %s", replyMediaID) + } + + externalUser := openidentity.ExternalUser{ + ExternalSource: enums.ExternalSourceThreads, + ExternalID: externalID, + ExternalName: name, + } + + conversation, err := ConversationService.Create(externalUser, channel.ID, channel.AIAgentID) + if err != nil { + return fmt.Errorf("create threads conversation failed: %w", err) + } + + clientMsgID := fmt.Sprintf("threads_%s", replyMediaID) + payloadMap := map[string]any{ + "threads_media_id": replyMediaID, + "threads_username": strings.TrimSpace(value.Username), + "threads_media_type": strings.TrimSpace(value.MediaType), + "threads_permalink": strings.TrimSpace(value.Permalink), + } + if value.RepliedTo != nil && strings.TrimSpace(value.RepliedTo.ID) != "" { + payloadMap["threads_reply_to_id"] = strings.TrimSpace(value.RepliedTo.ID) + } + if value.RootPost != nil && strings.TrimSpace(value.RootPost.ID) != "" { + payloadMap["threads_root_post_id"] = strings.TrimSpace(value.RootPost.ID) + } + payloadBytes, _ := json.Marshal(payloadMap) + + if _, err := MessageService.SendCustomerMessage( + conversation.ID, + clientMsgID, + enums.IMMessageTypeText, + strings.TrimSpace(value.Text), + string(payloadBytes), + externalUser, + ); err != nil { + return fmt.Errorf("send customer message failed: %w", err) + } + return nil +} diff --git a/internal/services/threads_inbound_service_test.go b/internal/services/threads_inbound_service_test.go new file mode 100644 index 00000000..cbddf395 --- /dev/null +++ b/internal/services/threads_inbound_service_test.go @@ -0,0 +1,145 @@ +package services + +import ( + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "fmt" + "testing" + "time" + + "agent-desk/internal/models" + "agent-desk/internal/pkg/dto" + "agent-desk/internal/pkg/enums" + "agent-desk/internal/repositories" + + "github.com/mlogclub/simple/sqls" +) + +const threadsTestAppSecret = "threads_app_secret_123" + +func signThreadsPayload(t *testing.T, secret string, payload []byte) string { + t.Helper() + mac := hmac.New(sha256.New, []byte(secret)) + mac.Write(payload) + return "sha256=" + hex.EncodeToString(mac.Sum(nil)) +} + +func TestThreadsInboundAndOutbound(t *testing.T) { + db := setupTikTokTestDB(t) + + now := time.Now() + aiAgent := &models.AIAgent{ + Name: "Threads AI Agent", + Status: enums.StatusOk, + PublishedRevisionID: 1, + AuditFields: models.AuditFields{CreatedAt: now, UpdatedAt: now}, + } + if err := db.Create(aiAgent).Error; err != nil { + t.Fatalf("create ai agent: %v", err) + } + + threadsConfig := dto.ThreadsChannelConfig{ + ThreadsUserID: "threads_biz_42", + Username: "crove_desk", + AccessToken: "test_threads_access_token", + WebhookVerifyToken: "threads_verify_token_9", + AppSecret: threadsTestAppSecret, + } + cfgBytes, _ := json.Marshal(threadsConfig) + + channel := &models.Channel{ + ChannelType: enums.ChannelTypeThreads, + ChannelID: "threads_channel_uuid_1", + AIAgentID: aiAgent.ID, + AIAgentRolloutPercent: 100, + Name: "Threads Support Channel", + ConfigJSON: string(cfgBytes), + Status: enums.StatusOk, + AuditFields: models.AuditFields{CreatedAt: now, UpdatedAt: now}, + } + if err := db.Create(channel).Error; err != nil { + t.Fatalf("create threads channel: %v", err) + } + + // Standard Meta envelope with a replies field change. + payload := []byte(fmt.Sprintf( + `{"object":"threads","entry":[{"id":"threads_biz_42","time":1725260000,"changes":[{"field":"replies","value":{"id":"reply_9001","username":"threads_customer","text":"Hello from Threads","media_type":"TEXT_POST","permalink":"https://www.threads.com/@threads_customer/post/Pp","replied_to":{"id":"root_post_1"},"root_post":{"id":"root_post_1","owner_id":"threads_biz_42"},"shortcode":"Pp","timestamp":"2026-09-07T10:33:16+0000"}}]}]}`, + )) + signature := signThreadsPayload(t, threadsTestAppSecret, payload) + + ctx := context.Background() + if err := ThreadsInboundService.HandleWebhook(ctx, "", signature, payload); err != nil { + t.Fatalf("HandleWebhook failed: %v", err) + } + + // Invalid signature must be rejected. + if err := ThreadsInboundService.HandleWebhook(ctx, "", "sha256=deadbeef", payload); err == nil { + t.Fatalf("expected invalid signature to be rejected") + } + + // Verify customer identity - username is the stable external id, so + // both replies from the same person map to one customer. + identity := repositories.CustomerIdentityRepository.FindOne(sqls.DB(), sqls.NewCnd(). + Eq("external_source", enums.ExternalSourceThreads). + Eq("external_id", "threads_customer")) + if identity == nil { + t.Fatalf("expected customer identity for threads_customer") + } + + // Verify conversation + conv := repositories.ConversationRepository.FindOne(sqls.DB(), sqls.NewCnd(). + Eq("customer_id", identity.CustomerID). + Eq("channel_id", channel.ID)) + if conv == nil { + t.Fatalf("expected conversation to be created") + } + + // Verify message + msg := repositories.MessageRepository.FindOne(sqls.DB(), sqls.NewCnd(). + Eq("conversation_id", conv.ID). + Eq("sender_type", enums.IMSenderTypeCustomer)) + if msg == nil { + t.Fatalf("expected customer message to be created") + } + if msg.Content != "Hello from Threads" { + t.Fatalf("expected message content 'Hello from Threads', got %s", msg.Content) + } + + // Verify outbox enqueue on agent reply + operator := &dto.AuthPrincipal{UserID: 1, Username: "tester"} + if _, err := MessageService.SendAIMessage(conv.ID, aiAgent.ID, "ai_threads_reply_1", enums.IMMessageTypeText, "Hi, how can we help?", "", operator); err != nil { + t.Fatalf("SendAIMessage failed: %v", err) + } + + replyMsg := repositories.MessageRepository.FindOne(sqls.DB(), sqls.NewCnd(). + Eq("conversation_id", conv.ID). + Eq("sender_type", enums.IMSenderTypeAI). + Desc("id")) + if replyMsg == nil { + t.Fatalf("expected agent reply message to be created") + } + outbox := ChannelMessageOutboxService.GetByMessageID(enums.ChannelTypeThreads, replyMsg.ID) + if outbox == nil { + t.Fatalf("expected threads outbox row for agent reply") + } + if outbox.ChannelType != enums.ChannelTypeThreads { + t.Fatalf("expected outbox channel type 'threads', got %s", outbox.ChannelType) + } + + // Topic/values envelope shape must also be parsed. + topicPayload := []byte(`{"app_id":"123456","topic":"moderate","target_id":"78901","time":1723226877,"subscription_id":"234567","values":{"value":{"id":"reply_9002","username":"threads_customer","text":"Second reply","media_type":"TEXT_POST","permalink":"https://www.threads.com/@threads_customer/post/Pq","replied_to":{"id":"reply_9001"},"root_post":{"id":"root_post_1"}},"field":"replies"}}`) + topicSignature := signThreadsPayload(t, threadsTestAppSecret, topicPayload) + if err := ThreadsInboundService.HandleWebhook(ctx, "", topicSignature, topicPayload); err != nil { + t.Fatalf("topic envelope HandleWebhook failed: %v", err) + } + // Same username must reuse the same customer identity (no fragmentation). + identities := repositories.CustomerIdentityRepository.Find(sqls.DB(), sqls.NewCnd(). + Eq("external_source", enums.ExternalSourceThreads). + Eq("external_id", "threads_customer")) + if len(identities) != 1 { + t.Fatalf("expected exactly 1 customer identity for threads_customer, got %d", len(identities)) + } +} diff --git a/internal/services/threads_outbound_service.go b/internal/services/threads_outbound_service.go new file mode 100644 index 00000000..42a7431b --- /dev/null +++ b/internal/services/threads_outbound_service.go @@ -0,0 +1,176 @@ +package services + +import ( + "context" + "encoding/json" + "log/slog" + "strings" + "time" + + "agent-desk/internal/models" + "agent-desk/internal/pkg/enums" + "agent-desk/internal/repositories" + "agent-desk/internal/services/storage" + "agent-desk/internal/threads" + + "github.com/mlogclub/simple/sqls" +) + +const ( + threadsOutboxBatchSize = 20 + threadsOutboxMaxRetry = 5 +) + +var ThreadsOutboundService = newThreadsOutboundService() + +func newThreadsOutboundService() *threadsOutboundService { + return &threadsOutboundService{} +} + +type threadsOutboundService struct{} + +func (s *threadsOutboundService) DispatchPendingOutbox() int { + return s.doDispatchPendingOutbox(threadsOutboxBatchSize) +} + +func (s *threadsOutboundService) doDispatchPendingOutbox(limit int) int { + if limit <= 0 { + limit = threadsOutboxBatchSize + } + items := ChannelMessageOutboxService.ListPending(enums.ChannelTypeThreads, limit) + if len(items) == 0 { + return 0 + } + + successCount := 0 + for i := range items { + if err := s.processOutbox(items[i].ID); err != nil { + slog.Warn("process threads outbox failed", + "outbox_id", items[i].ID, + "error", err, + ) + continue + } + successCount++ + } + return successCount +} + +func (s *threadsOutboundService) processOutbox(outboxID int64) error { + outbox := ChannelMessageOutboxService.Get(outboxID) + if outbox == nil { + return nil + } + if outbox.ChannelType != enums.ChannelTypeThreads { + return nil + } + if outbox.SendStatus == string(enums.ChannelMessageOutboxStatusSent) { + return nil + } + if outbox.NextRetryAt != nil && outbox.NextRetryAt.After(time.Now()) { + return nil + } + + if err := ChannelMessageOutboxService.Updates(outbox.ID, map[string]any{ + "send_status": string(enums.ChannelMessageOutboxStatusSending), + "updated_at": time.Now(), + }); err != nil { + return err + } + + message := MessageService.Get(outbox.MessageID) + if message == nil { + return s.markOutboxFailed(outbox, "message not found") + } + conversation := ConversationService.Get(outbox.ConversationID) + if conversation == nil { + return s.markOutboxFailed(outbox, "conversation not found") + } + + channel := ChannelService.Get(conversation.ChannelID) + if channel == nil || channel.Status != enums.StatusOk { + return s.markOutboxFailed(outbox, "threads channel not found or disabled") + } + cfg, err := ChannelService.ParseThreadsChannelConfig(channel.ConfigJSON) + if err != nil || cfg == nil || cfg.AccessToken == "" || cfg.ThreadsUserID == "" { + return s.markOutboxFailed(outbox, "threads credentials (access token / threads user id) not configured") + } + + // Threads replies target a media object, not a user. Use the latest + // customer message in the conversation as the reply anchor. + var replyTargetID string + lastCustomerMsg := repositories.MessageRepository.FindOne(sqls.DB(), sqls.NewCnd(). + Eq("conversation_id", conversation.ID). + Eq("sender_type", enums.IMSenderTypeCustomer). + Desc("id")) + if lastCustomerMsg != nil && lastCustomerMsg.Payload != "" { + var payloadMap map[string]any + if err := json.Unmarshal([]byte(lastCustomerMsg.Payload), &payloadMap); err == nil { + if id, ok := payloadMap["threads_media_id"].(string); ok && id != "" { + replyTargetID = id + } else if id, ok := payloadMap["threads_reply_to_id"].(string); ok && id != "" { + replyTargetID = id + } + } + } + if replyTargetID == "" { + return s.markOutboxFailed(outbox, "unable to resolve threads reply target") + } + + text := strings.TrimSpace(message.Content) + if message.MessageType == enums.IMMessageTypeImage || message.MessageType == enums.IMMessageTypeAttachment { + assetPayload, err := parseIMMessageAssetPayload(message.Payload) + if err == nil && assetPayload != nil { + assetPayload = hydrateIMMessageAssetPayload(assetPayload) + if assetPayload.Provider != "" && assetPayload.StorageKey != "" { + if provider, err := storage.NewProvider(assetPayload.Provider); err == nil { + fileURL := provider.GetSignedURL(assetPayload.StorageKey) + if fileURL != "" { + if text != "" { + text += "\n" + fileURL + } else { + text = fileURL + } + } + } + } + } + } + if text == "" { + return s.markOutboxFailed(outbox, "threads message has no text or resolvable media url") + } + + client := threads.NewClient(cfg.AccessToken) + ctx, cancel := context.WithTimeout(context.Background(), 20*time.Second) + defer cancel() + + if _, err := client.PublishTextReply(ctx, cfg.ThreadsUserID, text, replyTargetID); err != nil { + return s.markOutboxFailed(outbox, err.Error()) + } + + return ChannelMessageOutboxService.Updates(outbox.ID, map[string]any{ + "send_status": string(enums.ChannelMessageOutboxStatusSent), + "sent_at": time.Now(), + "updated_at": time.Now(), + }) +} + +func (s *threadsOutboundService) markOutboxFailed(outbox *models.ChannelMessageOutbox, errMsg string) error { + if outbox == nil { + return nil + } + retryCount := outbox.RetryCount + 1 + status := string(enums.ChannelMessageOutboxStatusFailed) + if retryCount >= threadsOutboxMaxRetry { + status = string(enums.ChannelMessageOutboxStatusIgnored) + } + nextRetryAt := time.Now().Add(time.Duration(retryCount*30) * time.Second) + + return ChannelMessageOutboxService.Updates(outbox.ID, map[string]any{ + "send_status": status, + "retry_count": retryCount, + "next_retry_at": &nextRetryAt, + "last_error": errMsg, + "updated_at": time.Now(), + }) +} diff --git a/internal/services/viber_inbound_service.go b/internal/services/viber_inbound_service.go new file mode 100644 index 00000000..cfc321a1 --- /dev/null +++ b/internal/services/viber_inbound_service.go @@ -0,0 +1,137 @@ +package services + +import ( + "context" + "encoding/json" + "fmt" + "strings" + + "agent-desk/internal/models" + "agent-desk/internal/pkg/enums" + "agent-desk/internal/pkg/errorsx" + "agent-desk/internal/pkg/openidentity" + "agent-desk/internal/viber" +) + +var ViberInboundService = newViberInboundService() + +func newViberInboundService() *viberInboundService { + return &viberInboundService{} +} + +type viberInboundService struct{} + +// HandleWebhook processes an incoming callback from Viber. +// +// Viber expects the HTTP response body of a conversation_started callback +// to carry the welcome message, so this method returns an optional JSON +// response body alongside the error. +func (s *viberInboundService) HandleWebhook(ctx context.Context, channelID string, signature string, rawPayload []byte) (string, error) { + channelID = strings.TrimSpace(channelID) + var channel *models.Channel + if channelID != "" { + channel = ChannelService.Take("channel_id = ? AND channel_type = ? AND status = ?", channelID, enums.ChannelTypeViber, enums.StatusOk) + } + if channel == nil { + channel = ChannelService.Take("channel_type = ? AND status = ?", enums.ChannelTypeViber, enums.StatusOk) + } + if channel == nil { + return "", errorsx.InvalidParam("viber channel not found or disabled") + } + + cfg, err := ChannelService.ParseViberChannelConfig(channel.ConfigJSON) + if err != nil || cfg == nil || cfg.AuthToken == "" { + return "", errorsx.InvalidParam("viber channel config invalid") + } + + // Every Viber callback is signed with the authentication token. + if !viber.VerifyWebhookSignature(cfg.AuthToken, signature, rawPayload) { + return "", errorsx.UnauthorizedI18n("error.auth.invalidSignature") + } + + var callback viber.Callback + if err := json.Unmarshal(rawPayload, &callback); err != nil { + return "", fmt.Errorf("unmarshal viber callback failed: %w", err) + } + + switch callback.Event { + case "conversation_started": + if strings.TrimSpace(cfg.WelcomeMessage) == "" { + return "", nil + } + welcome := map[string]any{ + "sender": map[string]any{ + "name": strings.TrimSpace(cfg.BotName), + "avatar": strings.TrimSpace(cfg.AvatarURL), + }, + "type": "text", + "text": strings.TrimSpace(cfg.WelcomeMessage), + } + body, err := json.Marshal(welcome) + if err != nil { + return "", err + } + return string(body), nil + case "message": + if err := s.processMessage(channel, &callback); err != nil { + return "", err + } + return "", nil + default: + // subscribed / unsubscribed / delivered / seen / failed are ignored. + return "", nil + } +} + +func (s *viberInboundService) processMessage(channel *models.Channel, callback *viber.Callback) error { + if callback.Sender == nil || callback.Message == nil { + return nil + } + if strings.TrimSpace(callback.Message.Type) != "text" { + return nil // Ignore non-text messages for now + } + text := strings.TrimSpace(callback.Message.Text) + if text == "" { + return nil + } + externalID := strings.TrimSpace(callback.Sender.ID) + if externalID == "" { + return nil + } + + name := strings.TrimSpace(callback.Sender.Name) + if name == "" { + name = fmt.Sprintf("Viber User %s", externalID) + } + + externalUser := openidentity.ExternalUser{ + ExternalSource: enums.ExternalSourceViber, + ExternalID: externalID, + ExternalName: name, + } + + conversation, err := ConversationService.Create(externalUser, channel.ID, channel.AIAgentID) + if err != nil { + return fmt.Errorf("create viber conversation failed: %w", err) + } + + clientMsgID := fmt.Sprintf("viber_%d", callback.MessageToken) + payloadMap := map[string]any{ + "viber_message_token": callback.MessageToken, + "viber_user_id": externalID, + "viber_message_type": callback.Message.Type, + } + payloadBytes, _ := json.Marshal(payloadMap) + + if _, err := MessageService.SendCustomerMessage( + conversation.ID, + clientMsgID, + enums.IMMessageTypeText, + text, + string(payloadBytes), + externalUser, + ); err != nil { + return fmt.Errorf("send customer message failed: %w", err) + } + return nil +} diff --git a/internal/services/viber_inbound_service_test.go b/internal/services/viber_inbound_service_test.go new file mode 100644 index 00000000..a7219980 --- /dev/null +++ b/internal/services/viber_inbound_service_test.go @@ -0,0 +1,138 @@ +package services + +import ( + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "fmt" + "strings" + "testing" + "time" + + "agent-desk/internal/models" + "agent-desk/internal/pkg/dto" + "agent-desk/internal/pkg/enums" + "agent-desk/internal/repositories" + + "github.com/mlogclub/simple/sqls" +) + +const viberTestAuthToken = "viber_auth_token_123" + +func signViberPayload(t *testing.T, token string, payload []byte) string { + t.Helper() + mac := hmac.New(sha256.New, []byte(token)) + mac.Write(payload) + return hex.EncodeToString(mac.Sum(nil)) +} + +func TestViberInboundAndOutbound(t *testing.T) { + db := setupTikTokTestDB(t) + + now := time.Now() + aiAgent := &models.AIAgent{ + Name: "Viber AI Agent", + Status: enums.StatusOk, + PublishedRevisionID: 1, + AuditFields: models.AuditFields{CreatedAt: now, UpdatedAt: now}, + } + if err := db.Create(aiAgent).Error; err != nil { + t.Fatalf("create ai agent: %v", err) + } + + viberConfig := dto.ViberChannelConfig{ + AuthToken: viberTestAuthToken, + BotName: "Crove Support", + WelcomeMessage: "Welcome! How can we help?", + } + cfgBytes, _ := json.Marshal(viberConfig) + + channel := &models.Channel{ + ChannelType: enums.ChannelTypeViber, + ChannelID: "viber_channel_uuid_1", + AIAgentID: aiAgent.ID, + AIAgentRolloutPercent: 100, + Name: "Viber Support Channel", + ConfigJSON: string(cfgBytes), + Status: enums.StatusOk, + AuditFields: models.AuditFields{CreatedAt: now, UpdatedAt: now}, + } + if err := db.Create(channel).Error; err != nil { + t.Fatalf("create viber channel: %v", err) + } + + payload := []byte(fmt.Sprintf( + `{"event":"message","timestamp":1457764199822,"message_token":4911,"sender":{"id":"viber_cust_777","name":"Viber Customer"},"message":{"type":"text","text":"Hello from Viber"}}`, + )) + signature := signViberPayload(t, viberTestAuthToken, payload) + + ctx := context.Background() + if _, err := ViberInboundService.HandleWebhook(ctx, "", signature, payload); err != nil { + t.Fatalf("HandleWebhook failed: %v", err) + } + + // Invalid signature must be rejected. + if _, err := ViberInboundService.HandleWebhook(ctx, "", "deadbeef", payload); err == nil { + t.Fatalf("expected invalid signature to be rejected") + } + + // Verify customer identity + identity := repositories.CustomerIdentityRepository.FindOne(sqls.DB(), sqls.NewCnd(). + Eq("external_source", enums.ExternalSourceViber). + Eq("external_id", "viber_cust_777")) + if identity == nil { + t.Fatalf("expected customer identity for viber_cust_777") + } + + // Verify conversation + message + conv := repositories.ConversationRepository.FindOne(sqls.DB(), sqls.NewCnd(). + Eq("customer_id", identity.CustomerID). + Eq("channel_id", channel.ID)) + if conv == nil { + t.Fatalf("expected conversation to be created") + } + + msg := repositories.MessageRepository.FindOne(sqls.DB(), sqls.NewCnd(). + Eq("conversation_id", conv.ID). + Eq("sender_type", enums.IMSenderTypeCustomer)) + if msg == nil { + t.Fatalf("expected customer message to be created") + } + if msg.Content != "Hello from Viber" { + t.Fatalf("expected message content 'Hello from Viber', got %s", msg.Content) + } + + // conversation_started returns the welcome message JSON body. + convStartedPayload := []byte(`{"event":"conversation_started","timestamp":1457764199822,"message_token":4910,"user":{"id":"viber_cust_777","name":"Viber Customer"}}`) + welcomeSignature := signViberPayload(t, viberTestAuthToken, convStartedPayload) + respBody, err := ViberInboundService.HandleWebhook(ctx, "", welcomeSignature, convStartedPayload) + if err != nil { + t.Fatalf("conversation_started HandleWebhook failed: %v", err) + } + if !strings.Contains(respBody, "Welcome! How can we help?") { + t.Fatalf("expected welcome message in response body, got %s", respBody) + } + + // Verify outbox enqueue on agent reply + operator := &dto.AuthPrincipal{UserID: 1, Username: "tester"} + if _, err := MessageService.SendAIMessage(conv.ID, aiAgent.ID, "ai_viber_reply_1", enums.IMMessageTypeText, "Hi, how can we help?", "", operator); err != nil { + t.Fatalf("SendAIMessage failed: %v", err) + } + + replyMsg := repositories.MessageRepository.FindOne(sqls.DB(), sqls.NewCnd(). + Eq("conversation_id", conv.ID). + Eq("sender_type", enums.IMSenderTypeAI). + Desc("id")) + if replyMsg == nil { + t.Fatalf("expected agent reply message to be created") + } + outbox := ChannelMessageOutboxService.GetByMessageID(enums.ChannelTypeViber, replyMsg.ID) + if outbox == nil { + t.Fatalf("expected viber outbox row for agent reply") + } + if outbox.ChannelType != enums.ChannelTypeViber { + t.Fatalf("expected outbox channel type 'viber', got %s", outbox.ChannelType) + } +} diff --git a/internal/services/viber_outbound_service.go b/internal/services/viber_outbound_service.go new file mode 100644 index 00000000..eef03dcb --- /dev/null +++ b/internal/services/viber_outbound_service.go @@ -0,0 +1,173 @@ +package services + +import ( + "context" + "log/slog" + "strings" + "time" + + "agent-desk/internal/models" + "agent-desk/internal/pkg/enums" + "agent-desk/internal/repositories" + "agent-desk/internal/services/storage" + "agent-desk/internal/viber" + + "github.com/mlogclub/simple/sqls" +) + +const ( + viberOutboxBatchSize = 20 + viberOutboxMaxRetry = 5 +) + +var ViberOutboundService = newViberOutboundService() + +func newViberOutboundService() *viberOutboundService { + return &viberOutboundService{} +} + +type viberOutboundService struct{} + +func (s *viberOutboundService) DispatchPendingOutbox() int { + return s.doDispatchPendingOutbox(viberOutboxBatchSize) +} + +func (s *viberOutboundService) doDispatchPendingOutbox(limit int) int { + if limit <= 0 { + limit = viberOutboxBatchSize + } + items := ChannelMessageOutboxService.ListPending(enums.ChannelTypeViber, limit) + if len(items) == 0 { + return 0 + } + + successCount := 0 + for i := range items { + if err := s.processOutbox(items[i].ID); err != nil { + slog.Warn("process viber outbox failed", + "outbox_id", items[i].ID, + "error", err, + ) + continue + } + successCount++ + } + return successCount +} + +func (s *viberOutboundService) processOutbox(outboxID int64) error { + outbox := ChannelMessageOutboxService.Get(outboxID) + if outbox == nil { + return nil + } + if outbox.ChannelType != enums.ChannelTypeViber { + return nil + } + if outbox.SendStatus == string(enums.ChannelMessageOutboxStatusSent) { + return nil + } + if outbox.NextRetryAt != nil && outbox.NextRetryAt.After(time.Now()) { + return nil + } + + if err := ChannelMessageOutboxService.Updates(outbox.ID, map[string]any{ + "send_status": string(enums.ChannelMessageOutboxStatusSending), + "updated_at": time.Now(), + }); err != nil { + return err + } + + message := MessageService.Get(outbox.MessageID) + if message == nil { + return s.markOutboxFailed(outbox, "message not found") + } + conversation := ConversationService.Get(outbox.ConversationID) + if conversation == nil { + return s.markOutboxFailed(outbox, "conversation not found") + } + + channel := ChannelService.Get(conversation.ChannelID) + if channel == nil || channel.Status != enums.StatusOk { + return s.markOutboxFailed(outbox, "viber channel not found or disabled") + } + cfg, err := ChannelService.ParseViberChannelConfig(channel.ConfigJSON) + if err != nil || cfg == nil || cfg.AuthToken == "" { + return s.markOutboxFailed(outbox, "viber credentials (auth token) not configured") + } + + // Resolve target Viber user ID (ExternalID) + var recipientID string + customerIdentity := repositories.CustomerIdentityRepository.FindOne(sqls.DB(), sqls.NewCnd(). + Eq("customer_id", conversation.CustomerID). + Eq("external_source", enums.ExternalSourceViber)) + if customerIdentity != nil { + recipientID = strings.TrimSpace(customerIdentity.ExternalID) + } + if recipientID == "" { + return s.markOutboxFailed(outbox, "unable to resolve recipient viber user id") + } + + text := strings.TrimSpace(message.Content) + if message.MessageType == enums.IMMessageTypeImage || message.MessageType == enums.IMMessageTypeAttachment { + assetPayload, err := parseIMMessageAssetPayload(message.Payload) + if err == nil && assetPayload != nil { + assetPayload = hydrateIMMessageAssetPayload(assetPayload) + if assetPayload.Provider != "" && assetPayload.StorageKey != "" { + if provider, err := storage.NewProvider(assetPayload.Provider); err == nil { + fileURL := provider.GetSignedURL(assetPayload.StorageKey) + if fileURL != "" { + if text != "" { + text += "\n" + fileURL + } else { + text = fileURL + } + } + } + } + } + } + if text == "" { + return s.markOutboxFailed(outbox, "viber message has no text or resolvable media url") + } + + // Viber requires a sender name on send_message; fall back to the + // channel name when no display name is configured. + senderName := strings.TrimSpace(cfg.BotName) + if senderName == "" { + senderName = strings.TrimSpace(channel.Name) + } + + client := viber.NewClient(cfg.AuthToken) + ctx, cancel := context.WithTimeout(context.Background(), 15*time.Second) + defer cancel() + + if _, err := client.SendTextMessage(ctx, recipientID, cfg.BotName, cfg.AvatarURL, text); err != nil { + return s.markOutboxFailed(outbox, err.Error()) + } + + return ChannelMessageOutboxService.Updates(outbox.ID, map[string]any{ + "send_status": string(enums.ChannelMessageOutboxStatusSent), + "sent_at": time.Now(), + "updated_at": time.Now(), + }) +} + +func (s *viberOutboundService) markOutboxFailed(outbox *models.ChannelMessageOutbox, errMsg string) error { + if outbox == nil { + return nil + } + retryCount := outbox.RetryCount + 1 + status := string(enums.ChannelMessageOutboxStatusFailed) + if retryCount >= viberOutboxMaxRetry { + status = string(enums.ChannelMessageOutboxStatusIgnored) + } + nextRetryAt := time.Now().Add(time.Duration(retryCount*30) * time.Second) + + return ChannelMessageOutboxService.Updates(outbox.ID, map[string]any{ + "send_status": status, + "retry_count": retryCount, + "next_retry_at": &nextRetryAt, + "last_error": errMsg, + "updated_at": time.Now(), + }) +} diff --git a/internal/threads/client.go b/internal/threads/client.go new file mode 100644 index 00000000..a73d284d --- /dev/null +++ b/internal/threads/client.go @@ -0,0 +1,149 @@ +package threads + +import ( + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "fmt" + "io" + "net/http" + "net/url" + "strings" + "time" +) + +const defaultBaseURL = "https://graph.threads.net/v1.0" + +type Client struct { + accessToken string + baseURL string + httpClient *http.Client +} + +func NewClient(accessToken string) *Client { + return &Client{ + accessToken: strings.TrimSpace(accessToken), + baseURL: defaultBaseURL, + httpClient: &http.Client{Timeout: 20 * time.Second}, + } +} + +func (c *Client) SetBaseURL(url string) { + if strings.TrimSpace(url) != "" { + c.baseURL = strings.TrimRight(strings.TrimSpace(url), "/") + } +} + +// VerifyWebhookSignature validates the X-Hub-Signature-256 header value. +// The signature is HMAC-SHA256 of the raw body keyed by the app secret, +// sent as "sha256=". +func VerifyWebhookSignature(appSecret string, signature string, payload []byte) bool { + secret := strings.TrimSpace(appSecret) + sig := strings.TrimSpace(signature) + if secret == "" || sig == "" { + return false + } + if strings.HasPrefix(sig, "sha256=") { + sig = sig[len("sha256="):] + } + mac := hmac.New(sha256.New, []byte(secret)) + mac.Write(payload) + expected := hex.EncodeToString(mac.Sum(nil)) + return hmac.Equal([]byte(expected), []byte(sig)) +} + +// PublishTextReply publishes a text reply to an existing Threads media +// object. The Threads API requires a two-step flow: create a media +// container, then publish it. +func (c *Client) PublishTextReply(ctx context.Context, threadsUserID string, text string, replyToID string) (*ContainerResponse, error) { + userID := strings.TrimSpace(threadsUserID) + if userID == "" { + return nil, fmt.Errorf("threads user id is required") + } + if strings.TrimSpace(text) == "" { + return nil, fmt.Errorf("threads text is required") + } + + container, err := c.createTextContainer(ctx, userID, text, strings.TrimSpace(replyToID)) + if err != nil { + return nil, err + } + + published, err := c.publishContainer(ctx, userID, container.ID) + if err != nil { + return nil, err + } + return published, nil +} + +func (c *Client) createTextContainer(ctx context.Context, threadsUserID string, text string, replyToID string) (*ContainerResponse, error) { + params := url.Values{} + params.Set("media_type", "TEXT") + params.Set("text", text) + params.Set("access_token", c.accessToken) + if replyToID != "" { + params.Set("reply_to_id", replyToID) + } + + var resp ContainerResponse + if err := c.doRequest(ctx, fmt.Sprintf("/%s/threads", threadsUserID), params, &resp); err != nil { + return nil, err + } + if resp.ID == "" { + return nil, fmt.Errorf("threads container creation returned no id") + } + return &resp, nil +} + +func (c *Client) publishContainer(ctx context.Context, threadsUserID string, containerID string) (*ContainerResponse, error) { + params := url.Values{} + params.Set("creation_id", containerID) + params.Set("access_token", c.accessToken) + + var resp ContainerResponse + if err := c.doRequest(ctx, fmt.Sprintf("/%s/threads_publish", threadsUserID), params, &resp); err != nil { + return nil, err + } + if resp.ID == "" { + return nil, fmt.Errorf("threads publish returned no id") + } + return &resp, nil +} + +func (c *Client) doRequest(ctx context.Context, path string, form url.Values, result any) error { + if c.accessToken == "" { + return fmt.Errorf("threads access token is required") + } + + endpoint := c.baseURL + path + + req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, strings.NewReader(form.Encode())) + if err != nil { + return fmt.Errorf("create threads request failed: %w", err) + } + req.Header.Set("Content-Type", "application/x-www-form-urlencoded") + + res, err := c.httpClient.Do(req) + if err != nil { + return fmt.Errorf("threads http request failed: %w", err) + } + defer res.Body.Close() + + bodyBytes, err := io.ReadAll(res.Body) + if err != nil { + return fmt.Errorf("read threads response failed: %w", err) + } + + if res.StatusCode < 200 || res.StatusCode >= 300 { + return fmt.Errorf("threads api error (%d): %s", res.StatusCode, string(bodyBytes)) + } + + if result != nil && len(bodyBytes) > 0 { + if err := json.Unmarshal(bodyBytes, result); err != nil { + return fmt.Errorf("unmarshal threads response failed: %w (body: %s)", err, string(bodyBytes)) + } + } + return nil +} diff --git a/internal/threads/client_test.go b/internal/threads/client_test.go new file mode 100644 index 00000000..e44dc84b --- /dev/null +++ b/internal/threads/client_test.go @@ -0,0 +1,74 @@ +package threads + +import ( + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/hex" + "net/http" + "net/http/httptest" + "strings" + "testing" +) + +func TestThreadsPublishTextReply(t *testing.T) { + var publishedContainerID string + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if !strings.HasSuffix(r.URL.Path, "/threads") && !strings.HasSuffix(r.URL.Path, "/threads_publish") { + t.Errorf("unexpected path %s", r.URL.Path) + } + w.Header().Set("Content-Type", "application/json") + if strings.HasSuffix(r.URL.Path, "/threads") { + if err := r.ParseForm(); err != nil { + t.Errorf("parse form failed: %v", err) + } + if r.Form.Get("media_type") != "TEXT" { + t.Errorf("expected media_type TEXT, got %s", r.Form.Get("media_type")) + } + if r.Form.Get("text") != "hello" { + t.Errorf("expected text hello, got %s", r.Form.Get("text")) + } + if r.Form.Get("reply_to_id") != "8901234" { + t.Errorf("expected reply_to_id 8901234, got %s", r.Form.Get("reply_to_id")) + } + publishedContainerID = "container_1" + w.Write([]byte(`{"id":"container_1"}`)) + return + } + if err := r.ParseForm(); err != nil { + t.Errorf("parse form failed: %v", err) + } + if r.Form.Get("creation_id") != publishedContainerID { + t.Errorf("expected creation_id %s, got %s", publishedContainerID, r.Form.Get("creation_id")) + } + w.Write([]byte(`{"id":"published_1"}`)) + })) + defer server.Close() + + client := NewClient("test_token") + client.SetBaseURL(server.URL) + + resp, err := client.PublishTextReply(context.Background(), "999", "hello", "8901234") + if err != nil { + t.Fatalf("PublishTextReply failed: %v", err) + } + if resp.ID != "published_1" { + t.Errorf("expected published id published_1, got %s", resp.ID) + } +} + +func TestThreadsVerifyWebhookSignature(t *testing.T) { + const secret = "app-secret" + body := []byte(`{"object":"threads","entry":[]}`) + + mac := hmac.New(sha256.New, []byte(secret)) + mac.Write(body) + valid := "sha256=" + hex.EncodeToString(mac.Sum(nil)) + + if !VerifyWebhookSignature(secret, valid, body) { + t.Errorf("expected valid signature to verify") + } + if VerifyWebhookSignature(secret, "sha256=deadbeef", body) { + t.Errorf("expected invalid signature to fail") + } +} diff --git a/internal/threads/types.go b/internal/threads/types.go new file mode 100644 index 00000000..08a1c2a2 --- /dev/null +++ b/internal/threads/types.go @@ -0,0 +1,72 @@ +package threads + +// WebhookPayload is the incoming webhook payload pushed by Meta. +// +// Threads webhooks are delivered in the standard Meta envelope +// (object/entry/changes) or, per the Threads webhook documentation, +// in a topic/values envelope. Both shapes are supported. +type WebhookPayload struct { + Object string `json:"object,omitempty"` + Entry []Entry `json:"entry,omitempty"` + Topic string `json:"topic,omitempty"` + Values *ValuesWrapper `json:"values,omitempty"` +} + +// Entry is one entry of the standard Meta webhook envelope. +type Entry struct { + ID string `json:"id,omitempty"` + Time int64 `json:"time,omitempty"` + Changes []Change `json:"changes,omitempty"` +} + +// Change is one field change of a Meta webhook entry. +type Change struct { + Field string `json:"field,omitempty"` // replies | mentions | publish | delete + Value *WebhookValue `json:"value,omitempty"` +} + +// ValuesWrapper is the values envelope of the topic-style payload. +type ValuesWrapper struct { + Field string `json:"field,omitempty"` // replies | mentions | publish | delete + Value *WebhookValue `json:"value,omitempty"` +} + +// WebhookValue carries the reply/post object of a webhook event. +type WebhookValue struct { + Event string `json:"event,omitempty"` // published | ... + ID string `json:"id,omitempty"` + MediaID string `json:"media_id,omitempty"` + Text string `json:"text,omitempty"` + Username string `json:"username,omitempty"` + MediaType string `json:"media_type,omitempty"` // TEXT_POST | IMAGE | VIDEO ... + Permalink string `json:"permalink,omitempty"` + Shortcode string `json:"shortcode,omitempty"` + Timestamp string `json:"timestamp,omitempty"` + RepliedTo *PostRef `json:"replied_to,omitempty"` + RootPost *PostRef `json:"root_post,omitempty"` + OwnerID string `json:"owner_id,omitempty"` +} + +// PostRef references another Threads media object. +type PostRef struct { + ID string `json:"id,omitempty"` + OwnerID string `json:"owner_id,omitempty"` + Username string `json:"username,omitempty"` +} + +// ReplyRef is the media id a reply targets. +type ReplyRef struct { + ID string `json:"id"` +} + +// CreateContainerRequest publishes a TEXT container via form parameters. +type CreateContainerRequest struct { + ThreadsUserID string + Text string + ReplyToID string +} + +// ContainerResponse is the response of the media container creation endpoint. +type ContainerResponse struct { + ID string `json:"id,omitempty"` +} diff --git a/internal/viber/client.go b/internal/viber/client.go new file mode 100644 index 00000000..fcf3885f --- /dev/null +++ b/internal/viber/client.go @@ -0,0 +1,119 @@ +package viber + +import ( + "bytes" + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "fmt" + "io" + "net/http" + "strings" + "time" +) + +const defaultBaseURL = "https://chatapi.viber.com" + +type Client struct { + authToken string + baseURL string + httpClient *http.Client +} + +func NewClient(authToken string) *Client { + return &Client{ + authToken: strings.TrimSpace(authToken), + baseURL: defaultBaseURL, + httpClient: &http.Client{Timeout: 15 * time.Second}, + } +} + +func (c *Client) SetBaseURL(url string) { + if strings.TrimSpace(url) != "" { + c.baseURL = strings.TrimRight(strings.TrimSpace(url), "/") + } +} + +// VerifyWebhookSignature validates the X-Viber-Content-Signature header value. +// The signature is HMAC-SHA256 of the raw body keyed by the authentication +// token, encoded as lowercase hex. +func VerifyWebhookSignature(authToken string, signature string, payload []byte) bool { + token := strings.TrimSpace(authToken) + sig := strings.TrimSpace(signature) + if token == "" || sig == "" { + return false + } + mac := hmac.New(sha256.New, []byte(token)) + mac.Write(payload) + expected := hex.EncodeToString(mac.Sum(nil)) + return hmac.Equal([]byte(expected), []byte(sig)) +} + +// SendTextMessage sends a text message to a Viber user. +func (c *Client) SendTextMessage(ctx context.Context, receiverID string, senderName string, senderAvatar string, text string) (*SendResponse, error) { + if strings.TrimSpace(receiverID) == "" { + return nil, fmt.Errorf("viber receiver id is required") + } + if strings.TrimSpace(text) == "" { + return nil, fmt.Errorf("viber message text is required") + } + + req := SendTextRequest{ + Receiver: strings.TrimSpace(receiverID), + Type: "text", + Text: text, + } + if strings.TrimSpace(senderName) != "" || strings.TrimSpace(senderAvatar) != "" { + req.Sender = &SenderRef{ + Name: strings.TrimSpace(senderName), + Avatar: strings.TrimSpace(senderAvatar), + } + } + + var resp SendResponse + if err := c.doRequest(ctx, "/pa/send_message", req, &resp); err != nil { + return nil, err + } + if resp.Status != 0 { + return nil, fmt.Errorf("viber sendMessage failed (%d): %s", resp.Status, resp.StatusMessage) + } + return &resp, nil +} + +func (c *Client) doRequest(ctx context.Context, path string, payload any, result any) error { + if c.authToken == "" { + return fmt.Errorf("viber auth token is required") + } + + endpoint := c.baseURL + path + + bodyBytes, err := json.Marshal(payload) + if err != nil { + return fmt.Errorf("marshal viber request failed: %w", err) + } + + req, err := http.NewRequestWithContext(ctx, http.MethodPost, endpoint, bytes.NewBuffer(bodyBytes)) + if err != nil { + return fmt.Errorf("create viber request failed: %w", err) + } + req.Header.Set("X-Viber-Auth-Token", c.authToken) + req.Header.Set("Content-Type", "application/json") + + res, err := c.httpClient.Do(req) + if err != nil { + return fmt.Errorf("viber http request failed: %w", err) + } + defer res.Body.Close() + + respBytes, err := io.ReadAll(res.Body) + if err != nil { + return fmt.Errorf("read viber response failed: %w", err) + } + + if err := json.Unmarshal(respBytes, result); err != nil { + return fmt.Errorf("unmarshal viber response failed: %w (body: %s)", err, string(respBytes)) + } + return nil +} diff --git a/internal/viber/client_test.go b/internal/viber/client_test.go new file mode 100644 index 00000000..19311cc0 --- /dev/null +++ b/internal/viber/client_test.go @@ -0,0 +1,80 @@ +package viber + +import ( + "context" + "crypto/hmac" + "crypto/sha256" + "encoding/hex" + "encoding/json" + "net/http" + "net/http/httptest" + "testing" +) + +func TestViberSendTextMessage(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + if r.URL.Path != "/pa/send_message" { + t.Errorf("expected path /pa/send_message, got %s", r.URL.Path) + } + if r.Header.Get("X-Viber-Auth-Token") != "test_token" { + t.Errorf("expected X-Viber-Auth-Token test_token, got %s", r.Header.Get("X-Viber-Auth-Token")) + } + var req SendTextRequest + if err := json.NewDecoder(r.Body).Decode(&req); err != nil { + t.Errorf("decode request failed: %v", err) + } + if req.Receiver != "01234567890=" { + t.Errorf("expected receiver 01234567890=, got %s", req.Receiver) + } + if req.Type != "text" || req.Text != "hello" { + t.Errorf("unexpected message: %+v", req) + } + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + w.Write([]byte(`{"status":0,"status_message":"ok","message_token":4911}`)) + })) + defer server.Close() + + client := NewClient("test_token") + client.SetBaseURL(server.URL) + + resp, err := client.SendTextMessage(context.Background(), "01234567890=", "", "", "hello") + if err != nil { + t.Fatalf("SendTextMessage failed: %v", err) + } + if resp.Status != 0 || resp.MessageToken != 4911 { + t.Errorf("unexpected response: %+v", resp) + } +} + +func TestViberSendTextMessageError(t *testing.T) { + server := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) { + w.Header().Set("Content-Type", "application/json") + w.WriteHeader(http.StatusOK) + w.Write([]byte(`{"status":5,"status_message":"Not a Viber user"}`)) + })) + defer server.Close() + + client := NewClient("test_token") + client.SetBaseURL(server.URL) + + if _, err := client.SendTextMessage(context.Background(), "unknown", "", "", "hello"); err == nil { + t.Fatalf("expected error for non-zero status") + } +} + +func TestViberVerifyWebhookSignature(t *testing.T) { + const token = "4453b0dcd47c3ae3-5e6c9866b2b1c3f7-oxv2lbqvbolcgtbe" + body := []byte(`{"event":"message","timestamp":1457764199822,"message_token":4911}`) + + mac := hmac.New(sha256.New, []byte(token)) + mac.Write(body) + valid := hex.EncodeToString(mac.Sum(nil)) + + if !VerifyWebhookSignature(token, valid, body) { + t.Errorf("expected valid signature to verify") + } + if VerifyWebhookSignature(token, "deadbeef", body) { + t.Errorf("expected invalid signature to fail") + } +} diff --git a/internal/viber/types.go b/internal/viber/types.go new file mode 100644 index 00000000..d07bfb73 --- /dev/null +++ b/internal/viber/types.go @@ -0,0 +1,57 @@ +package viber + +// Callback is the incoming callback payload pushed by Viber to the webhook. +type Callback struct { + Event string `json:"event,omitempty"` // message | conversation_started | subscribed | unsubscribed | delivered | seen | failed + Timestamp int64 `json:"timestamp,omitempty"` + MessageToken int64 `json:"message_token,omitempty"` + Sender *UserRef `json:"sender,omitempty"` + User *UserRef `json:"user,omitempty"` + Message *Msg `json:"message,omitempty"` + Silent bool `json:"silent,omitempty"` + Declined *FailedInfo `json:"declined_reason,omitempty"` +} + +// UserRef identifies a Viber user. +type UserRef struct { + ID string `json:"id,omitempty"` + Name string `json:"name,omitempty"` + Avatar string `json:"avatar,omitempty"` +} + +// Msg is the message object carried by a message event. +type Msg struct { + Type string `json:"type,omitempty"` // text | picture | video | file | contact | location | url ... + Text string `json:"text,omitempty"` + Media string `json:"media,omitempty"` + Contact *struct { + Name string `json:"name,omitempty"` + PhoneNumber string `json:"phone_number,omitempty"` + } `json:"contact,omitempty"` +} + +// FailedInfo describes why a message delivery failed. +type FailedInfo struct { + Description string `json:"description,omitempty"` +} + +// SendTextRequest is the request body of the send_message API. +type SendTextRequest struct { + Receiver string `json:"receiver"` + Type string `json:"type"` + Text string `json:"text"` + Sender *SenderRef `json:"sender,omitempty"` +} + +// SenderRef describes the sender displayed to the Viber user. +type SenderRef struct { + Name string `json:"name,omitempty"` + Avatar string `json:"avatar,omitempty"` +} + +// SendResponse is the response of the send_message API. +type SendResponse struct { + Status int `json:"status"` + StatusMessage string `json:"status_message,omitempty"` + MessageToken int64 `json:"message_token,omitempty"` +} diff --git a/web/app/(dashboard)/dashboard/channels/_components/edit.tsx b/web/app/(dashboard)/dashboard/channels/_components/edit.tsx index af995cdf..009a22b4 100644 --- a/web/app/(dashboard)/dashboard/channels/_components/edit.tsx +++ b/web/app/(dashboard)/dashboard/channels/_components/edit.tsx @@ -148,6 +148,29 @@ type TikTokChannelConfig = { webhookVerifyToken?: string } +type LineChannelConfig = { + channelId?: string + channelSecret?: string + channelAccessToken?: string + welcomeMessage?: string +} + +type ViberChannelConfig = { + authToken?: string + botName?: string + avatarUrl?: string + webhookSecret?: string + welcomeMessage?: string +} + +type ThreadsChannelConfig = { + threadsUserId?: string + username?: string + accessToken?: string + webhookVerifyToken?: string + appSecret?: string +} + function getDefaultWebChannelConfig(t: Translate): Required { return { title: t("channel.defaultTitleWeb"), @@ -162,7 +185,7 @@ function getDefaultWebChannelConfig(t: Translate): Required { function createSchema(t: Translate) { return z .object({ - channelType: z.enum(["web", "wechat_mp", "wxwork_kf", "telegram", "zalo_oa", "email", "discord", "messenger", "instagram", "whatsapp", "slack", "x", "tiktok"], t("channel.typeRequired")), + channelType: z.enum(["web", "wechat_mp", "wxwork_kf", "telegram", "zalo_oa", "email", "discord", "messenger", "instagram", "whatsapp", "slack", "x", "tiktok", "line", "viber", "threads"], t("channel.typeRequired")), aiAgentId: z.string().trim().regex(/^\d+$/, t("channel.agentRequired")), aiAgentRolloutPercent: z.coerce.number().int().min(1).max(100), name: z.string().trim().min(1, t("channel.nameRequired")), @@ -212,6 +235,20 @@ function createSchema(t: Translate) { tiktokOpenId: z.string().trim(), tiktokUsername: z.string().trim(), tiktokWebhookVerifyToken: z.string().trim(), + lineChannelId: z.string().trim(), + lineChannelSecret: z.string().trim(), + lineChannelAccessToken: z.string().trim(), + lineWelcomeMessage: z.string().trim(), + viberAuthToken: z.string().trim(), + viberBotName: z.string().trim(), + viberAvatarUrl: z.string().trim(), + viberWelcomeMessage: z.string().trim(), + viberWebhookSecret: z.string().trim(), + threadsUserId: z.string().trim(), + threadsUsername: z.string().trim(), + threadsAccessToken: z.string().trim(), + threadsAppSecret: z.string().trim(), + threadsWebhookVerifyToken: z.string().trim(), emailAddress: z.string().trim(), senderName: z.string().trim(), emailProvider: z.string().trim(), @@ -257,11 +294,46 @@ function createSchema(t: Translate) { message: "Zalo OA Access Token is required", }) } + if (values.channelType === "line" && !values.lineChannelAccessToken.trim()) { + ctx.addIssue({ + code: "custom", + path: ["lineChannelAccessToken"], + message: "LINE Channel Access Token is required", + }) + } + if (values.channelType === "line" && !values.lineChannelSecret.trim()) { + ctx.addIssue({ + code: "custom", + path: ["lineChannelSecret"], + message: "LINE Channel Secret is required", + }) + } + if (values.channelType === "viber" && !values.viberAuthToken.trim()) { + ctx.addIssue({ + code: "custom", + path: ["viberAuthToken"], + message: "Viber Auth Token is required", + }) + } + if (values.channelType === "threads" && !values.threadsAccessToken.trim()) { + ctx.addIssue({ + code: "custom", + path: ["threadsAccessToken"], + message: "Threads Access Token is required", + }) + } + if (values.channelType === "threads" && !values.threadsUserId.trim()) { + ctx.addIssue({ + code: "custom", + path: ["threadsUserId"], + message: "Threads User ID is required", + }) + } }) } type EditForm = { - channelType: "web" | "wechat_mp" | "wxwork_kf" | "telegram" | "zalo_oa" | "email" | "discord" | "messenger" | "instagram" | "whatsapp" | "slack" | "x" | "tiktok" + channelType: "web" | "wechat_mp" | "wxwork_kf" | "telegram" | "zalo_oa" | "email" | "discord" | "messenger" | "instagram" | "whatsapp" | "slack" | "x" | "tiktok" | "line" | "viber" | "threads" aiAgentId: string aiAgentRolloutPercent: number name: string @@ -311,6 +383,20 @@ type EditForm = { tiktokOpenId: string tiktokUsername: string tiktokWebhookVerifyToken: string + lineChannelId: string + lineChannelSecret: string + lineChannelAccessToken: string + lineWelcomeMessage: string + viberAuthToken: string + viberBotName: string + viberAvatarUrl: string + viberWelcomeMessage: string + viberWebhookSecret: string + threadsUserId: string + threadsUsername: string + threadsAccessToken: string + threadsAppSecret: string + threadsWebhookVerifyToken: string emailAddress: string senderName: string emailProvider: string @@ -381,6 +467,20 @@ function createEmptyForm(t: Translate): EditForm { tiktokOpenId: "", tiktokUsername: "", tiktokWebhookVerifyToken: "", + lineChannelId: "", + lineChannelSecret: "", + lineChannelAccessToken: "", + lineWelcomeMessage: "", + viberAuthToken: "", + viberBotName: "", + viberAvatarUrl: "", + viberWelcomeMessage: "", + viberWebhookSecret: "", + threadsUserId: "", + threadsUsername: "", + threadsAccessToken: "", + threadsAppSecret: "", + threadsWebhookVerifyToken: "", emailAddress: "help@crove.com", senderName: "Crove Desk Support", emailProvider: "brevo", @@ -620,6 +720,53 @@ function parseTikTokChannelConfig(configJson: string): TikTokChannelConfig { } } +function parseLineChannelConfig(configJson: string): LineChannelConfig { + if (!configJson.trim()) return {} + try { + const parsed = JSON.parse(configJson) as LineChannelConfig + return { + channelId: parsed.channelId?.trim() || "", + channelSecret: parsed.channelSecret?.trim() || "", + channelAccessToken: parsed.channelAccessToken?.trim() || "", + welcomeMessage: parsed.welcomeMessage?.trim() || "", + } + } catch { + return {} + } +} + +function parseViberChannelConfig(configJson: string): ViberChannelConfig { + if (!configJson.trim()) return {} + try { + const parsed = JSON.parse(configJson) as ViberChannelConfig + return { + authToken: parsed.authToken?.trim() || "", + botName: parsed.botName?.trim() || "", + avatarUrl: parsed.avatarUrl?.trim() || "", + webhookSecret: parsed.webhookSecret?.trim() || "", + welcomeMessage: parsed.welcomeMessage?.trim() || "", + } + } catch { + return {} + } +} + +function parseThreadsChannelConfig(configJson: string): ThreadsChannelConfig { + if (!configJson.trim()) return {} + try { + const parsed = JSON.parse(configJson) as ThreadsChannelConfig + return { + threadsUserId: parsed.threadsUserId?.trim() || "", + username: parsed.username?.trim() || "", + accessToken: parsed.accessToken?.trim() || "", + webhookVerifyToken: parsed.webhookVerifyToken?.trim() || "", + appSecret: parsed.appSecret?.trim() || "", + } + } catch { + return {} + } +} + function buildForm(item: AdminChannel | null, t: Translate): EditForm { if (!item) { return createEmptyForm(t) @@ -635,6 +782,9 @@ function buildForm(item: AdminChannel | null, t: Translate): EditForm { const isSlack = item.channelType === "slack" const isX = item.channelType === "x" const isTikTok = item.channelType === "tiktok" + const isLine = item.channelType === "line" + const isViber = item.channelType === "viber" + const isThreads = item.channelType === "threads" const webConfig = parseWebChannelConfig(item.configJson, t) const wechatConfig = isWechatMP ? parseWechatMPChannelConfig(item.configJson, t) @@ -669,6 +819,15 @@ function buildForm(item: AdminChannel | null, t: Translate): EditForm { const tiktokConfig = isTikTok ? parseTikTokChannelConfig(item.configJson) : null + const lineConfig = isLine + ? parseLineChannelConfig(item.configJson) + : null + const viberConfig = isViber + ? parseViberChannelConfig(item.configJson) + : null + const threadsConfig = isThreads + ? parseThreadsChannelConfig(item.configJson) + : null return { channelType: item.channelType === "wxwork_kf" @@ -691,7 +850,13 @@ function buildForm(item: AdminChannel | null, t: Translate): EditForm { ? "x" : item.channelType === "tiktok" ? "tiktok" - : item.channelType === "email" + : item.channelType === "line" + ? "line" + : item.channelType === "viber" + ? "viber" + : item.channelType === "threads" + ? "threads" + : item.channelType === "email" ? "email" : item.channelType === "wechat_mp" ? "wechat_mp" @@ -745,6 +910,20 @@ function buildForm(item: AdminChannel | null, t: Translate): EditForm { tiktokOpenId: tiktokConfig?.openId ?? "", tiktokUsername: tiktokConfig?.username ?? "", tiktokWebhookVerifyToken: tiktokConfig?.webhookVerifyToken ?? "", + lineChannelId: lineConfig?.channelId ?? "", + lineChannelSecret: lineConfig?.channelSecret ?? "", + lineChannelAccessToken: lineConfig?.channelAccessToken ?? "", + lineWelcomeMessage: lineConfig?.welcomeMessage ?? "", + viberAuthToken: viberConfig?.authToken ?? "", + viberBotName: viberConfig?.botName ?? "", + viberAvatarUrl: viberConfig?.avatarUrl ?? "", + viberWelcomeMessage: viberConfig?.welcomeMessage ?? "", + viberWebhookSecret: viberConfig?.webhookSecret ?? "", + threadsUserId: threadsConfig?.threadsUserId ?? "", + threadsUsername: threadsConfig?.username ?? "", + threadsAccessToken: threadsConfig?.accessToken ?? "", + threadsAppSecret: threadsConfig?.appSecret ?? "", + threadsWebhookVerifyToken: threadsConfig?.webhookVerifyToken ?? "", emailAddress: emailConfig?.emailAddress || "help@crove.com", senderName: emailConfig?.senderName || "Crove Desk Support", emailProvider: emailConfig?.provider || "brevo", @@ -857,6 +1036,29 @@ function buildPayload(form: EditForm, status: number, t: Translate): CreateAdmin username: form.tiktokUsername.trim(), webhookVerifyToken: form.tiktokWebhookVerifyToken.trim(), }) + : channelType === "line" + ? JSON.stringify({ + channelId: form.lineChannelId.trim(), + channelSecret: form.lineChannelSecret.trim(), + channelAccessToken: form.lineChannelAccessToken.trim(), + welcomeMessage: form.lineWelcomeMessage.trim(), + }) + : channelType === "viber" + ? JSON.stringify({ + authToken: form.viberAuthToken.trim(), + botName: form.viberBotName.trim(), + avatarUrl: form.viberAvatarUrl.trim(), + welcomeMessage: form.viberWelcomeMessage.trim(), + webhookSecret: form.viberWebhookSecret.trim(), + }) + : channelType === "threads" + ? JSON.stringify({ + threadsUserId: form.threadsUserId.trim(), + username: form.threadsUsername.trim(), + accessToken: form.threadsAccessToken.trim(), + appSecret: form.threadsAppSecret.trim(), + webhookVerifyToken: form.threadsWebhookVerifyToken.trim(), + }) : channelType === "wechat_mp" ? JSON.stringify(webLikeConfig) : JSON.stringify({ @@ -1103,6 +1305,9 @@ function ChannelFormBody({ { value: "slack", label: t("channel.typeSlack") }, { value: "x", label: t("channel.typeX") }, { value: "tiktok", label: t("channel.typeTikTok") }, + { value: "line", label: t("channel.typeLine") }, + { value: "viber", label: t("channel.typeViber") }, + { value: "threads", label: t("channel.typeThreads") }, { value: "telegram", label: t("channel.typeTelegram") }, { value: "zalo_oa", label: t("channel.typeZaloOa") }, { value: "wechat_mp", label: t("channel.typeWechatMp") }, @@ -1962,6 +2167,210 @@ function ChannelFormBody({ ) : null} + {channelType === "line" ? ( +
+
+
{t("channel.lineConnectTitle")}
+
{t("channel.lineConnectDescription")}
+
+ {t("channel.inboundWebhookUrl")}: /api/third/line/webhook +
+
+ +
+ + {t("channel.lineChannelId")} + + + + + + + + {t("channel.lineChannelSecret")} + + + + + +
+ + + {t("channel.lineChannelAccessToken")} + + + + + + + + {t("channel.welcomeMessageLabel")} + + + + + +
+ ) : null} + + {channelType === "viber" ? ( +
+
+
{t("channel.viberConnectTitle")}
+
{t("channel.viberConnectDescription")}
+
+ {t("channel.inboundWebhookUrl")}: /api/third/viber/webhook +
+
+ +
+ + {t("channel.viberAuthToken")} + + + + + + + + {t("channel.viberBotName")} + + + + + +
+ + + {t("channel.viberAvatarUrl")} + + + + + + + + {t("channel.welcomeMessageLabel")} + + + + + +
+ ) : null} + + {channelType === "threads" ? ( +
+
+
{t("channel.threadsConnectTitle")}
+
{t("channel.threadsConnectDescription")}
+
+ {t("channel.inboundWebhookUrl")}: /api/third/threads/webhook +
+
+ +
+ + {t("channel.threadsUserId")} + + + + + + + + {t("channel.threadsUsername")} + + + + + +
+ + + {t("channel.threadsAccessToken")} + + + + + + + + {t("channel.threadsAppSecret")} + + + + + + + + {t("channel.threadsWebhookVerifyToken")} + + + + +
+ ) : null} + {channelType === "wxwork_kf" ? ( {t("channel.wxworkAccount")} diff --git a/web/app/(dashboard)/dashboard/channels/page.tsx b/web/app/(dashboard)/dashboard/channels/page.tsx index b02da777..0d42803a 100644 --- a/web/app/(dashboard)/dashboard/channels/page.tsx +++ b/web/app/(dashboard)/dashboard/channels/page.tsx @@ -8,10 +8,12 @@ import { InstagramIcon, MailIcon, MessageCircleIcon, + MessageCircleMoreIcon, MessagesSquareIcon, MessageSquareMoreIcon, PhoneIcon, SendIcon, + SmartphoneIcon, VideoIcon, } from "lucide-react" @@ -59,6 +61,15 @@ function getChannelTypeLabel(channelType: string, t: (key: string) => string) { if (channelType === "tiktok") { return t("channel.typeTikTok") } + if (channelType === "line") { + return t("channel.typeLine") + } + if (channelType === "viber") { + return t("channel.typeViber") + } + if (channelType === "threads") { + return t("channel.typeThreads") + } if (channelType === "wechat_mp") { return t("channel.typeWechatMp") } @@ -109,6 +120,15 @@ function ChannelIcon({ channelType }: { channelType: string }) { if (channelType === "tiktok") { return } + if (channelType === "line") { + return + } + if (channelType === "viber") { + return + } + if (channelType === "threads") { + return + } if (channelType === "wechat_mp") { return } @@ -141,6 +161,9 @@ export default function DashboardChannelsPage() { { value: "slack", label: t("channel.typeSlack") }, { value: "x", label: t("channel.typeX") }, { value: "tiktok", label: t("channel.typeTikTok") }, + { value: "line", label: t("channel.typeLine") }, + { value: "viber", label: t("channel.typeViber") }, + { value: "threads", label: t("channel.typeThreads") }, { value: "telegram", label: t("channel.typeTelegram") }, { value: "zalo_oa", label: t("channel.typeZaloOa") }, { value: "wechat_mp", label: t("channel.typeWechatMp") }, diff --git a/web/components/channel-icon.tsx b/web/components/channel-icon.tsx index 1a2d9652..2e504acb 100644 --- a/web/components/channel-icon.tsx +++ b/web/components/channel-icon.tsx @@ -1,10 +1,18 @@ import { + AtSignIcon, + Gamepad2Icon, GlobeIcon, + HashIcon, + InstagramIcon, MailIcon, MessageCircleIcon, + MessageCircleMoreIcon, MessagesSquareIcon, MessageSquareMoreIcon, + PhoneIcon, SendIcon, + SmartphoneIcon, + VideoIcon, } from "lucide-react" export type ChannelIconProps = { @@ -24,6 +32,22 @@ export function ChannelIcon({ channelType, className = "size-3.5" }: ChannelIcon return case "messenger": return + case "instagram": + return + case "whatsapp": + return + case "slack": + return + case "x": + return + case "tiktok": + return + case "line": + return + case "viber": + return + case "threads": + return case "wxwork_kf": return case "wechat_mp": diff --git a/web/lib/generated/enums.ts b/web/lib/generated/enums.ts index 3b796297..575930b4 100644 --- a/web/lib/generated/enums.ts +++ b/web/lib/generated/enums.ts @@ -90,6 +90,7 @@ export enum ExternalSource { TikTok = "tiktok", Line = "line", Viber = "viber", + Threads = "threads", } export const ExternalSourceLabels: Record = { [ExternalSource.Guest]: "访客", @@ -108,6 +109,7 @@ export const ExternalSourceLabels: Record = { [ExternalSource.TikTok]: "TikTok", [ExternalSource.Line]: "LINE", [ExternalSource.Viber]: "Viber", + [ExternalSource.Threads]: "Threads", } export enum Gender { diff --git a/web/messages/en-US.json b/web/messages/en-US.json index 8f6478ae..b9c5df9a 100644 --- a/web/messages/en-US.json +++ b/web/messages/en-US.json @@ -637,6 +637,9 @@ "typeSlack": "Slack Workspace", "typeX": "X (Twitter)", "typeTikTok": "TikTok Messaging", + "typeLine": "LINE Official Account", + "typeViber": "Viber Business Bot", + "typeThreads": "Meta Threads", "typeWechatMp": "WeChat Official Account", "typeWxworkKf": "WeCom Customer Service", "emailAddress": "Support Email Address", @@ -717,6 +720,25 @@ "tiktokAccessToken": "Business Access Token", "tiktokClientKey": "App Client Key", "tiktokClientSecret": "App Client Secret", + "lineConnectTitle": "LINE Official Account Connection", + "lineConnectDescription": "Create a Messaging API channel in the LINE Developers Console, then paste the credentials here. Register the webhook endpoint below in your LINE channel settings. Inbound user messages will route directly to your agent inbox and AI.", + "lineChannelId": "LINE Channel ID", + "lineChannelSecret": "Channel Secret", + "lineChannelAccessToken": "Channel Access Token", + "viberConnectTitle": "Viber Business Bot Connection", + "viberConnectDescription": "Create a Viber bot to obtain the Auth Token, then register the webhook endpoint below (https required). Inbound customer messages will be automatically synchronized with your agent workbench and AI agent.", + "viberAuthToken": "Viber Auth Token", + "viberBotName": "Sender Display Name", + "viberAvatarUrl": "Sender Avatar URL (Optional)", + "welcomeMessageLabel": "Welcome Message (Optional)", + "threadsConnectTitle": "Meta Threads Connection", + "threadsConnectDescription": "Configure a Meta app with the Threads API permissions (threads_basic, threads_manage_replies, threads_read_replies), subscribe the replies webhook field with the endpoint below, then paste the long-lived access token and Threads user ID here. Customer replies to your Threads posts will be ingested and answered.", + "threadsUserId": "Threads User ID", + "threadsUsername": "Threads @Username", + "threadsAccessToken": "Threads Access Token", + "threadsAppSecret": "Meta App Secret", + "threadsWebhookVerifyToken": "Webhook Verify Token (Auto-generated)", + "threadsVerifyTokenHint": "Generated after saving - paste into your Meta app webhook settings", "loadFailed": "Could not load channels.", "created": "Channel created: {name}", "updated": "Channel updated: {name}", diff --git a/web/messages/vi-VN.json b/web/messages/vi-VN.json index 6c4d1880..ee7a7375 100644 --- a/web/messages/vi-VN.json +++ b/web/messages/vi-VN.json @@ -644,6 +644,9 @@ "typeSlack": "Slack Workspace", "typeX": "X (Twitter)", "typeTikTok": "TikTok Direct Messaging", + "typeLine": "LINE Official Account", + "typeViber": "Viber Business Bot", + "typeThreads": "Meta Threads", "typeWechatMp": "WeChat Official Account", "typeWxworkKf": "WeCom Customer Service", "emailAddress": "Địa chỉ Email Hỗ trợ", @@ -724,6 +727,25 @@ "tiktokAccessToken": "Business Access Token", "tiktokClientKey": "App Client Key", "tiktokClientSecret": "App Client Secret", + "lineConnectTitle": "Kết nối LINE Official Account", + "lineConnectDescription": "Tạo kênh Messaging API trong LINE Developers Console, dán thông tin xác thực vào bên dưới và đăng ký endpoint Webhook bên dưới trong cài đặt kênh LINE. Tin nhắn khách hàng sẽ tự động đồng bộ với workbench và AI Agent.", + "lineChannelId": "LINE Channel ID", + "lineChannelSecret": "Channel Secret", + "lineChannelAccessToken": "Channel Access Token", + "viberConnectTitle": "Kết nối Viber Business Bot", + "viberConnectDescription": "Tạo Viber Bot để lấy Auth Token, sau đó đăng ký endpoint Webhook bên dưới (yêu cầu https). Tin nhắn khách hàng sẽ tự động đồng bộ với workbench và AI Agent.", + "viberAuthToken": "Viber Auth Token", + "viberBotName": "Tên Người gửi Hiển thị", + "viberAvatarUrl": "URL Ảnh đại diện Người gửi (Tùy chọn)", + "welcomeMessageLabel": "Tin nhắn Chào mừng (Tùy chọn)", + "threadsConnectTitle": "Kết nối Meta Threads", + "threadsConnectDescription": "Tạo Meta app với quyền Threads API (threads_basic, threads_manage_replies, threads_read_replies), đăng ký trường Webhook replies với endpoint bên dưới, sau đó dán Access Token dài hạn và Threads User ID. Phản hồi của khách hàng trên bài viết Threads của bạn sẽ được tiếp nhận và xử lý tự động.", + "threadsUserId": "Threads User ID", + "threadsUsername": "Threads @Username", + "threadsAccessToken": "Threads Access Token", + "threadsAppSecret": "Meta App Secret", + "threadsWebhookVerifyToken": "Webhook Verify Token (Tự động tạo)", + "threadsVerifyTokenHint": "Tạo sau khi lưu - dán vào cài đặt Webhook của Meta app", "loadFailed": "Could not load channels.", "created": "Channel created: {name}", "updated": "Channel updated: {name}", diff --git a/web/messages/zh-CN.json b/web/messages/zh-CN.json index 32712ea9..af20116b 100644 --- a/web/messages/zh-CN.json +++ b/web/messages/zh-CN.json @@ -637,6 +637,9 @@ "typeSlack": "Slack Workspace", "typeX": "X (Twitter)", "typeTikTok": "TikTok 企业私信", + "typeLine": "LINE 公众号", + "typeViber": "Viber 商业机器人", + "typeThreads": "Meta Threads", "typeWechatMp": "微信公众号", "typeWxworkKf": "企业微信客服", "emailAddress": "支持邮箱地址", @@ -717,6 +720,25 @@ "tiktokAccessToken": "Business Access Token", "tiktokClientKey": "App Client Key", "tiktokClientSecret": "App Client Secret", + "lineConnectTitle": "LINE 公众号接入", + "lineConnectDescription": "在 LINE Developers Console 创建 Messaging API 渠道,将凭据填入下方,并在 LINE 渠道设置中注册下方的 Webhook 端点。用户消息将自动同步到坐席工作台和 AI Agent。", + "lineChannelId": "LINE Channel ID", + "lineChannelSecret": "Channel Secret", + "lineChannelAccessToken": "Channel Access Token", + "viberConnectTitle": "Viber 商业机器人接入", + "viberConnectDescription": "创建 Viber Bot 获取 Auth Token,并在 Viber 后台注册下方 Webhook 端点(需 https)。客户消息将自动同步到坐席工作台和 AI Agent。", + "viberAuthToken": "Viber Auth Token", + "viberBotName": "发件人显示名称", + "viberAvatarUrl": "发件人头像 URL(可选)", + "welcomeMessageLabel": "欢迎消息(可选)", + "threadsConnectTitle": "Meta Threads 接入", + "threadsConnectDescription": "创建具备 Threads API 权限(threads_basic、threads_manage_replies、threads_read_replies)的 Meta 应用,使用下方端点订阅 replies Webhook 字段,然后填入长期 Access Token 和 Threads 用户 ID。客户对你 Threads 帖子的回复将被接入并自动处理。", + "threadsUserId": "Threads 用户 ID", + "threadsUsername": "Threads @账号", + "threadsAccessToken": "Threads Access Token", + "threadsAppSecret": "Meta App Secret", + "threadsWebhookVerifyToken": "Webhook 验证令牌(自动生成)", + "threadsVerifyTokenHint": "保存后生成 - 粘贴到 Meta 应用的 Webhook 设置中", "loadFailed": "加载接入渠道失败", "created": "已创建接入渠道:{name}", "updated": "已更新接入渠道:{name}",