diff --git a/internal/handler/api/v1/conversation_handler.go b/internal/handler/api/v1/conversation_handler.go index e9117fe7..e26d94db 100644 --- a/internal/handler/api/v1/conversation_handler.go +++ b/internal/handler/api/v1/conversation_handler.go @@ -6,6 +6,7 @@ import ( "github.com/gin-gonic/gin" + "github.com/gochat/gochat/internal/model" "github.com/gochat/gochat/internal/search" "github.com/gochat/gochat/internal/service" "github.com/gochat/gochat/pkg/pagination" @@ -52,7 +53,7 @@ func (h *ConversationHandler) List(c *gin.Context) { // Support status filter via query param status := c.Query("status") - var items []interface{} + var conversations []model.Conversation var total int64 if status != "" { @@ -61,7 +62,7 @@ func (h *ConversationHandler) List(c *gin.Context) { handleServiceError(c, svcErr) return } - items = toInterfaceSlice(result) + conversations = result total = count } else { result, count, svcErr := h.conversationSvc.ListByAccount(c.Request.Context(), accountID, p.Offset, p.PerPage) @@ -69,11 +70,11 @@ func (h *ConversationHandler) List(c *gin.Context) { handleServiceError(c, svcErr) return } - items = toInterfaceSlice(result) + conversations = result total = count } - response.OKWithMeta(c, items, p.Page, p.PerPage, total) + c.JSON(http.StatusOK, serializeConversationList(c.Request.Context(), h.conversationSvc.DB(), conversations, total)) } // Create creates a new conversation. @@ -97,7 +98,7 @@ func (h *ConversationHandler) Create(c *gin.Context) { handleServiceError(c, svcErr) return } - response.Created(c, conversation) + c.JSON(http.StatusOK, serializeConversation(c.Request.Context(), h.conversationSvc.DB(), conversation)) } // @Summary Get a single conversation @@ -126,12 +127,12 @@ func (h *ConversationHandler) Get(c *gin.Context) { return } - conversation, svcErr := h.conversationSvc.GetByAccountAndID(c.Request.Context(), accountID, conversationID) + conversation, svcErr := h.conversationSvc.GetByAccountAndDisplayIDOrID(c.Request.Context(), accountID, conversationID) if svcErr != nil { handleServiceError(c, svcErr) return } - response.OK(c, conversation) + c.JSON(http.StatusOK, serializeConversation(c.Request.Context(), h.conversationSvc.DB(), conversation)) } // Update updates a conversation (status, priority). @@ -155,12 +156,17 @@ func (h *ConversationHandler) Update(c *gin.Context) { return } - conversation, svcErr := h.conversationSvc.Update(c.Request.Context(), accountID, conversationID, req) + conversation, svcErr := h.conversationSvc.GetByAccountAndDisplayIDOrID(c.Request.Context(), accountID, conversationID) if svcErr != nil { handleServiceError(c, svcErr) return } - response.OK(c, conversation) + conversation, svcErr = h.conversationSvc.Update(c.Request.Context(), accountID, conversation.ID, req) + if svcErr != nil { + handleServiceError(c, svcErr) + return + } + c.JSON(http.StatusOK, serializeConversation(c.Request.Context(), h.conversationSvc.DB(), conversation)) } // Delete soft-deletes a conversation. @@ -177,7 +183,12 @@ func (h *ConversationHandler) Delete(c *gin.Context) { return } - if svcErr := h.conversationSvc.Delete(c.Request.Context(), accountID, conversationID); svcErr != nil { + conversation, ok := h.resolveConversationRoute(c, accountID, conversationID) + if !ok { + return + } + + if svcErr := h.conversationSvc.Delete(c.Request.Context(), accountID, conversation.ID); svcErr != nil { handleServiceError(c, svcErr) return } @@ -217,12 +228,16 @@ func (h *ConversationHandler) AssignAgent(c *gin.Context) { return } - conversation, svcErr := h.conversationSvc.AssignAgent(c.Request.Context(), accountID, conversationID, req.AssigneeID) + conversation, ok := h.resolveConversationRoute(c, accountID, conversationID) + if !ok { + return + } + conversation, svcErr := h.conversationSvc.AssignAgent(c.Request.Context(), accountID, conversation.ID, req.AssigneeID) if svcErr != nil { handleServiceError(c, svcErr) return } - response.OK(c, conversation) + c.JSON(http.StatusOK, serializeUserFromDB(c.Request.Context(), h.conversationSvc.DB(), req.AssigneeID, accountID)) } // @Summary Toggle conversation status @@ -261,12 +276,21 @@ func (h *ConversationHandler) ToggleStatus(c *gin.Context) { return } - conversation, svcErr := h.conversationSvc.ToggleStatus(c.Request.Context(), accountID, conversationID, req) + conversation, ok := h.resolveConversationRoute(c, accountID, conversationID) + if !ok { + return + } + conversation, svcErr := h.conversationSvc.ToggleStatus(c.Request.Context(), accountID, conversation.ID, req) if svcErr != nil { handleServiceError(c, svcErr) return } - response.OK(c, conversation) + c.JSON(http.StatusOK, gin.H{"meta": gin.H{}, "payload": gin.H{ + "success": true, + "conversation_id": conversationDisplayID(conversation), + "current_status": conversation.Status, + "snoozed_until": conversation.SnoozedUntil, + }}) } // @Summary Mute a conversation @@ -297,12 +321,16 @@ func (h *ConversationHandler) Mute(c *gin.Context) { return } - conversation, svcErr := h.conversationSvc.Mute(c.Request.Context(), accountID, conversationID) + conversation, ok := h.resolveConversationRoute(c, accountID, conversationID) + if !ok { + return + } + _, svcErr := h.conversationSvc.Mute(c.Request.Context(), accountID, conversation.ID) if svcErr != nil { handleServiceError(c, svcErr) return } - response.OK(c, conversation) + c.Status(http.StatusOK) } // @Summary Unmute a conversation @@ -333,12 +361,16 @@ func (h *ConversationHandler) Unmute(c *gin.Context) { return } - conversation, svcErr := h.conversationSvc.Unmute(c.Request.Context(), accountID, conversationID) + conversation, ok := h.resolveConversationRoute(c, accountID, conversationID) + if !ok { + return + } + _, svcErr := h.conversationSvc.Unmute(c.Request.Context(), accountID, conversation.ID) if svcErr != nil { handleServiceError(c, svcErr) return } - response.OK(c, conversation) + c.Status(http.StatusOK) } // UpdateLabels updates the labels on a conversation. @@ -362,12 +394,16 @@ func (h *ConversationHandler) UpdateLabels(c *gin.Context) { return } - conversation, svcErr := h.conversationSvc.UpdateLabels(c.Request.Context(), accountID, conversationID, req.Labels) + conversation, ok := h.resolveConversationRoute(c, accountID, conversationID) + if !ok { + return + } + conversation, svcErr := h.conversationSvc.UpdateLabels(c.Request.Context(), accountID, conversation.ID, req.Labels) if svcErr != nil { handleServiceError(c, svcErr) return } - response.OK(c, conversation) + c.JSON(http.StatusOK, serializeConversation(c.Request.Context(), h.conversationSvc.DB(), conversation)) } // @Summary Search conversations @@ -410,7 +446,7 @@ func (h *ConversationHandler) Search(c *gin.Context) { return } - response.OKWithMeta(c, toInterfaceSlice(conversations), p.Page, p.PerPage, total) + c.JSON(http.StatusOK, serializeConversationList(c.Request.Context(), h.conversationSvc.DB(), conversations, total)) } // Filter retrieves conversations matching advanced filter criteria. @@ -439,9 +475,14 @@ func (h *ConversationHandler) Filter(c *gin.Context) { return } - // 1:1 Chatwoot ConversationFinder response shape: - // {conversations: [...], count: {mine_count, assigned_count, unassigned_count, all_count}} - response.OKWithMeta(c, toInterfaceSlice(result.Conversations), p.Page, p.PerPage, result.Count.AllCount) + payload := serializeConversationList(c.Request.Context(), h.conversationSvc.DB(), result.Conversations, result.Count.AllCount) + payload.Data.Meta = chatwootConversationCounts{ + MineCount: result.Count.MineCount, + AssignedCount: result.Count.AssignedCount, + UnassignedCount: result.Count.UnassignedCount, + AllCount: result.Count.AllCount, + } + c.JSON(http.StatusOK, payload.Data) } // UpdatePriority updates the priority of a conversation. @@ -466,12 +507,16 @@ func (h *ConversationHandler) UpdatePriority(c *gin.Context) { return } - conversation, svcErr := h.conversationSvc.UpdatePriority(c.Request.Context(), accountID, conversationID, req.Priority) + conversation, ok := h.resolveConversationRoute(c, accountID, conversationID) + if !ok { + return + } + conversation, svcErr := h.conversationSvc.UpdatePriority(c.Request.Context(), accountID, conversation.ID, req.Priority) if svcErr != nil { handleServiceError(c, svcErr) return } - response.OK(c, conversation) + c.JSON(http.StatusOK, serializeConversation(c.Request.Context(), h.conversationSvc.DB(), conversation)) } // TogglePriority updates a conversation priority using Chatwoot's member action path. @@ -501,7 +546,11 @@ func (h *ConversationHandler) TogglePriority(c *gin.Context) { priority = *req.Priority } - if _, svcErr := h.conversationSvc.UpdatePriority(c.Request.Context(), accountID, conversationID, priority); svcErr != nil { + conversation, ok := h.resolveConversationRoute(c, accountID, conversationID) + if !ok { + return + } + if _, svcErr := h.conversationSvc.UpdatePriority(c.Request.Context(), accountID, conversation.ID, priority); svcErr != nil { handleServiceError(c, svcErr) return } @@ -511,21 +560,30 @@ func (h *ConversationHandler) TogglePriority(c *gin.Context) { // ListMessages lists messages in a conversation. // GET /api/v1/accounts/:account_id/conversations/:conversation_id/messages func (h *ConversationHandler) ListMessages(c *gin.Context) { + accountID, err := parseUintParam(c, "account_id") + if err != nil { + response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrBadRequest, "invalid account_id") + return + } conversationID, err := parseUintParam(c, "conversation_id") if err != nil { response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrBadRequest, "invalid conversation_id") return } + conversation, ok := h.resolveConversationRoute(c, accountID, conversationID) + if !ok { + return + } p := pagination.Parse(c) - messages, total, svcErr := h.conversationSvc.ListMessages(c.Request.Context(), conversationID, p.Offset, p.PerPage) + messages, _, svcErr := h.conversationSvc.ListMessages(c.Request.Context(), conversation.ID, p.Offset, p.PerPage) if svcErr != nil { handleServiceError(c, svcErr) return } - response.OKWithMeta(c, toInterfaceSlice(messages), p.Page, p.PerPage, total) + c.JSON(http.StatusOK, serializeMessageIndex(c.Request.Context(), h.conversationSvc.DB(), conversation, messages)) } // Meta retrieves aggregated conversation metadata (status counts, label counts) for an account. @@ -563,13 +621,17 @@ func (h *ConversationHandler) Unread(c *gin.Context) { return } - conversation, svcErr := h.conversationSvc.MarkUnread(c.Request.Context(), accountID, conversationID) + conversation, ok := h.resolveConversationRoute(c, accountID, conversationID) + if !ok { + return + } + conversation, svcErr := h.conversationSvc.MarkUnread(c.Request.Context(), accountID, conversation.ID) if svcErr != nil { handleServiceError(c, svcErr) return } - response.OK(c, conversation) + c.JSON(http.StatusOK, serializeConversation(c.Request.Context(), h.conversationSvc.DB(), conversation)) } // Transcript sends a conversation transcript via email. @@ -596,7 +658,11 @@ func (h *ConversationHandler) Transcript(c *gin.Context) { return } - if svcErr := h.conversationSvc.SendTranscript(c.Request.Context(), accountID, conversationID, req.Email); svcErr != nil { + conversation, ok := h.resolveConversationRoute(c, accountID, conversationID) + if !ok { + return + } + if svcErr := h.conversationSvc.SendTranscript(c.Request.Context(), accountID, conversation.ID, req.Email); svcErr != nil { handleServiceError(c, svcErr) return } @@ -628,13 +694,17 @@ func (h *ConversationHandler) UpdateCustomAttributes(c *gin.Context) { return } - conversation, svcErr := h.conversationSvc.UpdateCustomAttributes(c.Request.Context(), accountID, conversationID, req.CustomAttributes) + conversation, ok := h.resolveConversationRoute(c, accountID, conversationID) + if !ok { + return + } + conversation, svcErr := h.conversationSvc.UpdateCustomAttributes(c.Request.Context(), accountID, conversation.ID, req.CustomAttributes) if svcErr != nil { handleServiceError(c, svcErr) return } - response.OK(c, conversation) + c.JSON(http.StatusOK, serializeConversation(c.Request.Context(), h.conversationSvc.DB(), conversation)) } // ListAttachments returns paginated message attachments for a conversation. @@ -654,13 +724,17 @@ func (h *ConversationHandler) ListAttachments(c *gin.Context) { } p := pagination.Parse(c) - attachments, total, svcErr := h.messageSvc.ListAttachments(c.Request.Context(), accountID, conversationID, p.Offset, p.PerPage) + conversation, ok := h.resolveConversationRoute(c, accountID, conversationID) + if !ok { + return + } + attachments, _, svcErr := h.messageSvc.ListAttachments(c.Request.Context(), accountID, conversation.ID, p.Offset, p.PerPage) if svcErr != nil { handleServiceError(c, svcErr) return } - response.OKWithMeta(c, toInterfaceSlice(attachments), p.Page, p.PerPage, total) + c.JSON(http.StatusOK, gin.H{"payload": attachments}) } // ToggleTyping toggles the typing status for an agent in a conversation. @@ -687,13 +761,17 @@ func (h *ConversationHandler) ToggleTyping(c *gin.Context) { return } - svcErr := h.conversationSvc.ToggleTyping(c.Request.Context(), accountID, conversationID, req.TypingStatus) + conversation, ok := h.resolveConversationRoute(c, accountID, conversationID) + if !ok { + return + } + svcErr := h.conversationSvc.ToggleTyping(c.Request.Context(), accountID, conversation.ID, req.TypingStatus) if svcErr != nil { handleServiceError(c, svcErr) return } - response.OK(c, gin.H{"typing_status": req.TypingStatus}) + c.Status(http.StatusOK) } // UpdateLastSeen updates the agent's last seen timestamp for a conversation. @@ -712,13 +790,17 @@ func (h *ConversationHandler) UpdateLastSeen(c *gin.Context) { return } - svcErr := h.conversationSvc.UpdateLastSeen(c.Request.Context(), accountID, conversationID) + conversation, ok := h.resolveConversationRoute(c, accountID, conversationID) + if !ok { + return + } + svcErr := h.conversationSvc.UpdateLastSeen(c.Request.Context(), accountID, conversation.ID) if svcErr != nil { handleServiceError(c, svcErr) return } - response.OK(c, nil) + c.Status(http.StatusOK) } // @Summary Assign a team to a conversation @@ -761,13 +843,21 @@ func (h *ConversationHandler) AssignTeam(c *gin.Context) { return } - conversation, svcErr := h.conversationSvc.AssignTeam(c.Request.Context(), accountID, conversationID, req.AgentID, req.TeamID) + conversation, ok := h.resolveConversationRoute(c, accountID, conversationID) + if !ok { + return + } + conversation, svcErr := h.conversationSvc.AssignTeam(c.Request.Context(), accountID, conversation.ID, req.AgentID, req.TeamID) if svcErr != nil { handleServiceError(c, svcErr) return } - response.OK(c, conversation) + if req.AgentID != nil { + c.JSON(http.StatusOK, serializeUserFromDB(c.Request.Context(), h.conversationSvc.DB(), *req.AgentID, accountID)) + return + } + c.JSON(http.StatusOK, serializeConversation(c.Request.Context(), h.conversationSvc.DB(), conversation)) } // handleServiceError maps service-layer errors to appropriate HTTP responses. @@ -798,6 +888,15 @@ func toInterfaceSlice[T any](slice []T) []interface{} { return result } +func (h *ConversationHandler) resolveConversationRoute(c *gin.Context, accountID, routeID uint) (*model.Conversation, bool) { + conversation, svcErr := h.conversationSvc.GetByAccountAndDisplayIDOrID(c.Request.Context(), accountID, routeID) + if svcErr != nil { + handleServiceError(c, svcErr) + return nil, false + } + return conversation, true +} + // UnreadCounts returns unread conversation counts grouped by inbox, label, and team. // GET /api/v1/accounts/:account_id/conversations/unread_counts // Reference: Chatwoot app/controllers/api/v1/accounts/conversations/unread_counts_controller.rb diff --git a/internal/handler/api/v1/conversation_handler_crud_test.go b/internal/handler/api/v1/conversation_handler_crud_test.go index fce6fd17..1db44df0 100644 --- a/internal/handler/api/v1/conversation_handler_crud_test.go +++ b/internal/handler/api/v1/conversation_handler_crud_test.go @@ -194,18 +194,18 @@ func (s *ConversationCrudTestSuite) TestList_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) var resp struct { - Success bool `json:"success"` - Data interface{} `json:"data"` - Meta struct { - Page int `json:"page"` - PerPage int `json:"per_page"` - TotalCount int64 `json:"total_count"` - } `json:"meta"` + Data struct { + Meta struct { + AllCount int64 `json:"all_count"` + } `json:"meta"` + Payload []map[string]interface{} `json:"payload"` + } `json:"data"` } err := json.Unmarshal(w.Body.Bytes(), &resp) assert.NoError(s.T(), err) - assert.True(s.T(), resp.Success) - assert.Equal(s.T(), int64(1), resp.Meta.TotalCount) + assert.Equal(s.T(), int64(1), resp.Data.Meta.AllCount) + assert.Len(s.T(), resp.Data.Payload, 1) + assert.NotNil(s.T(), resp.Data.Payload[0]["meta"]) } func (s *ConversationCrudTestSuite) TestList_WithStatusFilter() { @@ -216,15 +216,15 @@ func (s *ConversationCrudTestSuite) TestList_WithStatusFilter() { assert.Equal(s.T(), http.StatusOK, w.Code) var resp struct { - Success bool `json:"success"` - Data interface{} `json:"data"` - Meta struct { - TotalCount int64 `json:"total_count"` - } `json:"meta"` + Data struct { + Meta struct { + AllCount int64 `json:"all_count"` + } `json:"meta"` + } `json:"data"` } err := json.Unmarshal(w.Body.Bytes(), &resp) assert.NoError(s.T(), err) - assert.True(s.T(), resp.Success) + assert.Equal(s.T(), int64(1), resp.Data.Meta.AllCount) } func (s *ConversationCrudTestSuite) TestList_InvalidAccountID() { @@ -247,16 +247,13 @@ func (s *ConversationCrudTestSuite) TestCreate_Success() { req.Header.Set("Content-Type", "application/json") s.router.ServeHTTP(w, req) - assert.Equal(s.T(), http.StatusCreated, w.Code) + assert.Equal(s.T(), http.StatusOK, w.Code) - var resp struct { - Success bool `json:"success"` - Data map[string]interface{} `json:"data"` - } + var resp map[string]interface{} err := json.Unmarshal(w.Body.Bytes(), &resp) assert.NoError(s.T(), err) - assert.True(s.T(), resp.Success) - assert.NotNil(s.T(), resp.Data["id"]) + assert.NotNil(s.T(), resp["id"]) + assert.NotNil(s.T(), resp["meta"]) } func (s *ConversationCrudTestSuite) TestCreate_InvalidAccountID() { @@ -304,14 +301,11 @@ func (s *ConversationCrudTestSuite) TestGet_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) - var resp struct { - Success bool `json:"success"` - Data map[string]interface{} `json:"data"` - } + var resp map[string]interface{} err := json.Unmarshal(w.Body.Bytes(), &resp) assert.NoError(s.T(), err) - assert.True(s.T(), resp.Success) - assert.Equal(s.T(), float64(s.testConv.ID), resp.Data["id"]) + assert.Equal(s.T(), float64(s.testConv.ID), resp["id"]) + assert.NotNil(s.T(), resp["messages"]) } func (s *ConversationCrudTestSuite) TestGet_InvalidAccountID() { @@ -352,14 +346,10 @@ func (s *ConversationCrudTestSuite) TestUpdate_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) - var resp struct { - Success bool `json:"success"` - Data map[string]interface{} `json:"data"` - } + var resp map[string]interface{} err := json.Unmarshal(w.Body.Bytes(), &resp) assert.NoError(s.T(), err) - assert.True(s.T(), resp.Success) - assert.Equal(s.T(), "resolved", resp.Data["status"]) + assert.Equal(s.T(), "resolved", resp["status"]) } func (s *ConversationCrudTestSuite) TestUpdate_InvalidAccountID() { @@ -466,13 +456,10 @@ func (s *ConversationCrudTestSuite) TestAssignAgent_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) - var resp struct { - Success bool `json:"success"` - Data map[string]interface{} `json:"data"` - } + var resp map[string]interface{} err := json.Unmarshal(w.Body.Bytes(), &resp) assert.NoError(s.T(), err) - assert.True(s.T(), resp.Success) + assert.Equal(s.T(), float64(user.ID), resp["id"]) } func (s *ConversationCrudTestSuite) TestAssignAgent_InvalidAccountID() { @@ -522,13 +509,13 @@ func (s *ConversationCrudTestSuite) TestToggleStatus_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) var resp struct { - Success bool `json:"success"` - Data map[string]interface{} `json:"data"` + Payload struct { + CurrentStatus string `json:"current_status"` + } `json:"payload"` } err := json.Unmarshal(w.Body.Bytes(), &resp) assert.NoError(s.T(), err) - assert.True(s.T(), resp.Success) - assert.Equal(s.T(), "resolved", resp.Data["status"]) + assert.Equal(s.T(), "resolved", resp.Payload.CurrentStatus) } func (s *ConversationCrudTestSuite) TestToggleStatus_InvalidAccountID() { @@ -598,13 +585,7 @@ func (s *ConversationCrudTestSuite) TestMute_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) - var resp struct { - Success bool `json:"success"` - Data map[string]interface{} `json:"data"` - } - err := json.Unmarshal(w.Body.Bytes(), &resp) - assert.NoError(s.T(), err) - assert.True(s.T(), resp.Success) + assert.Empty(s.T(), w.Body.String()) } func (s *ConversationCrudTestSuite) TestMute_InvalidAccountID() { @@ -643,13 +624,7 @@ func (s *ConversationCrudTestSuite) TestUnmute_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) - var resp struct { - Success bool `json:"success"` - Data map[string]interface{} `json:"data"` - } - err := json.Unmarshal(w.Body.Bytes(), &resp) - assert.NoError(s.T(), err) - assert.True(s.T(), resp.Success) + assert.Empty(s.T(), w.Body.String()) } func (s *ConversationCrudTestSuite) TestUnmute_InvalidAccountID() { @@ -689,13 +664,10 @@ func (s *ConversationCrudTestSuite) TestUpdateLabels_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) - var resp struct { - Success bool `json:"success"` - Data map[string]interface{} `json:"data"` - } + var resp map[string]interface{} err := json.Unmarshal(w.Body.Bytes(), &resp) assert.NoError(s.T(), err) - assert.True(s.T(), resp.Success) + assert.ElementsMatch(s.T(), []interface{}{"support", "bug"}, resp["labels"]) } func (s *ConversationCrudTestSuite) TestUpdateLabels_InvalidAccountID() { @@ -754,15 +726,15 @@ func (s *ConversationCrudTestSuite) TestSearch_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) var resp struct { - Success bool `json:"success"` - Data interface{} `json:"data"` - Meta struct { - TotalCount int64 `json:"total_count"` + Meta struct { + AllCount int64 `json:"all_count"` } `json:"meta"` + Payload []map[string]interface{} `json:"payload"` } err := json.Unmarshal(w.Body.Bytes(), &resp) assert.NoError(s.T(), err) - assert.True(s.T(), resp.Success) + assert.Equal(s.T(), int64(1), resp.Meta.AllCount) + assert.Len(s.T(), resp.Payload, 1) } func (s *ConversationCrudTestSuite) TestSearch_InvalidAccountID() { @@ -785,7 +757,7 @@ func (s *ConversationCrudTestSuite) TestSearch_MissingQuery() { func (s *ConversationCrudTestSuite) TestFilter_Success() { body, _ := json.Marshal(map[string]interface{}{ - "status": "open", + "status": "open", "assignee_type": "all", }) w := httptest.NewRecorder() @@ -796,15 +768,15 @@ func (s *ConversationCrudTestSuite) TestFilter_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) var resp struct { - Success bool `json:"success"` - Data interface{} `json:"data"` - Meta struct { - TotalCount int64 `json:"total_count"` + Meta struct { + AllCount int64 `json:"all_count"` } `json:"meta"` + Payload []map[string]interface{} `json:"payload"` } err := json.Unmarshal(w.Body.Bytes(), &resp) assert.NoError(s.T(), err) - assert.True(s.T(), resp.Success) + assert.Equal(s.T(), int64(1), resp.Meta.AllCount) + assert.Len(s.T(), resp.Payload, 1) } func (s *ConversationCrudTestSuite) TestFilter_InvalidAccountID() { @@ -841,14 +813,10 @@ func (s *ConversationCrudTestSuite) TestUpdatePriority_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) - var resp struct { - Success bool `json:"success"` - Data map[string]interface{} `json:"data"` - } + var resp map[string]interface{} err := json.Unmarshal(w.Body.Bytes(), &resp) assert.NoError(s.T(), err) - assert.True(s.T(), resp.Success) - assert.Equal(s.T(), "urgent", resp.Data["priority"]) + assert.Equal(s.T(), "urgent", resp["priority"]) } func (s *ConversationCrudTestSuite) TestUpdatePriority_InvalidAccountID() { @@ -925,13 +893,10 @@ func (s *ConversationCrudTestSuite) TestAssignTeam_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) - var resp struct { - Success bool `json:"success"` - Data map[string]interface{} `json:"data"` - } + var resp map[string]interface{} err := json.Unmarshal(w.Body.Bytes(), &resp) assert.NoError(s.T(), err) - assert.True(s.T(), resp.Success) + assert.Equal(s.T(), float64(s.testConv.ID), resp["id"]) } func (s *ConversationCrudTestSuite) TestAssignTeam_WithAgentID() { @@ -1009,4 +974,4 @@ func (s *ConversationCrudTestSuite) TestAssignTeam_NotFound() { func TestConversationCrudTestSuite(t *testing.T) { suite.Run(t, new(ConversationCrudTestSuite)) -} \ No newline at end of file +} diff --git a/internal/handler/api/v1/conversation_handler_test.go b/internal/handler/api/v1/conversation_handler_test.go index b35f6b55..e27d41d8 100644 --- a/internal/handler/api/v1/conversation_handler_test.go +++ b/internal/handler/api/v1/conversation_handler_test.go @@ -13,7 +13,6 @@ import ( "github.com/gin-gonic/gin" "github.com/stretchr/testify/assert" "github.com/stretchr/testify/suite" - "gorm.io/datatypes" "gorm.io/driver/sqlite" "gorm.io/gorm" "gorm.io/gorm/logger" @@ -209,12 +208,12 @@ func (s *ConversationHandlerTestSuite) TestUnread_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) var resp struct { - Success bool `json:"success"` - Data interface{} `json:"data"` + ID uint `json:"id"` + Status string `json:"status"` } err := json.Unmarshal(w.Body.Bytes(), &resp) assert.NoError(s.T(), err) - assert.True(s.T(), resp.Success) + assert.Equal(s.T(), s.testConv.ID, resp.ID) } func (s *ConversationHandlerTestSuite) TestUnread_InvalidAccountID() { @@ -330,16 +329,13 @@ func (s *ConversationHandlerTestSuite) TestUpdateCustomAttributes_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) var resp struct { - Success bool `json:"success"` - Data struct { - ID uint `json:"id"` - CustomAttributes datatypes.JSON `json:"custom_attributes"` - } `json:"data"` + ID uint `json:"id"` + CustomAttributes map[string]interface{} `json:"custom_attributes"` } err := json.Unmarshal(w.Body.Bytes(), &resp) assert.NoError(s.T(), err) - assert.True(s.T(), resp.Success) - assert.NotNil(s.T(), resp.Data.CustomAttributes) + assert.Equal(s.T(), s.testConv.ID, resp.ID) + assert.Equal(s.T(), "vip_customer", resp.CustomAttributes["priority_reason"]) } func (s *ConversationHandlerTestSuite) TestUpdateCustomAttributes_InvalidAccountID() { @@ -465,4 +461,4 @@ func (s *ConversationHandlerTestSuite) TestUpdateLastSeen_InvalidConversationID( // Run the test suite func TestConversationHandlerTestSuite(t *testing.T) { suite.Run(t, new(ConversationHandlerTestSuite)) -} \ No newline at end of file +} diff --git a/internal/handler/api/v1/conversation_serializer.go b/internal/handler/api/v1/conversation_serializer.go new file mode 100644 index 00000000..030f54df --- /dev/null +++ b/internal/handler/api/v1/conversation_serializer.go @@ -0,0 +1,387 @@ +package v1 + +import ( + "context" + "encoding/json" + "strings" + "time" + + "github.com/gochat/gochat/internal/model" + "gorm.io/datatypes" + "gorm.io/gorm" +) + +type chatwootConversationListResponse struct { + Data chatwootConversationListData `json:"data"` +} + +type chatwootConversationListData struct { + Meta chatwootConversationCounts `json:"meta"` + Payload []chatwootConversationPayload `json:"payload"` +} + +type chatwootConversationCounts struct { + MineCount int64 `json:"mine_count"` + AssignedCount int64 `json:"assigned_count"` + UnassignedCount int64 `json:"unassigned_count"` + AllCount int64 `json:"all_count"` +} + +type chatwootConversationPayload struct { + Meta chatwootConversationMeta `json:"meta"` + ID uint `json:"id"` + Messages []chatwootMessagePayload `json:"messages"` + AccountID uint `json:"account_id"` + UUID string `json:"uuid"` + AdditionalAttributes map[string]any `json:"additional_attributes"` + AgentLastSeenAt int64 `json:"agent_last_seen_at"` + AssigneeLastSeenAt int64 `json:"assignee_last_seen_at"` + CanReply bool `json:"can_reply"` + ContactLastSeenAt int64 `json:"contact_last_seen_at"` + CustomAttributes map[string]any `json:"custom_attributes"` + InboxID uint `json:"inbox_id"` + Labels []string `json:"labels"` + Muted bool `json:"muted"` + SnoozedUntil *int64 `json:"snoozed_until"` + Status string `json:"status"` + CreatedAt int64 `json:"created_at"` + UpdatedAt float64 `json:"updated_at"` + Timestamp int64 `json:"timestamp"` + FirstReplyCreatedAt int64 `json:"first_reply_created_at"` + UnreadCount int64 `json:"unread_count"` + LastNonActivityMessage *chatwootMessagePayload `json:"last_non_activity_message"` + LastActivityAt int64 `json:"last_activity_at"` + Priority string `json:"priority"` + WaitingSince int64 `json:"waiting_since"` + SlaPolicyID *uint `json:"sla_policy_id"` +} + +type chatwootConversationMeta struct { + Sender map[string]any `json:"sender"` + Channel string `json:"channel"` + Assignee map[string]any `json:"assignee,omitempty"` + AssigneeType string `json:"assignee_type,omitempty"` + Team map[string]any `json:"team,omitempty"` + HMACVerified *bool `json:"hmac_verified,omitempty"` +} + +type chatwootMessageIndexResponse struct { + Meta chatwootMessageIndexMeta `json:"meta"` + Payload []chatwootMessagePayload `json:"payload"` +} + +type chatwootMessageIndexMeta struct { + Labels []string `json:"labels"` + AdditionalAttrs map[string]any `json:"additional_attributes"` + Contact map[string]any `json:"contact"` + Assignee map[string]any `json:"assignee,omitempty"` + AgentLastSeenAt int64 `json:"agent_last_seen_at"` + AssigneeLastSeenAt int64 `json:"assignee_last_seen_at"` +} + +type chatwootMessagePayload struct { + ID uint `json:"id"` + Content string `json:"content"` + InboxID uint `json:"inbox_id"` + EchoID string `json:"echo_id,omitempty"` + ConversationID uint `json:"conversation_id"` + MessageType int `json:"message_type"` + ContentType string `json:"content_type"` + Status string `json:"status"` + ContentAttributes map[string]any `json:"content_attributes"` + CreatedAt int64 `json:"created_at"` + Private bool `json:"private"` + SourceID string `json:"source_id"` + Sender map[string]any `json:"sender,omitempty"` + Attachments []any `json:"attachments,omitempty"` +} + +func serializeConversationList(ctx context.Context, db *gorm.DB, conversations []model.Conversation, total int64) chatwootConversationListResponse { + payload := make([]chatwootConversationPayload, 0, len(conversations)) + for i := range conversations { + payload = append(payload, serializeConversation(ctx, db, &conversations[i])) + } + return chatwootConversationListResponse{Data: chatwootConversationListData{ + Meta: chatwootConversationCounts{AllCount: total}, + Payload: payload, + }} +} + +func serializeConversation(ctx context.Context, db *gorm.DB, conversation *model.Conversation) chatwootConversationPayload { + var lastMessage *model.Message + if db != nil { + var msg model.Message + if err := db.WithContext(ctx).Where("account_id = ? AND conversation_id = ?", conversation.AccountID, conversation.ID).Order("id DESC").First(&msg).Error; err == nil { + lastMessage = &msg + } + } + + messages := []chatwootMessagePayload{} + var lastNonActivity *chatwootMessagePayload + if lastMessage != nil { + serialized := serializeMessage(ctx, db, lastMessage, conversation) + messages = append(messages, serialized) + if lastMessage.MessageType != "activity" { + lastNonActivity = &serialized + } + } + + return chatwootConversationPayload{ + Meta: serializeConversationMeta(ctx, db, conversation), + ID: conversationDisplayID(conversation), + Messages: messages, + AccountID: conversation.AccountID, + UUID: conversation.UUID, + AdditionalAttributes: jsonObject(conversation.AdditionalAttributes), + AgentLastSeenAt: int64Value(conversation.AgentLastSeenAt), + AssigneeLastSeenAt: int64Value(conversation.AssigneeLastSeenAt), + CanReply: true, + ContactLastSeenAt: int64Value(conversation.ContactLastSeenAt), + CustomAttributes: jsonObject(conversation.CustomAttributes), + InboxID: conversation.InboxID, + Labels: labelList(conversation.Labels), + Muted: conversation.Muted, + SnoozedUntil: conversation.SnoozedUntil, + Status: conversation.Status, + CreatedAt: conversation.CreatedAt.Unix(), + UpdatedAt: float64(conversation.UpdatedAt.UnixNano()) / float64(time.Second), + Timestamp: int64Value(conversation.LastActivityAt), + FirstReplyCreatedAt: int64Value(conversation.FirstReplyCreatedAt), + UnreadCount: unreadCount(ctx, db, conversation), + LastNonActivityMessage: lastNonActivity, + LastActivityAt: int64Value(conversation.LastActivityAt), + Priority: conversation.Priority, + WaitingSince: int64Value(conversation.WaitingSince), + SlaPolicyID: conversation.SlaPolicyID, + } +} + +func serializeConversationMeta(ctx context.Context, db *gorm.DB, conversation *model.Conversation) chatwootConversationMeta { + meta := chatwootConversationMeta{Channel: conversation.ChannelType} + if db == nil { + meta.Sender = map[string]any{"id": conversation.ContactID} + return meta + } + + var contact model.Contact + if err := db.WithContext(ctx).First(&contact, conversation.ContactID).Error; err == nil { + meta.Sender = serializeContact(&contact) + } else { + meta.Sender = map[string]any{"id": conversation.ContactID} + } + + if conversation.AssigneeID != nil && *conversation.AssigneeID != 0 { + var user model.User + if err := db.WithContext(ctx).First(&user, *conversation.AssigneeID).Error; err == nil { + meta.Assignee = serializeUser(&user, conversation.AccountID) + meta.AssigneeType = "User" + } + } + + if conversation.ContactInboxID != nil && *conversation.ContactInboxID != 0 { + var contactInbox model.ContactInbox + if err := db.WithContext(ctx).First(&contactInbox, *conversation.ContactInboxID).Error; err == nil { + meta.HMACVerified = &contactInbox.HMACVerified + } + } + return meta +} + +func serializeMessageIndex(ctx context.Context, db *gorm.DB, conversation *model.Conversation, messages []model.Message) chatwootMessageIndexResponse { + payload := make([]chatwootMessagePayload, 0, len(messages)) + for i := range messages { + payload = append(payload, serializeMessage(ctx, db, &messages[i], conversation)) + } + + meta := chatwootMessageIndexMeta{ + Labels: labelList(conversation.Labels), + AdditionalAttrs: jsonObject(conversation.AdditionalAttributes), + Contact: map[string]any{"id": conversation.ContactID}, + AgentLastSeenAt: int64Value(conversation.AgentLastSeenAt), + AssigneeLastSeenAt: int64Value(conversation.AssigneeLastSeenAt), + } + if db != nil { + var contact model.Contact + if err := db.WithContext(ctx).First(&contact, conversation.ContactID).Error; err == nil { + meta.Contact = serializeContact(&contact) + } + if conversation.AssigneeID != nil && *conversation.AssigneeID != 0 { + var user model.User + if err := db.WithContext(ctx).First(&user, *conversation.AssigneeID).Error; err == nil { + meta.Assignee = serializeUser(&user, conversation.AccountID) + } + } + } + + return chatwootMessageIndexResponse{Meta: meta, Payload: payload} +} + +func serializeMessage(ctx context.Context, db *gorm.DB, message *model.Message, conversation *model.Conversation) chatwootMessagePayload { + conversationID := message.ConversationID + if conversation != nil { + conversationID = conversationDisplayID(conversation) + } + payload := chatwootMessagePayload{ + ID: message.ID, + Content: message.Content, + InboxID: message.InboxID, + EchoID: message.EchoID, + ConversationID: conversationID, + MessageType: messageTypeValue(message.MessageType), + ContentType: nonEmpty(message.ContentType, "text"), + Status: nonEmpty(message.Status, "sent"), + ContentAttributes: jsonObject(message.ContentAttributes), + CreatedAt: message.CreatedAt.Unix(), + Private: message.Private, + SourceID: message.SourceID, + } + if db != nil && message.SenderID != nil && *message.SenderID != 0 { + senderType := strings.ToLower(message.SenderType) + if senderType == "contact" { + var contact model.Contact + if err := db.WithContext(ctx).First(&contact, *message.SenderID).Error; err == nil { + payload.Sender = serializeContact(&contact) + } + } else { + var user model.User + if err := db.WithContext(ctx).First(&user, *message.SenderID).Error; err == nil { + payload.Sender = serializeUser(&user, message.AccountID) + } + } + } + return payload +} + +func serializeContact(contact *model.Contact) map[string]any { + return map[string]any{ + "additional_attributes": jsonObject(contact.AdditionalAttributes), + "availability_status": "offline", + "email": contact.Email, + "id": contact.ID, + "name": contact.Name, + "phone_number": contact.PhoneNumber, + "blocked": contact.Blocked, + "identifier": contact.Identifier, + "thumbnail": contact.AvatarURL, + "custom_attributes": jsonObject(contact.CustomAttributes), + "last_activity_at": int64Value(contact.LastActivityAt), + "created_at": contact.CreatedAt.Unix(), + } +} + +func serializeUser(user *model.User, accountID uint) map[string]any { + return map[string]any{ + "id": user.ID, + "account_id": accountID, + "availability_status": availabilityStatus(user.Available), + "auto_offline": false, + "confirmed": user.ConfirmedAt != nil, + "email": user.Email, + "provider": nonEmpty(user.Provider, "email"), + "available_name": nonEmpty(user.DisplayName, user.Name), + "name": user.Name, + "role": nonEmpty(user.Role, "agent"), + "thumbnail": user.AvatarURL, + } +} + +func serializeUserFromDB(ctx context.Context, db *gorm.DB, userID uint, accountID uint) any { + if userID == 0 || db == nil { + return nil + } + var user model.User + if err := db.WithContext(ctx).First(&user, userID).Error; err != nil { + return nil + } + return serializeUser(&user, accountID) +} + +func conversationDisplayID(conversation *model.Conversation) uint { + if conversation.DisplayID != nil && *conversation.DisplayID != 0 { + return *conversation.DisplayID + } + return conversation.ID +} + +func messageTypeValue(value string) int { + switch strings.ToLower(value) { + case "incoming": + return 0 + case "outgoing", "private_note": + return 1 + case "activity": + return 2 + case "template": + return 3 + default: + return 1 + } +} + +func labelList(labels string) []string { + labels = strings.TrimSpace(labels) + if labels == "" { + return []string{} + } + if strings.HasPrefix(labels, "[") { + var list []string + if err := json.Unmarshal([]byte(labels), &list); err == nil { + return list + } + } + parts := strings.Split(labels, ",") + result := make([]string, 0, len(parts)) + for _, part := range parts { + label := strings.TrimSpace(part) + if label != "" { + result = append(result, label) + } + } + return result +} + +func jsonObject(raw datatypes.JSON) map[string]any { + if len(raw) == 0 || string(raw) == "null" { + return map[string]any{} + } + var value map[string]any + if err := json.Unmarshal(raw, &value); err != nil || value == nil { + return map[string]any{} + } + return value +} + +func int64Value(value *int64) int64 { + if value == nil { + return 0 + } + return *value +} + +func nonEmpty(value, fallback string) string { + if value == "" { + return fallback + } + return value +} + +func availabilityStatus(available bool) string { + if available { + return "online" + } + return "offline" +} + +func unreadCount(ctx context.Context, db *gorm.DB, conversation *model.Conversation) int64 { + if db == nil { + return 0 + } + query := db.WithContext(ctx).Model(&model.Message{}). + Where("account_id = ? AND conversation_id = ? AND message_type = ?", conversation.AccountID, conversation.ID, "incoming") + if conversation.AgentLastSeenAt != nil { + query = query.Where("created_at > ?", time.Unix(*conversation.AgentLastSeenAt, 0)) + } + var count int64 + _ = query.Count(&count).Error + return count +} diff --git a/internal/handler/api/v1/message_handler.go b/internal/handler/api/v1/message_handler.go index 2c3899d9..12da639b 100644 --- a/internal/handler/api/v1/message_handler.go +++ b/internal/handler/api/v1/message_handler.go @@ -1,6 +1,7 @@ package v1 import ( + "encoding/json" "net/http" "strings" @@ -10,6 +11,7 @@ import ( "github.com/gochat/gochat/internal/service" "github.com/gochat/gochat/pkg/pagination" "github.com/gochat/gochat/pkg/response" + "gorm.io/datatypes" ) // MessageHandler handles message-related API endpoints. @@ -43,21 +45,31 @@ func NewMessageHandler(svc *service.MessageService) *MessageHandler { // GET /api/v1/accounts/:account_id/conversations/:conversation_id/messages // Reference: Chatwoot conversations#messages (index) func (h *MessageHandler) List(c *gin.Context) { + accountID, err := parseUintParam(c, "account_id") + if err != nil { + response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrBadRequest, "invalid account_id") + return + } conversationID, err := parseUintParam(c, "conversation_id") if err != nil { response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrBadRequest, "invalid conversation_id") return } - - p := pagination.Parse(c) - - messages, total, svcErr := h.svc.ListByConversation(c.Request.Context(), conversationID, p.Offset, p.PerPage) + conversation, svcErr := h.svc.ResolveConversationForRoute(c.Request.Context(), accountID, conversationID) if svcErr != nil { handleServiceError(c, svcErr) return } - response.OKWithMeta(c, toInterfaceSlice(messages), p.Page, p.PerPage, total) + p := pagination.Parse(c) + + messages, _, svcErr := h.svc.ListByConversation(c.Request.Context(), conversation.ID, p.Offset, p.PerPage) + if svcErr != nil { + handleServiceError(c, svcErr) + return + } + + c.JSON(http.StatusOK, serializeMessageIndex(c.Request.Context(), h.svc.DB(), conversation, messages)) } // @Summary Create a message in a conversation @@ -92,7 +104,7 @@ func (h *MessageHandler) Create(c *gin.Context) { } var req service.CreateMessageRequest - if err := c.ShouldBindJSON(&req); err != nil { + if err := bindCreateMessageRequest(c, &req); err != nil { response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrValidation, err.Error()) return } @@ -104,7 +116,12 @@ func (h *MessageHandler) Create(c *gin.Context) { response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrBadRequest, "conversation_id is required") return } - req.ConversationID = conversationID + conversation, resolveErr := h.svc.ResolveConversationForRoute(c.Request.Context(), accountID, conversationID) + if resolveErr != nil { + handleServiceError(c, resolveErr) + return + } + req.ConversationID = conversation.ID } message, svcErr := h.svc.Create(c.Request.Context(), accountID, userID, req) @@ -112,7 +129,8 @@ func (h *MessageHandler) Create(c *gin.Context) { handleServiceError(c, svcErr) return } - response.Created(c, message) + conversation, _ := h.svc.ResolveConversationForRoute(c.Request.Context(), accountID, req.ConversationID) + c.JSON(http.StatusOK, serializeMessage(c.Request.Context(), h.svc.DB(), message, conversation)) } // @Summary Get a single message @@ -178,7 +196,8 @@ func (h *MessageHandler) Update(c *gin.Context) { handleServiceError(c, svcErr) return } - response.OK(c, message) + conversation, _ := h.svc.ResolveConversationForRoute(c.Request.Context(), accountID, message.ConversationID) + c.JSON(http.StatusOK, serializeMessage(c.Request.Context(), h.svc.DB(), message, conversation)) } // @Summary Delete a message @@ -264,7 +283,8 @@ func (h *MessageHandler) Retry(c *gin.Context) { handleServiceError(c, svcErr) return } - response.OK(c, message) + conversation, _ := h.svc.ResolveConversationForRoute(c.Request.Context(), accountID, message.ConversationID) + c.JSON(http.StatusOK, serializeMessage(c.Request.Context(), h.svc.DB(), message, conversation)) } // Translate translates a message's content to a target language using LLM. @@ -304,3 +324,34 @@ func (h *MessageHandler) Translate(c *gin.Context) { } response.OK(c, result) } + +func bindCreateMessageRequest(c *gin.Context, req *service.CreateMessageRequest) error { + if strings.HasPrefix(c.ContentType(), "multipart/form-data") { + if err := c.Request.ParseMultipartForm(32 << 20); err != nil { + return err + } + req.Content = c.PostForm("content") + req.MessageType = c.PostForm("message_type") + req.ContentType = c.PostForm("content_type") + req.SourceID = c.PostForm("source_id") + req.EchoID = c.PostForm("echo_id") + req.Private = strings.EqualFold(c.PostForm("private"), "true") || c.PostForm("private") == "1" + if raw := c.PostForm("content_attributes"); raw != "" { + req.ContentAttributes = []byte(raw) + } + return nil + } + + var raw map[string]json.RawMessage + if err := c.ShouldBindJSON(&raw); err != nil { + return err + } + bytes, _ := json.Marshal(raw) + if err := json.Unmarshal(bytes, req); err != nil { + return err + } + if value, ok := raw["content_attributes"]; ok && string(value) != "null" { + req.ContentAttributes = datatypes.JSON(value) + } + return nil +} diff --git a/internal/handler/api/v1/message_handler_test.go b/internal/handler/api/v1/message_handler_test.go index 3d8be5ca..44cbf2bd 100644 --- a/internal/handler/api/v1/message_handler_test.go +++ b/internal/handler/api/v1/message_handler_test.go @@ -84,7 +84,7 @@ func (s *MessageHandlerTestSuite) SetupSuite() { ) s.Require().NoError(err) -// Wire up repos, services, handlers + // Wire up repos, services, handlers msgRepo := repository.NewMessageRepo(db) s.dispatcher = channel.NewDispatcher() s.mockLLM = &mockMsgHandlerLLMProvider{ @@ -173,7 +173,7 @@ func (s *MessageHandlerTestSuite) SetupTest() { s.Require().NoError(s.db.Create(msg).Error) s.testMessage = msg -// Reset mock LLM to default success response + // Reset mock LLM to default success response s.mockLLM.chatResponse = &llm.ChatResponse{ Choices: []llm.ChatChoice{ { @@ -199,7 +199,7 @@ func TestMessageHandlerTestSuite(t *testing.T) { suite.Run(t, new(MessageHandlerTestSuite)) } -// --- Helper for building message list URL --- +// --- Helper for building message list URL --- func msgListURL(accountID, convID uint) string { return fmt.Sprintf("/api/v1/accounts/%d/conversations/%d/messages/", accountID, convID) @@ -217,7 +217,7 @@ func msgTranslateURL(accountID, convID, msgID uint) string { return fmt.Sprintf("/api/v1/accounts/%d/conversations/%d/messages/%d/translate", accountID, convID, msgID) } -// --- List Tests --- +// --- List Tests --- func (s *MessageHandlerTestSuite) TestList_Success() { w := httptest.NewRecorder() @@ -228,10 +228,12 @@ func (s *MessageHandlerTestSuite) TestList_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) var resp map[string]interface{} json.Unmarshal(w.Body.Bytes(), &resp) - assert.True(s.T(), resp["success"].(bool)) - data, ok := resp["data"].([]interface{}) + data, ok := resp["payload"].([]interface{}) assert.True(s.T(), ok) assert.GreaterOrEqual(s.T(), len(data), 1) + meta, ok := resp["meta"].(map[string]interface{}) + assert.True(s.T(), ok) + assert.NotNil(s.T(), meta["contact"]) } func (s *MessageHandlerTestSuite) TestList_Empty() { @@ -246,7 +248,7 @@ func (s *MessageHandlerTestSuite) TestList_Empty() { assert.Equal(s.T(), http.StatusOK, w.Code) var resp map[string]interface{} json.Unmarshal(w.Body.Bytes(), &resp) - data, ok := resp["data"].([]interface{}) + data, ok := resp["payload"].([]interface{}) assert.True(s.T(), ok) assert.Equal(s.T(), 0, len(data)) } @@ -259,7 +261,7 @@ func (s *MessageHandlerTestSuite) TestList_InvalidConversationID() { assert.Equal(s.T(), http.StatusBadRequest, w.Code) } -// --- Create Tests --- +// --- Create Tests --- func (s *MessageHandlerTestSuite) TestCreate_Success() { payload := map[string]interface{}{ @@ -276,13 +278,39 @@ func (s *MessageHandlerTestSuite) TestCreate_Success() { req.Header.Set("Content-Type", "application/json") s.router.ServeHTTP(w, req) - assert.Equal(s.T(), http.StatusCreated, w.Code) + assert.Equal(s.T(), http.StatusOK, w.Code) var resp map[string]interface{} json.Unmarshal(w.Body.Bytes(), &resp) - assert.True(s.T(), resp["success"].(bool)) - data := resp["data"].(map[string]interface{}) - assert.NotNil(s.T(), data["id"]) - assert.Equal(s.T(), "New message", data["content"]) + assert.NotNil(s.T(), resp["id"]) + assert.Equal(s.T(), "New message", resp["content"]) + assert.Equal(s.T(), float64(1), resp["message_type"]) + assert.Equal(s.T(), float64(s.testConv.ID), resp["conversation_id"]) +} + +func (s *MessageHandlerTestSuite) TestCreate_ChatwootFrontendPayloadDefaultsOutgoing() { + payload := map[string]interface{}{ + "content": "Frontend payload", + "private": true, + "echo_id": "tmp-123", + "content_attributes": map[string]interface{}{"submitted_values": []interface{}{}}, + } + body, _ := json.Marshal(payload) + + w := httptest.NewRecorder() + url := msgListURL(s.testAccount.ID, s.testConv.ID) + req, _ := http.NewRequest("POST", url, bytes.NewReader(body)) + req.Header.Set("Content-Type", "application/json") + s.router.ServeHTTP(w, req) + + assert.Equal(s.T(), http.StatusOK, w.Code) + var resp map[string]interface{} + json.Unmarshal(w.Body.Bytes(), &resp) + assert.Equal(s.T(), "Frontend payload", resp["content"]) + assert.Equal(s.T(), true, resp["private"]) + assert.Equal(s.T(), "tmp-123", resp["echo_id"]) + assert.Equal(s.T(), float64(1), resp["message_type"]) + assert.Equal(s.T(), "text", resp["content_type"]) + assert.NotNil(s.T(), resp["content_attributes"]) } func (s *MessageHandlerTestSuite) TestCreate_MissingContent() { @@ -316,7 +344,7 @@ func (s *MessageHandlerTestSuite) TestCreate_InvalidConversationID() { assert.Equal(s.T(), http.StatusBadRequest, w.Code) } -// --- Get Tests --- +// --- Get Tests --- func (s *MessageHandlerTestSuite) TestGet_Success() { w := httptest.NewRecorder() @@ -350,7 +378,7 @@ func (s *MessageHandlerTestSuite) TestGet_InvalidID() { assert.Equal(s.T(), http.StatusBadRequest, w.Code) } -// --- Update Tests --- +// --- Update Tests --- func (s *MessageHandlerTestSuite) TestUpdate_Success() { payload := map[string]interface{}{ @@ -367,9 +395,7 @@ func (s *MessageHandlerTestSuite) TestUpdate_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) var resp map[string]interface{} json.Unmarshal(w.Body.Bytes(), &resp) - assert.True(s.T(), resp["success"].(bool)) - data := resp["data"].(map[string]interface{}) - assert.Equal(s.T(), "Updated content", data["content"]) + assert.Equal(s.T(), "Updated content", resp["content"]) } func (s *MessageHandlerTestSuite) TestUpdate_NotFound() { @@ -387,7 +413,7 @@ func (s *MessageHandlerTestSuite) TestUpdate_NotFound() { assert.Equal(s.T(), http.StatusNotFound, w.Code) } -// --- Delete Tests --- +// --- Delete Tests --- func (s *MessageHandlerTestSuite) TestDelete_Success() { w := httptest.NewRecorder() @@ -407,7 +433,7 @@ func (s *MessageHandlerTestSuite) TestDelete_NotFound() { assert.Equal(s.T(), http.StatusNotFound, w.Code) } -// --- Retry Tests --- +// --- Retry Tests --- func (s *MessageHandlerTestSuite) TestRetry_Success() { w := httptest.NewRecorder() @@ -418,9 +444,8 @@ func (s *MessageHandlerTestSuite) TestRetry_Success() { assert.Equal(s.T(), http.StatusOK, w.Code) var resp map[string]interface{} json.Unmarshal(w.Body.Bytes(), &resp) - assert.True(s.T(), resp["success"].(bool)) - data := resp["data"].(map[string]interface{}) - assert.Equal(s.T(), float64(s.testMessage.ID), data["id"]) + assert.Equal(s.T(), float64(s.testMessage.ID), resp["id"]) + assert.Equal(s.T(), float64(1), resp["message_type"]) } func (s *MessageHandlerTestSuite) TestRetry_InvalidAccountID() { @@ -464,7 +489,7 @@ func (s *MessageHandlerTestSuite) TestRetry_AccountMismatch() { assert.Equal(s.T(), http.StatusNotFound, w.Code) } -// --- Translate Tests --- +// --- Translate Tests --- func (s *MessageHandlerTestSuite) TestTranslate_Success() { payload := map[string]interface{}{ @@ -579,7 +604,7 @@ func (s *MessageHandlerTestSuite) TestTranslate_LLMError() { } func (s *MessageHandlerTestSuite) TestTranslate_EmptyChoices() { -// Set LLM to return response with no choices + // Set LLM to return response with no choices s.mockLLM.chatResponse = &llm.ChatResponse{ Choices: []llm.ChatChoice{}, } @@ -637,4 +662,4 @@ func TestDirectUploadHandler_InvalidContentType(t *testing.T) { r.ServeHTTP(w, req) assert.Equal(t, http.StatusBadRequest, w.Code) -} \ No newline at end of file +} diff --git a/internal/model/message.go b/internal/model/message.go index 77343180..a415060c 100644 --- a/internal/model/message.go +++ b/internal/model/message.go @@ -7,21 +7,22 @@ import ( // Message represents a message within a conversation. type Message struct { Base - ConversationID uint `gorm:"index;not null" json:"conversation_id"` - AccountID uint `gorm:"index;not null" json:"account_id"` - InboxID uint `gorm:"index;not null" json:"inbox_id"` - SenderID *uint `gorm:"index" json:"sender_id,omitempty"` - SenderType string `gorm:"size:50" json:"sender_type"` // contact, agent, bot - Content string `gorm:"type:text" json:"content"` - ContentType string `gorm:"size:50;default:text" json:"content_type"` // text, input, input_csat, file, image, etc - Status string `gorm:"size:50;default:sent" json:"status"` // sent, delivered, read, failed - Private bool `gorm:"default:false" json:"private"` - SourceID string `gorm:"size:255" json:"source_id,omitempty"` - MessageType string `gorm:"size:50;default:incoming" json:"message_type"` // incoming, outgoing, activity, template - External bool `gorm:"default:false" json:"external"` - ContentAttributes datatypes.JSON `gorm:"type:jsonb" json:"content_attributes,omitempty"` // attachments, mentions, etc + ConversationID uint `gorm:"index;not null" json:"conversation_id"` + AccountID uint `gorm:"index;not null" json:"account_id"` + InboxID uint `gorm:"index;not null" json:"inbox_id"` + SenderID *uint `gorm:"index" json:"sender_id,omitempty"` + SenderType string `gorm:"size:50" json:"sender_type"` // contact, agent, bot + Content string `gorm:"type:text" json:"content"` + ContentType string `gorm:"size:50;default:text" json:"content_type"` // text, input, input_csat, file, image, etc + Status string `gorm:"size:50;default:sent" json:"status"` // sent, delivered, read, failed + Private bool `gorm:"default:false" json:"private"` + EchoID string `gorm:"size:255" json:"echo_id,omitempty"` + SourceID string `gorm:"size:255" json:"source_id,omitempty"` + MessageType string `gorm:"size:50;default:incoming" json:"message_type"` // incoming, outgoing, activity, template + External bool `gorm:"default:false" json:"external"` + ContentAttributes datatypes.JSON `gorm:"type:jsonb" json:"content_attributes,omitempty"` // attachments, mentions, etc AdditionalAttributes datatypes.JSON `gorm:"type:jsonb" json:"additional_attributes,omitempty"` - ExternalSourceIDs datatypes.JSON `gorm:"type:jsonb" json:"external_source_ids,omitempty"` // external platform IDs + ExternalSourceIDs datatypes.JSON `gorm:"type:jsonb" json:"external_source_ids,omitempty"` // external platform IDs } -func (Message) TableName() string { return "messages" } \ No newline at end of file +func (Message) TableName() string { return "messages" } diff --git a/internal/repository/conversation_repo.go b/internal/repository/conversation_repo.go index b60b6a2c..c7276b9a 100644 --- a/internal/repository/conversation_repo.go +++ b/internal/repository/conversation_repo.go @@ -48,6 +48,23 @@ func (r *ConversationRepo) FindByAccountAndID(ctx context.Context, accountID, id return &conversation, nil } +// FindByAccountAndDisplayIDOrID retrieves a conversation using Chatwoot's +// account-scoped display_id route semantics, falling back to primary key for +// legacy GoChat data and tests that predate display_id. +func (r *ConversationRepo) FindByAccountAndDisplayIDOrID(ctx context.Context, accountID, routeID uint) (*model.Conversation, error) { + var conversation model.Conversation + err := r.db.WithContext(ctx).Where("account_id = ? AND display_id = ?", accountID, routeID).First(&conversation).Error + if err == nil { + return &conversation, nil + } + + err = r.db.WithContext(ctx).Where("account_id = ? AND id = ?", accountID, routeID).First(&conversation).Error + if err != nil { + return nil, err + } + return &conversation, nil +} + // FindByAccount retrieves all conversations for an account. func (r *ConversationRepo) FindByAccount(ctx context.Context, accountID uint, offset, limit int) ([]model.Conversation, int64, error) { var conversations []model.Conversation @@ -239,6 +256,16 @@ func (r *ConversationRepo) Search(ctx context.Context, accountID uint, query str // Create inserts a new conversation. func (r *ConversationRepo) Create(ctx context.Context, conversation *model.Conversation) error { + if conversation.DisplayID == nil || *conversation.DisplayID == 0 { + var next uint + if err := r.db.WithContext(ctx).Model(&model.Conversation{}). + Select("COALESCE(MAX(display_id), 0) + 1"). + Where("account_id = ?", conversation.AccountID). + Scan(&next).Error; err != nil { + return err + } + conversation.DisplayID = &next + } return r.db.WithContext(ctx).Create(conversation).Error } diff --git a/internal/service/conversation_service.go b/internal/service/conversation_service.go index f7e75859..58731f80 100644 --- a/internal/service/conversation_service.go +++ b/internal/service/conversation_service.go @@ -15,6 +15,7 @@ import ( pkgvalidator "github.com/gochat/gochat/pkg/validator" "gorm.io/datatypes" + "gorm.io/gorm" ) // ConversationService implements business logic for Conversation operations. @@ -39,6 +40,13 @@ func (s *ConversationService) SetSearchIndexer(indexer SearchIndexer) { s.searchIndexer = indexer } +func (s *ConversationService) DB() *gorm.DB { + if s == nil || s.repo == nil { + return nil + } + return s.repo.DB() +} + func (s *ConversationService) indexConversation(ctx context.Context, conversation *model.Conversation) { if s.searchIndexer != nil { logSearchIndexError("conversation", conversation.ID, s.searchIndexer.IndexConversation(ctx, conversation)) @@ -115,6 +123,11 @@ func (s *ConversationService) GetByAccountAndID(ctx context.Context, accountID, return s.repo.FindByAccountAndID(ctx, accountID, id) } +// GetByAccountAndDisplayIDOrID resolves Chatwoot conversation route IDs. +func (s *ConversationService) GetByAccountAndDisplayIDOrID(ctx context.Context, accountID, routeID uint) (*model.Conversation, error) { + return s.repo.FindByAccountAndDisplayIDOrID(ctx, accountID, routeID) +} + // CreateConversationRequest is the DTO for creating a conversation. // Reference: Chatwoot app/controllers/api/v1/accounts/conversations_controller.rb #create // Supports creating conversation with an initial message (like Chatwoot's ConversationBuilder). diff --git a/internal/service/message_service.go b/internal/service/message_service.go index c5028bfb..c030cd0b 100644 --- a/internal/service/message_service.go +++ b/internal/service/message_service.go @@ -3,6 +3,7 @@ package service import ( "context" "fmt" + "strings" "github.com/gochat/gochat/internal/channel" "github.com/gochat/gochat/internal/llm" @@ -11,6 +12,9 @@ import ( "github.com/gochat/gochat/internal/search" applogger "github.com/gochat/gochat/pkg/logger" pkgvalidator "github.com/gochat/gochat/pkg/validator" + + "gorm.io/datatypes" + "gorm.io/gorm" ) // MessageService implements business logic for Message operations. @@ -31,6 +35,13 @@ func (s *MessageService) SetSearchIndexer(indexer SearchIndexer) { s.searchIndexer = indexer } +func (s *MessageService) DB() *gorm.DB { + if s == nil || s.repo == nil { + return nil + } + return s.repo.DB() +} + func (s *MessageService) indexMessage(ctx context.Context, message *model.Message) { if s.searchIndexer != nil { logSearchIndexError("message", message.ID, s.searchIndexer.IndexMessage(ctx, message)) @@ -62,6 +73,18 @@ func (s *MessageService) ListByConversation(ctx context.Context, conversationID return s.repo.FindByConversation(ctx, conversationID, offset, limit) } +func (s *MessageService) ResolveConversationForRoute(ctx context.Context, accountID, routeID uint) (*model.Conversation, error) { + var conversation model.Conversation + db := s.repo.DB().WithContext(ctx) + if err := db.Where("account_id = ? AND display_id = ?", accountID, routeID).First(&conversation).Error; err == nil { + return &conversation, nil + } + if err := db.Where("account_id = ? AND id = ?", accountID, routeID).First(&conversation).Error; err != nil { + return nil, err + } + return &conversation, nil +} + // GetByID retrieves a single message. func (s *MessageService) GetByID(ctx context.Context, id uint) (*model.Message, error) { return s.repo.FindByID(ctx, id) @@ -85,34 +108,51 @@ func (s *MessageService) Search(ctx context.Context, accountID uint, query strin // CreateMessageRequest is the DTO for creating a message. // Reference: Chatwoot app/controllers/api/v1/accounts/conversations/messages_controller.rb #create type CreateMessageRequest struct { - ConversationID uint `json:"conversation_id" validate:"required"` - Content string `json:"content" validate:"required,min=1"` - MessageType string `json:"message_type" validate:"required,oneof=outgoing incoming activity template private_note"` - ContentType string `json:"content_type,omitempty" validate:"omitempty,oneof=text input_text input_email input_phone select card private_note"` - Private bool `json:"private,omitempty"` - SourceID string `json:"source_id,omitempty"` // Chatwoot: source_id for message origin (user/agent/bot) + ConversationID uint `json:"conversation_id" validate:"required"` + Content string `json:"content" validate:"required,min=1"` + MessageType string `json:"message_type,omitempty"` + ContentType string `json:"content_type,omitempty"` + Private bool `json:"private,omitempty"` + SourceID string `json:"source_id,omitempty"` + EchoID string `json:"echo_id,omitempty"` + ContentAttributes datatypes.JSON `json:"content_attributes,omitempty"` } // Create creates a new message. func (s *MessageService) Create(ctx context.Context, accountID uint, userID uint, req CreateMessageRequest) (*model.Message, error) { + req.MessageType = normalizeMessageType(req.MessageType) + if req.ContentType == "" { + req.ContentType = "text" + } + if !validMessageType(req.MessageType) { + return nil, fmt.Errorf("invalid message_type") + } + if !validContentType(req.ContentType) { + return nil, fmt.Errorf("invalid content_type") + } if err := pkgvalidator.ValidateStruct(req); err != nil { return nil, err } - message := &model.Message{ - AccountID: accountID, - ConversationID: req.ConversationID, - Content: req.Content, - MessageType: req.MessageType, - ContentType: req.ContentType, - SenderID: &userID, - SenderType: "user", - Private: req.Private, - SourceID: req.SourceID, + var conversation model.Conversation + if err := s.repo.DB().WithContext(ctx).Where("account_id = ? AND id = ?", accountID, req.ConversationID).First(&conversation).Error; err != nil { + return nil, err } - if req.ContentType == "" { - message.ContentType = "text" + message := &model.Message{ + AccountID: accountID, + ConversationID: req.ConversationID, + InboxID: conversation.InboxID, + Content: req.Content, + MessageType: req.MessageType, + ContentType: req.ContentType, + SenderID: &userID, + SenderType: "user", + Private: req.Private, + SourceID: req.SourceID, + EchoID: req.EchoID, + Status: "sent", + ContentAttributes: req.ContentAttributes, } // Chatwoot: when message_type is "private_note", force Private=true and ContentType="private_note" @@ -142,6 +182,41 @@ func (s *MessageService) Create(ctx context.Context, accountID uint, userID uint return message, nil } +func normalizeMessageType(value string) string { + switch strings.ToLower(strings.TrimSpace(value)) { + case "", "1", "outgoing": + return "outgoing" + case "0", "incoming": + return "incoming" + case "2", "activity": + return "activity" + case "3", "template": + return "template" + case "private_note": + return "private_note" + default: + return value + } +} + +func validMessageType(value string) bool { + switch value { + case "incoming", "outgoing", "activity", "template", "private_note": + return true + default: + return false + } +} + +func validContentType(value string) bool { + switch value { + case "text", "input_text", "input_email", "input_phone", "select", "card", "private_note", "input_csat", "file", "image", "audio", "video", "voice_call": + return true + default: + return false + } +} + // UpdateMessageRequest is the DTO for updating a message. type UpdateMessageRequest struct { Content string `json:"content,omitempty" validate:"omitempty,min=1"` diff --git a/migrations/000018_add_conversation_message_parity_fields.down.sql b/migrations/000018_add_conversation_message_parity_fields.down.sql new file mode 100644 index 00000000..be35a355 --- /dev/null +++ b/migrations/000018_add_conversation_message_parity_fields.down.sql @@ -0,0 +1,26 @@ +DROP INDEX IF EXISTS idx_messages_echo_id; +DROP INDEX IF EXISTS idx_conversations_account_display_id; + +ALTER TABLE messages DROP COLUMN IF EXISTS external_source_ids; +ALTER TABLE messages DROP COLUMN IF EXISTS additional_attributes; +ALTER TABLE messages DROP COLUMN IF EXISTS content_attributes; +ALTER TABLE messages DROP COLUMN IF EXISTS echo_id; +ALTER TABLE messages DROP COLUMN IF EXISTS status; + +ALTER TABLE conversations DROP COLUMN IF EXISTS resumed_at; +ALTER TABLE conversations DROP COLUMN IF EXISTS resolved_at; +ALTER TABLE conversations DROP COLUMN IF EXISTS muted; +ALTER TABLE conversations DROP COLUMN IF EXISTS first_reply_created_at; +ALTER TABLE conversations DROP COLUMN IF EXISTS last_activity_at; +ALTER TABLE conversations DROP COLUMN IF EXISTS waiting_since; +ALTER TABLE conversations DROP COLUMN IF EXISTS contact_last_seen_at; +ALTER TABLE conversations DROP COLUMN IF EXISTS assignee_last_seen_at; +ALTER TABLE conversations DROP COLUMN IF EXISTS agent_last_seen_at; +ALTER TABLE conversations DROP COLUMN IF EXISTS custom_attributes; +ALTER TABLE conversations DROP COLUMN IF EXISTS additional_attributes; +ALTER TABLE conversations DROP COLUMN IF EXISTS snoozed_until; +ALTER TABLE conversations DROP COLUMN IF EXISTS sla_policy_id; +ALTER TABLE conversations DROP COLUMN IF EXISTS campaign_id; +ALTER TABLE conversations DROP COLUMN IF EXISTS team_id; +ALTER TABLE conversations DROP COLUMN IF EXISTS contact_inbox_id; +ALTER TABLE conversations DROP COLUMN IF EXISTS display_id; diff --git a/migrations/000018_add_conversation_message_parity_fields.up.sql b/migrations/000018_add_conversation_message_parity_fields.up.sql new file mode 100644 index 00000000..df7b1cda --- /dev/null +++ b/migrations/000018_add_conversation_message_parity_fields.up.sql @@ -0,0 +1,35 @@ +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS display_id INTEGER; +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS contact_inbox_id INTEGER; +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS team_id INTEGER; +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS campaign_id INTEGER; +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS sla_policy_id INTEGER; +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS snoozed_until BIGINT; +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS additional_attributes JSONB DEFAULT '{}'; +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS custom_attributes JSONB DEFAULT '{}'; +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS agent_last_seen_at BIGINT; +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS assignee_last_seen_at BIGINT; +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS contact_last_seen_at BIGINT; +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS waiting_since BIGINT; +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS last_activity_at BIGINT; +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS first_reply_created_at BIGINT; +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS muted BOOLEAN DEFAULT FALSE; +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS resolved_at TIMESTAMP WITH TIME ZONE; +ALTER TABLE conversations ADD COLUMN IF NOT EXISTS resumed_at TIMESTAMP WITH TIME ZONE; + +ALTER TABLE messages ADD COLUMN IF NOT EXISTS status VARCHAR(50) DEFAULT 'sent'; +ALTER TABLE messages ADD COLUMN IF NOT EXISTS echo_id VARCHAR(255); +ALTER TABLE messages ADD COLUMN IF NOT EXISTS content_attributes JSONB DEFAULT '{}'; +ALTER TABLE messages ADD COLUMN IF NOT EXISTS additional_attributes JSONB DEFAULT '{}'; +ALTER TABLE messages ADD COLUMN IF NOT EXISTS external_source_ids JSONB DEFAULT '{}'; + +UPDATE conversations +SET display_id = numbered.display_id +FROM ( + SELECT id, ROW_NUMBER() OVER (PARTITION BY account_id ORDER BY id)::integer AS display_id + FROM conversations + WHERE display_id IS NULL OR display_id = 0 +) AS numbered +WHERE conversations.id = numbered.id; + +CREATE INDEX IF NOT EXISTS idx_conversations_account_display_id ON conversations(account_id, display_id) WHERE deleted_at IS NULL; +CREATE INDEX IF NOT EXISTS idx_messages_echo_id ON messages(echo_id) WHERE deleted_at IS NULL;