From 8cd564e1da1ad3b1f2dc4b6357b3c5abc6219c00 Mon Sep 17 00:00:00 2001 From: Rogee Date: Thu, 4 Jun 2026 22:43:35 +0800 Subject: [PATCH] feat(widget): support chatwoot direct upload attachments --- docs/CHATWOOT_PARITY_DEVELOPMENT_PLAN.md | 5 +- docs/parity/gochat_routes.txt | 4 +- internal/app/bootstrap.go | 2 +- internal/handler/api/v1/upload_handler.go | 29 +- .../handler/api/v1/upload_handler_test.go | 88 +++++- internal/handler/widget/widget_handler.go | 96 +++++- .../handler/widget/widget_handler_test.go | 75 +++++ internal/router/router.go | 2 + internal/service/upload_service.go | 289 ++++++++++++++++-- internal/service/widget_service.go | 71 ++++- 10 files changed, 617 insertions(+), 44 deletions(-) diff --git a/docs/CHATWOOT_PARITY_DEVELOPMENT_PLAN.md b/docs/CHATWOOT_PARITY_DEVELOPMENT_PLAN.md index 0499d378..8e0726a6 100644 --- a/docs/CHATWOOT_PARITY_DEVELOPMENT_PLAN.md +++ b/docs/CHATWOOT_PARITY_DEVELOPMENT_PLAN.md @@ -402,7 +402,7 @@ Tracking table: | P6.3 | Conversation APIs | `docs/ROUTE_GAP_ANALYSIS.md`, conversation handlers/services | Implement frontend-critical filters, assignment, status, snooze, merge, bulk actions. | Todo | | P6.4 | Message APIs | `docs/ROUTE_GAP_ANALYSIS.md`, message handlers/services | Implement create/list/delete, private notes, attachments, source attribution, events. | Todo | | P6.5 | Inbox APIs | `docs/ROUTE_GAP_ANALYSIS.md`, inbox handlers/services | Implement CRUD, assignable agents, avatar, campaigns, channel settings, reset secret. | Todo | -| P6.6 | Widget/public APIs | `docs/ROUTE_GAP_ANALYSIS.md`, widget/channel provider code, `chatwootParityStub` routes | Finish direct uploads and deeper public CSAT parity; public inbox/contact/conversation/message core flow is now handler-backed. | Doing | +| P6.6 | Widget/public APIs | `docs/ROUTE_GAP_ANALYSIS.md`, widget/channel provider code, `chatwootParityStub` routes | Finish deeper public CSAT parity; public inbox/contact/conversation/message core flow and widget direct uploads/attachments are now handler-backed. | Doing | | P6.7 | Webhook ingress | `internal/router/router.go`, `internal/handler/webhook/*`, channel providers | Replace generic placeholder with provider-specific verified ingestion and dispatch. | Todo | Widget/public subtracking: @@ -413,7 +413,7 @@ Widget/public subtracking: | P6.6b | `/api/v1/widget` campaigns, events, inbox members, labels | Chatwoot widget campaigns/events/labels controllers and serializers | Done | | P6.6c | `/api/v1/widget` message update, transcript, `contact/set_user`, Dyte participant | Chatwoot widget message/contact/transcript/integration behavior | Done | | P6.6d | `/public/api/v1/inboxes` contact/conversation/message core flow | `reference/chatwoot/app/controllers/public/api/v1/inboxes/*` and matching jbuilder views | Done | -| P6.6e | Widget direct uploads and attachments | Chatwoot active storage/direct upload and attachment payloads | Todo | +| P6.6e | Widget direct uploads and attachments | Chatwoot active storage/direct upload and attachment payloads | Done | | P6.6f | Public CSAT deep behavior | Chatwoot CSAT survey controller/listener and message locking rules | Todo | ## Phase 7: Verification Harness @@ -478,3 +478,4 @@ Verification milestone gates: - 2026-06-04: Continued Phase 3/6 widget behavior parity. Replaced more `/api/v1/widget` stubs with handlers for `campaigns`, `events`, `inbox_members`, `labels`, and label removal. Inbox member payload now follows Chatwoot `{payload: [...]}` shape; campaigns return enabled inbox campaigns with trigger rules; events validate website/contact token context and return `204`; labels mutate the latest widget conversation only when the label exists in the account. Added focused widget handler coverage for available agents, events, and label add/remove. Remaining widget stubs: message update, transcript, `contact/set_user`, and Dyte participant integration; public inbox/contact/conversation/message routes are still placeholder-backed. - 2026-06-04: Completed the remaining `/api/v1/widget` stub burn-down. Message update now persists submitted email/form values and identifies the contact; `contact/set_user` validates identifier HMAC, supports verified contact identification, and returns `widget_auth_token` when the contact context changes; conversation transcript returns Chatwoot-compatible status behavior around missing conversations; Dyte participant endpoint validates integration messages and returns a meeting token payload. Added model/repository support for contact inbox HMAC verification and identifier lookup. Public inbox/contact/conversation/message routes remain the next P6.6 placeholder group. - 2026-06-04: Replaced `/public/api/v1/inboxes` contact/conversation/message placeholders with real Chatwoot public API handlers. API inboxes now resolve through `Channel::Api` identifiers, public contacts create/update by `source_id` with optional identifier HMAC verification, public conversations enforce the same verified-contact visibility split, and public messages support create/list/update submitted values. Added focused public API handler flow coverage. Verified `go test ./...`, regenerated `docs/parity/gochat_routes.txt` (`TOTAL: 791`), and regenerated `docs/parity/route_parity.md` (`251 exact, 0 missing`). Remaining P6.6 work is direct uploads/attachments and deeper public CSAT behavior. +- 2026-06-04: Completed P6.6e widget direct upload/attachment parity for the reused Chatwoot widget frontend. `/api/v1/widget/direct_uploads` now accepts ActiveStorage metadata with `website_token` + `X-Auth-Token`, returns the raw `signed_id/direct_upload` blob shape expected by `DirectUpload`, supports the follow-up PUT body upload, and attaches `message[attachments][]` signed IDs to incoming widget messages. Message create/list payloads now include Chatwoot-style attachment fields (`data_url`, `thumb_url`, `file_type`, extension, size). Added focused handler coverage for ActiveStorage create/PUT and multipart attachment-only message send/list. Verified focused package tests and regenerated route artifacts; route dump now reports `TOTAL: 793`, while tracked parity remains `251 exact, 0 missing`. Remaining P6.6 work is deeper public CSAT behavior. diff --git a/docs/parity/gochat_routes.txt b/docs/parity/gochat_routes.txt index 94de541a..8840ca96 100644 --- a/docs/parity/gochat_routes.txt +++ b/docs/parity/gochat_routes.txt @@ -778,6 +778,7 @@ PUT /api/v1/profile PUT /api/v1/profile/avatar PUT /api/v1/profile/set_active_account PUT /api/v1/widget/contact +PUT /api/v1/widget/direct_uploads/:upload_uuid PUT /api/v1/widget/messages/:message_id PUT /platform/api/v1/agent_bots/:id PUT /platform/api/v1/agent_bots/:id/avatar @@ -788,5 +789,6 @@ PUT /public/api/v1/csat_survey/:id PUT /public/api/v1/inboxes/:inbox_id/contacts/:contact_id PUT /public/api/v1/inboxes/:inbox_id/contacts/:contact_id/conversations/:conversation_id/messages/:message_id PUT /webhooks/:channel_type/:identifier +PUT /widget/direct_uploads/:upload_uuid TRACE /webhooks/:channel_type/:identifier -TOTAL: 791 +TOTAL: 793 diff --git a/internal/app/bootstrap.go b/internal/app/bootstrap.go index 1356d8d6..370e6702 100644 --- a/internal/app/bootstrap.go +++ b/internal/app/bootstrap.go @@ -675,7 +675,7 @@ func Bootstrap(env string) (*App, error) { // Upload: DirectUpload repo + service + handler (account-level + widget direct uploads) directUploadRepo := repository.NewDirectUploadRepo(db) - uploadService := service.NewUploadService(directUploadRepo, cfg) + uploadService := service.NewUploadService(directUploadRepo, cfg).WithWidgetAuth(inboxRepo, contactInboxRepo) uploadHandler := v1.NewUploadHandler(uploadService) // Step 9: Wire handlers (HTTP presentation layer) diff --git a/internal/handler/api/v1/upload_handler.go b/internal/handler/api/v1/upload_handler.go index a60943b2..9a278450 100644 --- a/internal/handler/api/v1/upload_handler.go +++ b/internal/handler/api/v1/upload_handler.go @@ -2,6 +2,7 @@ package v1 import ( "net/http" + "strings" "github.com/gin-gonic/gin" @@ -48,6 +49,23 @@ func (h *UploadHandler) Upload(c *gin.Context) { // DirectUpload handles POST /api/v1/widget/direct_uploads — widget direct file upload. // Reference: Chatwoot POST /widget/direct_uploads func (h *UploadHandler) DirectUpload(c *gin.Context) { + if strings.Contains(c.GetHeader("Content-Type"), "application/json") { + var req service.ActiveStorageDirectUploadRequest + if err := c.ShouldBindJSON(&req); err != nil { + response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrValidation, "invalid direct upload metadata") + return + } + req.WebsiteToken = c.Query("website_token") + req.AuthToken = c.GetHeader("X-Auth-Token") + result, svcErr := h.svc.CreateWidgetDirectUpload(c.Request.Context(), req) + if svcErr != nil { + handleServiceError(c, svcErr) + return + } + c.JSON(http.StatusOK, result) + return + } + fileHeader, err := c.FormFile("file") if err != nil { response.AbortWithStatusError(c, http.StatusBadRequest, response.ErrValidation, "file is required") @@ -65,6 +83,15 @@ func (h *UploadHandler) DirectUpload(c *gin.Context) { response.OK(c, result) } +func (h *UploadHandler) CompleteWidgetDirectUpload(c *gin.Context) { + result, svcErr := h.svc.CompleteWidgetDirectUpload(c.Request.Context(), c.Param("upload_uuid"), c.Request.Body) + if svcErr != nil { + handleServiceError(c, svcErr) + return + } + response.OK(c, result) +} + // AccountDirectUpload handles POST /api/v1/accounts/:id/direct_uploads — account-level staged upload. // Returns a blob/UUID for later attachment to messages. // Reference: Chatwoot POST /api/v1/accounts/:account_id/direct_uploads @@ -90,4 +117,4 @@ func (h *UploadHandler) AccountDirectUpload(c *gin.Context) { } response.OK(c, result) -} \ No newline at end of file +} diff --git a/internal/handler/api/v1/upload_handler_test.go b/internal/handler/api/v1/upload_handler_test.go index cee23760..c39e075e 100644 --- a/internal/handler/api/v1/upload_handler_test.go +++ b/internal/handler/api/v1/upload_handler_test.go @@ -6,11 +6,20 @@ import ( "mime/multipart" "net/http" "net/http/httptest" + "os" + "path/filepath" "testing" "github.com/gin-gonic/gin" "github.com/stretchr/testify/assert" + "github.com/stretchr/testify/require" + "gorm.io/driver/sqlite" + "gorm.io/gorm" + "gorm.io/gorm/logger" + "github.com/gochat/gochat/internal/config" + "github.com/gochat/gochat/internal/model" + "github.com/gochat/gochat/internal/repository" "github.com/gochat/gochat/internal/service" ) @@ -32,6 +41,11 @@ func setupUploadHandlerRouter(h *UploadHandler) *gin.Engine { // Widget direct upload route widget := r.Group("/widget") widget.POST("/direct_uploads", h.DirectUpload) + widget.PUT("/direct_uploads/:upload_uuid", h.CompleteWidgetDirectUpload) + + chatwootWidget := r.Group("/api/v1/widget") + chatwootWidget.POST("/direct_uploads", h.DirectUpload) + chatwootWidget.PUT("/direct_uploads/:upload_uuid", h.CompleteWidgetDirectUpload) return r } @@ -71,7 +85,7 @@ func TestUploadHandler_DirectUpload_NoFile(t *testing.T) { r := setupUploadHandlerRouter(h) req, _ := http.NewRequest("POST", "/widget/direct_uploads", nil) - req.Header.Set("Content-Type", "application/json") + req.Header.Set("Content-Type", "multipart/form-data") w := httptest.NewRecorder() r.ServeHTTP(w, req) @@ -79,6 +93,76 @@ func TestUploadHandler_DirectUpload_NoFile(t *testing.T) { assert.Equal(t, http.StatusBadRequest, w.Code) } +func TestUploadHandler_WidgetActiveStorageDirectUploadFlow(t *testing.T) { + gin.SetMode(gin.TestMode) + tmpDir := t.TempDir() + db, err := gorm.Open(sqlite.Open("file::memory:"), &gorm.Config{Logger: logger.Default.LogMode(logger.Silent)}) + require.NoError(t, err) + require.NoError(t, db.AutoMigrate( + &model.Account{}, + &model.Inbox{}, + &model.Contact{}, + &model.ContactInbox{}, + &model.DirectUpload{}, + )) + + account := &model.Account{Name: "Widget Upload Org", Status: "active"} + require.NoError(t, db.Create(account).Error) + channelConfig, err := json.Marshal(service.WebWidgetConfig{WebsiteToken: "upload_ws_token"}) + require.NoError(t, err) + inbox := &model.Inbox{AccountID: account.ID, Name: "Upload Widget", ChannelType: "web_widget", Enabled: true, ChannelConfig: string(channelConfig)} + require.NoError(t, db.Create(inbox).Error) + contact := &model.Contact{AccountID: account.ID, Name: "Uploader"} + require.NoError(t, db.Create(contact).Error) + contactInbox := &model.ContactInbox{ContactID: contact.ID, InboxID: inbox.ID, PubsubToken: "upload_pubsub_token"} + require.NoError(t, db.Create(contactInbox).Error) + + uploadSvc := service.NewUploadService(repository.NewDirectUploadRepo(db), &config.Config{ + Storage: config.StorageConfig{LocalPath: tmpDir, MaxFileSize: 50 << 20}, + }).WithWidgetAuth(repository.NewInboxRepo(db), repository.NewContactInboxRepo(db)) + router := setupUploadHandlerRouter(NewUploadHandler(uploadSvc)) + + metadataBody, err := json.Marshal(map[string]any{ + "blob": map[string]any{ + "filename": "visitor.png", + "byte_size": 11, + "checksum": "checksum-token", + "content_type": "image/png", + "metadata": map[string]any{"identified": true}, + }, + }) + require.NoError(t, err) + + wCreate := httptest.NewRecorder() + reqCreate, _ := http.NewRequest("POST", "/api/v1/widget/direct_uploads?website_token=upload_ws_token", bytes.NewReader(metadataBody)) + reqCreate.Header.Set("Content-Type", "application/json") + reqCreate.Header.Set("X-Auth-Token", "upload_pubsub_token") + router.ServeHTTP(wCreate, reqCreate) + require.Equal(t, http.StatusOK, wCreate.Code) + + var createResp map[string]any + require.NoError(t, json.Unmarshal(wCreate.Body.Bytes(), &createResp)) + signedID, ok := createResp["signed_id"].(string) + require.True(t, ok) + require.NotEmpty(t, signedID) + assert.Equal(t, "visitor.png", createResp["filename"]) + directUpload := createResp["direct_upload"].(map[string]any) + assert.Equal(t, "/api/v1/widget/direct_uploads/"+signedID, directUpload["url"]) + + wPut := httptest.NewRecorder() + reqPut, _ := http.NewRequest("PUT", directUpload["url"].(string), bytes.NewReader([]byte("hello image"))) + reqPut.Header.Set("Content-Type", "image/png") + router.ServeHTTP(wPut, reqPut) + require.Equal(t, http.StatusOK, wPut.Code) + + var upload model.DirectUpload + require.NoError(t, db.Where("upload_uuid = ?", signedID).First(&upload).Error) + assert.Equal(t, account.ID, upload.AccountID) + storedBytes, err := os.ReadFile(filepath.Join(tmpDir, "widget_direct", signedID+".png")) + require.NoError(t, err) + assert.Equal(t, []byte("hello image"), storedBytes) +} + func TestUploadHandler_AccountDirectUpload_NoFile(t *testing.T) { // Create handler with nil service — we only test validation before service call h := &UploadHandler{svc: nil} @@ -163,4 +247,4 @@ func TestUploadHandler_ResponseStructure(t *testing.T) { assert.Equal(t, "/uploads/account/1/test.png", parsed["file_url"]) assert.Equal(t, "/uploads/account/1/test.png", parsed["thumb_url"]) assert.Equal(t, "pending", parsed["status"]) -} \ No newline at end of file +} diff --git a/internal/handler/widget/widget_handler.go b/internal/handler/widget/widget_handler.go index b220d3f3..c3e9c8e1 100644 --- a/internal/handler/widget/widget_handler.go +++ b/internal/handler/widget/widget_handler.go @@ -113,8 +113,11 @@ func (h *WidgetHandler) Config(c *gin.Context) { "welcome_tagline": resp.WidgetConfig.WelcomeTagline, "website_name": resp.InboxName, }, - "contact": contact, - "global_config": gin.H{}, + "contact": contact, + "global_config": gin.H{ + "directUploadsEnabled": true, + "maximumFileUploadSize": 40, + }, }) } @@ -149,7 +152,11 @@ func (h *WidgetHandler) SendMessage(c *gin.Context) { c.JSON(http.StatusOK, resp) return } - c.JSON(http.StatusOK, widgetMessagePayload(resp.Message, resp.ConversationID)) + payload := widgetMessagePayload(resp.Message, resp.ConversationID) + if len(resp.Attachments) > 0 { + payload["attachments"] = widgetAttachmentPayloads(resp.Attachments) + } + c.JSON(http.StatusOK, payload) } func (h *WidgetHandler) UpdateMessage(c *gin.Context) { @@ -216,7 +223,11 @@ func (h *WidgetHandler) GetLatestMessages(c *gin.Context) { payload := make([]gin.H, 0, len(messages)) for _, msg := range messages { - payload = append(payload, widgetMessagePayload(msg, msg.ConversationID)) + messagePayload := widgetMessagePayload(msg, msg.ConversationID) + if attachments, err := h.widgetService.GetMessageAttachments(c.Request.Context(), msg.ID); err == nil && len(attachments) > 0 { + messagePayload["attachments"] = widgetAttachmentPayloads(attachments) + } + payload = append(payload, messagePayload) } meta := gin.H{"total": total, "offset": offset, "limit": limit} if conversation != nil && conversation.ContactLastSeenAt != nil { @@ -988,12 +999,40 @@ func widgetTokenFromRequest(c *gin.Context) string { } func bindWidgetSendMessageRequest(c *gin.Context) (service.WidgetSendMessageRequest, error) { + if strings.Contains(c.GetHeader("Content-Type"), "multipart/form-data") { + if err := c.Request.ParseMultipartForm(32 << 20); err != nil { + return service.WidgetSendMessageRequest{}, err + } + form := c.Request.MultipartForm + content := firstFormValue(form.Value, "content", "message[content]") + contentType := firstFormValue(form.Value, "content_type", "message[content_type]") + var conversationID *uint + if rawID := firstFormValue(form.Value, "conversation_id", "message[conversation_id]"); rawID != "" { + if id, err := strconv.ParseUint(rawID, 10, 64); err == nil && id > 0 { + value := uint(id) + conversationID = &value + } + } + attachments := form.Value["message[attachments][]"] + if len(attachments) == 0 { + attachments = form.Value["attachments[]"] + } + return service.WidgetSendMessageRequest{ + Content: content, + ContentType: contentType, + ConversationID: conversationID, + AttachmentIDs: attachments, + }, nil + } + var body struct { - Content string `json:"content"` - ContentType string `json:"content_type"` - ConversationID *uint `json:"conversation_id"` + Content string `json:"content"` + ContentType string `json:"content_type"` + ConversationID *uint `json:"conversation_id"` + Attachments []string `json:"attachments"` Message struct { - Content string `json:"content"` + Content string `json:"content"` + Attachments []string `json:"attachments"` } `json:"message"` } if err := c.ShouldBindJSON(&body); err != nil { @@ -1003,13 +1042,27 @@ func bindWidgetSendMessageRequest(c *gin.Context) (service.WidgetSendMessageRequ if content == "" { content = body.Message.Content } + attachments := body.Attachments + if len(attachments) == 0 { + attachments = body.Message.Attachments + } return service.WidgetSendMessageRequest{ Content: content, ContentType: body.ContentType, ConversationID: body.ConversationID, + AttachmentIDs: attachments, }, nil } +func firstFormValue(values map[string][]string, keys ...string) string { + for _, key := range keys { + if list := values[key]; len(list) > 0 { + return list[0] + } + } + return "" +} + func widgetMessagePayload(message model.Message, conversationID uint) gin.H { return gin.H{ "id": message.ID, @@ -1025,6 +1078,33 @@ func widgetMessagePayload(message model.Message, conversationID uint) gin.H { } } +func widgetAttachmentPayloads(attachments []model.Attachment) []gin.H { + payload := make([]gin.H, 0, len(attachments)) + for _, attachment := range attachments { + payload = append(payload, gin.H{ + "id": attachment.ID, + "message_id": attachment.MessageID, + "thumb_url": attachment.ThumbURL, + "data_url": attachment.FileURL, + "file_size": attachment.FileSize, + "file_type": attachment.FileType, + "extension": strings.TrimPrefix(strings.ToLower(attachmentExtension(attachment.FileName)), "."), + "width": attachment.Width, + "height": attachment.Height, + "created_at": attachment.CreatedAt.Unix(), + }) + } + return payload +} + +func attachmentExtension(filename string) string { + idx := strings.LastIndex(filename, ".") + if idx == -1 { + return "" + } + return filename[idx:] +} + func widgetConversationPayload(conversation model.Conversation) gin.H { return gin.H{ "id": conversation.ID, diff --git a/internal/handler/widget/widget_handler_test.go b/internal/handler/widget/widget_handler_test.go index d5372a35..1b757f14 100644 --- a/internal/handler/widget/widget_handler_test.go +++ b/internal/handler/widget/widget_handler_test.go @@ -7,10 +7,12 @@ import ( "crypto/sha256" "encoding/hex" "encoding/json" + "mime/multipart" "net/http" "net/http/httptest" "strconv" "testing" + "time" "github.com/gin-gonic/gin" "github.com/stretchr/testify/assert" @@ -55,6 +57,8 @@ func setupWidgetHandlerTest(t *testing.T) (*gorm.DB, *gin.Engine, *WidgetHandler &model.ContactInbox{}, &model.Conversation{}, &model.Message{}, + &model.Attachment{}, + &model.DirectUpload{}, &model.WidgetThemeConfig{}, &model.PreChatForm{}, &model.WidgetFileUpload{}, @@ -353,6 +357,77 @@ func TestWidgetHandler_ChatwootMessages_AuthTokenAndNestedPayload(t *testing.T) require.Equal(t, http.StatusOK, wContact.Code) } +func TestWidgetHandler_ChatwootMessageDirectUploadAttachment(t *testing.T) { + db, router, _ := setupWidgetHandlerTest(t) + account, _ := seedWidgetHandlerData(t, db) + + wConfig := httptest.NewRecorder() + reqConfig, _ := http.NewRequest("POST", "/api/v1/widget/config?website_token=handler_ws_token_123", nil) + router.ServeHTTP(wConfig, reqConfig) + require.Equal(t, http.StatusOK, wConfig.Code) + + var configResp map[string]interface{} + require.NoError(t, json.Unmarshal(wConfig.Body.Bytes(), &configResp)) + authToken := configResp["contact"].(map[string]interface{})["pubsub_token"].(string) + + upload := &model.DirectUpload{ + UploadUUID: "signed-widget-upload-1", + AccountID: account.ID, + Status: model.DirectUploadStatusPending, + Source: model.DirectUploadSourceWidget, + OriginalName: "screenshot.png", + FileType: "image", + MimeType: "image/png", + FileSize: 12, + FileURL: "/uploads/widget_direct/signed-widget-upload-1.png", + ThumbURL: "/uploads/widget_direct/signed-widget-upload-1.png", + ExpiresAt: time.Now().Add(time.Hour), + } + require.NoError(t, db.Create(upload).Error) + + body := &bytes.Buffer{} + writer := multipart.NewWriter(body) + require.NoError(t, writer.WriteField("message[attachments][]", upload.UploadUUID)) + require.NoError(t, writer.Close()) + + wMessage := httptest.NewRecorder() + reqMessage, _ := http.NewRequest("POST", "/api/v1/widget/messages", body) + reqMessage.Header.Set("Content-Type", writer.FormDataContentType()) + reqMessage.Header.Set("X-Auth-Token", authToken) + router.ServeHTTP(wMessage, reqMessage) + require.Equal(t, http.StatusOK, wMessage.Code) + + var messageResp map[string]interface{} + require.NoError(t, json.Unmarshal(wMessage.Body.Bytes(), &messageResp)) + assert.Empty(t, messageResp["content"]) + attachments := messageResp["attachments"].([]interface{}) + require.Len(t, attachments, 1) + attachmentPayload := attachments[0].(map[string]interface{}) + assert.Equal(t, "/uploads/widget_direct/signed-widget-upload-1.png", attachmentPayload["data_url"]) + assert.Equal(t, "image", attachmentPayload["file_type"]) + + var attachment model.Attachment + require.NoError(t, db.Where("file_name = ?", "screenshot.png").First(&attachment).Error) + assert.Equal(t, upload.FileURL, attachment.FileURL) + require.NoError(t, db.First(upload, upload.ID).Error) + assert.Equal(t, model.DirectUploadStatusCompleted, upload.Status) + + wIndex := httptest.NewRecorder() + reqIndex, _ := http.NewRequest("GET", "/api/v1/widget/messages", nil) + reqIndex.Header.Set("X-Auth-Token", authToken) + router.ServeHTTP(wIndex, reqIndex) + require.Equal(t, http.StatusOK, wIndex.Code) + + var indexResp map[string]interface{} + require.NoError(t, json.Unmarshal(wIndex.Body.Bytes(), &indexResp)) + payload := indexResp["payload"].([]interface{}) + require.Len(t, payload, 1) + indexedMessage := payload[0].(map[string]interface{}) + indexedAttachments := indexedMessage["attachments"].([]interface{}) + require.Len(t, indexedAttachments, 1) + assert.Equal(t, "/uploads/widget_direct/signed-widget-upload-1.png", indexedAttachments[0].(map[string]interface{})["data_url"]) +} + func TestWidgetHandler_ChatwootMessageUpdate_SubmitsEmail(t *testing.T) { db, router, _ := setupWidgetHandlerTest(t) _, inbox := seedWidgetHandlerData(t, db) diff --git a/internal/router/router.go b/internal/router/router.go index b4ce4d40..e65ff2dc 100644 --- a/internal/router/router.go +++ b/internal/router/router.go @@ -227,6 +227,7 @@ func RegisterRoutes( // Widget direct file upload — visitor uploads file before conversation starts // Reference: Chatwoot POST /widget/direct_uploads widget.POST("/direct_uploads", handlers.Upload.DirectUpload) + widget.PUT("/direct_uploads/:upload_uuid", handlers.Upload.CompleteWidgetDirectUpload) // Chatwoot widget API routes — public + CORS for reused Chatwoot frontend/widget. // Reference: Chatwoot namespace :api/:v1/:widget at /api/v1/widget/*. @@ -1702,6 +1703,7 @@ func registerPlatformTokenRoutes(g *gin.RouterGroup, h *Handlers) { // Widget behavior is backed by Chatwoot-compatible handlers; public inbox APIs are tracked separately. func registerChatwootWidgetRoutes(g *gin.RouterGroup, h *Handlers) { g.POST("/direct_uploads", h.Upload.DirectUpload) + g.PUT("/direct_uploads/:upload_uuid", h.Upload.CompleteWidgetDirectUpload) g.POST("/config", h.Widget.Config) g.GET("/campaigns", h.Widget.ListCampaigns) g.POST("/events", h.Widget.CreateEvent) diff --git a/internal/service/upload_service.go b/internal/service/upload_service.go index 773b7234..598f4c4e 100644 --- a/internal/service/upload_service.go +++ b/internal/service/upload_service.go @@ -2,6 +2,9 @@ package service import ( "context" + "crypto/rand" + "encoding/hex" + "encoding/json" "errors" "fmt" "io" @@ -22,6 +25,8 @@ import ( // UploadService handles file uploads for both account-level and widget direct uploads. type UploadService struct { directUploadRepo *repository.DirectUploadRepo + inboxRepo *repository.InboxRepo + contactInboxRepo *repository.ContactInboxRepo cfg *config.Config } @@ -33,6 +38,14 @@ func NewUploadService(directUploadRepo *repository.DirectUploadRepo, cfg *config } } +// WithWidgetAuth wires the widget session repositories used by Chatwoot's +// website_token + X-Auth-Token direct upload guard. +func (s *UploadService) WithWidgetAuth(inboxRepo *repository.InboxRepo, contactInboxRepo *repository.ContactInboxRepo) *UploadService { + s.inboxRepo = inboxRepo + s.contactInboxRepo = contactInboxRepo + return s +} + // --- DTOs --- // AccountUploadRequest is the DTO for account-level file upload. @@ -51,6 +64,39 @@ type AccountDirectUploadRequest struct { FileHeader *multipart.FileHeader `json:"-"` } +type ActiveStorageDirectUploadRequest struct { + WebsiteToken string `json:"-"` + AuthToken string `json:"-"` + Blob ActiveStorageBlobParams `json:"blob"` +} + +type ActiveStorageBlobParams struct { + Filename string `json:"filename"` + ByteSize int64 `json:"byte_size"` + Checksum string `json:"checksum"` + ContentType string `json:"content_type"` + Metadata map[string]any `json:"metadata"` +} + +type ActiveStorageDirectUploadResponse struct { + ID uint `json:"id"` + Key string `json:"key"` + Filename string `json:"filename"` + ContentType string `json:"content_type"` + Metadata map[string]any `json:"metadata"` + ServiceName string `json:"service_name"` + ByteSize int64 `json:"byte_size"` + Checksum string `json:"checksum"` + CreatedAt time.Time `json:"created_at"` + SignedID string `json:"signed_id"` + DirectUpload ActiveStorageUploadURL `json:"direct_upload"` +} + +type ActiveStorageUploadURL struct { + URL string `json:"url"` + Headers map[string]string `json:"headers"` +} + // UploadResponse is the unified response DTO for upload endpoints. type UploadResponse struct { UploadID uint `json:"upload_id"` @@ -110,6 +156,144 @@ func (s *UploadService) WidgetDirectUpload(ctx context.Context, req WidgetDirect return s.processUpload(ctx, 0, req.FileHeader, model.DirectUploadSourceWidget) } +func (s *UploadService) CreateWidgetDirectUpload(ctx context.Context, req ActiveStorageDirectUploadRequest) (*ActiveStorageDirectUploadResponse, error) { + accountID, err := s.validateWidgetUploadSession(ctx, req.WebsiteToken, req.AuthToken) + if err != nil { + return nil, err + } + if req.Blob.Filename == "" { + return nil, errors.New("filename is required") + } + if req.Blob.ByteSize <= 0 { + return nil, errors.New("byte_size is required") + } + mimeType := req.Blob.ContentType + if mimeType == "" || mimeType == "application/octet-stream" { + mimeType = detectUploadMIMEFromFilename(req.Blob.Filename) + } + fileCategory := categorizeUploadMIME(mimeType) + if fileCategory == "" { + return nil, fmt.Errorf("unsupported file type: %s", mimeType) + } + if !isUploadMIMEAllowed(fileCategory, mimeType) { + return nil, fmt.Errorf("MIME type %s is not allowed for category %s", mimeType, fileCategory) + } + maxSize := model.WidgetUploadMaxSizeByType[fileCategory] + if maxSize == 0 { + maxSize = int64(s.cfg.Storage.MaxFileSize) + } + if req.Blob.ByteSize > maxSize { + return nil, fmt.Errorf("file size %d exceeds maximum %d for type %s", req.Blob.ByteSize, maxSize, fileCategory) + } + + uploadUUID := uuid.New().String() + fileURL, thumbURL := s.directUploadURL(model.DirectUploadSourceWidget, 0, uploadUUID, req.Blob.Filename) + metadata := req.Blob.Metadata + if metadata == nil { + metadata = map[string]any{} + } + metadata["checksum"] = req.Blob.Checksum + metadata["active_storage_key"] = randomStorageKey() + metadataJSON, _ := json.Marshal(metadata) + + upload := &model.DirectUpload{ + UploadUUID: uploadUUID, + AccountID: accountID, + Status: model.DirectUploadStatusPending, + Source: model.DirectUploadSourceWidget, + OriginalName: req.Blob.Filename, + FileType: fileCategory, + MimeType: mimeType, + FileSize: req.Blob.ByteSize, + FileURL: fileURL, + ThumbURL: thumbURL, + Metadata: metadataJSON, + ExpiresAt: time.Now().Add(24 * time.Hour), + } + if err := s.directUploadRepo.Create(ctx, upload); err != nil { + return nil, fmt.Errorf("failed to create direct upload: %w", err) + } + + return &ActiveStorageDirectUploadResponse{ + ID: upload.ID, + Key: fmt.Sprint(metadata["active_storage_key"]), + Filename: upload.OriginalName, + ContentType: upload.MimeType, + Metadata: metadata, + ServiceName: "gochat_local", + ByteSize: upload.FileSize, + Checksum: req.Blob.Checksum, + CreatedAt: upload.CreatedAt, + SignedID: upload.UploadUUID, + DirectUpload: ActiveStorageUploadURL{ + URL: "/api/v1/widget/direct_uploads/" + upload.UploadUUID, + Headers: map[string]string{ + "Content-Type": upload.MimeType, + }, + }, + }, nil +} + +func (s *UploadService) validateWidgetUploadSession(ctx context.Context, websiteToken, authToken string) (uint, error) { + if websiteToken == "" { + return 0, errors.New("website_token is required") + } + if authToken == "" { + return 0, errors.New("widget auth token is required") + } + if s.inboxRepo == nil || s.contactInboxRepo == nil { + return 0, nil + } + inbox, err := s.inboxRepo.FindByWebsiteToken(ctx, websiteToken) + if err != nil { + return 0, fmt.Errorf("invalid website_token: %w", err) + } + if !inbox.Enabled { + return 0, errors.New("inbox is disabled") + } + contactInbox, err := s.contactInboxRepo.FindByPubsubToken(ctx, authToken) + if err != nil { + return 0, fmt.Errorf("invalid widget auth token: %w", err) + } + if contactInbox.InboxID != inbox.ID { + return 0, errors.New("widget auth token does not belong to this inbox") + } + return inbox.AccountID, nil +} + +func (s *UploadService) CompleteWidgetDirectUpload(ctx context.Context, uploadUUID string, body io.Reader) (*UploadResponse, error) { + if uploadUUID == "" { + return nil, errors.New("upload_uuid is required") + } + upload, err := s.directUploadRepo.FindByUUID(ctx, uploadUUID) + if err != nil { + return nil, fmt.Errorf("direct upload not found: %w", err) + } + if upload.Source != model.DirectUploadSourceWidget { + return nil, errors.New("direct upload source mismatch") + } + if time.Now().After(upload.ExpiresAt) { + upload.Status = model.DirectUploadStatusExpired + _ = s.directUploadRepo.Update(ctx, upload) + return nil, errors.New("direct upload has expired") + } + if err := s.saveReaderToDisk(upload.FileURL, body); err != nil { + return nil, fmt.Errorf("failed to save direct upload: %w", err) + } + return &UploadResponse{ + UploadID: upload.ID, + UploadUUID: upload.UploadUUID, + OriginalName: upload.OriginalName, + FileType: upload.FileType, + MimeType: upload.MimeType, + FileSize: upload.FileSize, + FileURL: upload.FileURL, + ThumbURL: upload.ThumbURL, + Status: string(upload.Status), + ExpiresAt: upload.ExpiresAt, + }, nil +} + // --- Internal helpers --- func (s *UploadService) processUpload(ctx context.Context, accountID uint, fileHeader *multipart.FileHeader, source model.DirectUploadSource) (*UploadResponse, error) { @@ -152,17 +336,17 @@ func (s *UploadService) processUpload(ctx context.Context, accountID uint, fileH expiryDuration := 24 * time.Hour uploadUUID := uuid.New().String() upload := &model.DirectUpload{ - UploadUUID: uploadUUID, - AccountID: accountID, - Status: model.DirectUploadStatusPending, - Source: source, + UploadUUID: uploadUUID, + AccountID: accountID, + Status: model.DirectUploadStatusPending, + Source: source, OriginalName: fileHeader.Filename, - FileType: fileCategory, - MimeType: mimeType, - FileSize: fileHeader.Size, - FileURL: fileURL, - ThumbURL: thumbURL, - ExpiresAt: time.Now().Add(expiryDuration), + FileType: fileCategory, + MimeType: mimeType, + FileSize: fileHeader.Size, + FileURL: fileURL, + ThumbURL: thumbURL, + ExpiresAt: time.Now().Add(expiryDuration), } if err := s.directUploadRepo.Create(ctx, upload); err != nil { @@ -194,19 +378,7 @@ func (s *UploadService) saveFileToDisk(accountID uint, source model.DirectUpload localPath = "./uploads" } - // Determine subdirectory based on source - subDir := "account" - if source == model.DirectUploadSourceWidget { - subDir = "widget_direct" - } - - // Build directory path: uploads/// - dirPath := filepath.Join(localPath, subDir) - if accountID > 0 { - dirPath = filepath.Join(dirPath, fmt.Sprintf("%d", accountID)) - } - - // Create directory if it doesn't exist + dirPath := s.uploadDir(source, accountID) if err := os.MkdirAll(dirPath, 0755); err != nil { return "", "", fmt.Errorf("failed to create upload directory: %w", err) } @@ -237,8 +409,7 @@ func (s *UploadService) saveFileToDisk(accountID uint, source model.DirectUpload return "", "", fmt.Errorf("failed to copy file content: %w", err) } - // Build relative URL path - fileURL := fmt.Sprintf("/uploads/%s/%s/%s", subDir, fmt.Sprintf("%d", accountID), fileName) + fileURL := s.uploadURL(source, accountID, fileName) thumbURL := "" // For images, we reference the same path (thumbnail generation can be added later) @@ -250,6 +421,72 @@ func (s *UploadService) saveFileToDisk(accountID uint, source model.DirectUpload return fileURL, thumbURL, nil } +func (s *UploadService) saveReaderToDisk(fileURL string, body io.Reader) error { + localPath := s.cfg.Storage.LocalPath + if localPath == "" { + localPath = "./uploads" + } + relative := strings.TrimPrefix(fileURL, "/uploads/") + fullPath := filepath.Join(localPath, relative) + if err := os.MkdirAll(filepath.Dir(fullPath), 0755); err != nil { + return err + } + dst, err := os.Create(fullPath) + if err != nil { + return err + } + defer dst.Close() + _, err = io.Copy(dst, body) + return err +} + +func (s *UploadService) directUploadURL(source model.DirectUploadSource, accountID uint, uploadUUID, filename string) (string, string) { + ext := filepath.Ext(filename) + fileName := uploadUUID + ext + fileURL := s.uploadURL(source, accountID, fileName) + thumbURL := "" + if strings.HasPrefix(detectUploadMIMEFromFilename(filename), "image/") { + thumbURL = fileURL + } + return fileURL, thumbURL +} + +func (s *UploadService) uploadDir(source model.DirectUploadSource, accountID uint) string { + localPath := s.cfg.Storage.LocalPath + if localPath == "" { + localPath = "./uploads" + } + subDir := uploadSubDir(source) + parts := []string{localPath, subDir} + if accountID > 0 { + parts = append(parts, fmt.Sprintf("%d", accountID)) + } + return filepath.Join(parts...) +} + +func (s *UploadService) uploadURL(source model.DirectUploadSource, accountID uint, fileName string) string { + subDir := uploadSubDir(source) + if accountID > 0 { + return fmt.Sprintf("/uploads/%s/%d/%s", subDir, accountID, fileName) + } + return fmt.Sprintf("/uploads/%s/%s", subDir, fileName) +} + +func uploadSubDir(source model.DirectUploadSource) string { + if source == model.DirectUploadSourceWidget { + return "widget_direct" + } + return "account" +} + +func randomStorageKey() string { + buf := make([]byte, 16) + if _, err := rand.Read(buf); err != nil { + return uuid.New().String() + } + return hex.EncodeToString(buf) +} + // CleanupExpiredUploads removes expired direct upload records and their files. func (s *UploadService) CleanupExpiredUploads(ctx context.Context) (int64, error) { count, err := s.directUploadRepo.BatchDeleteExpired(ctx, time.Now()) @@ -343,4 +580,4 @@ func isUploadMIMEAllowed(category string, mimeType string) bool { return true } return false -} \ No newline at end of file +} diff --git a/internal/service/widget_service.go b/internal/service/widget_service.go index 4a83c6ae..362b9be3 100644 --- a/internal/service/widget_service.go +++ b/internal/service/widget_service.go @@ -119,12 +119,14 @@ type WidgetSendMessageRequest struct { Content string `json:"content" validate:"required"` ContentType string `json:"content_type,omitempty"` // default: text ConversationID *uint `json:"conversation_id,omitempty"` // nil → create new conversation + AttachmentIDs []string } // WidgetSendMessageResponse is returned after sending a message. type WidgetSendMessageResponse struct { - ConversationID uint `json:"conversation_id"` - Message model.Message `json:"message"` + ConversationID uint `json:"conversation_id"` + Message model.Message `json:"message"` + Attachments []model.Attachment `json:"attachments,omitempty"` } type WidgetContactUpdate struct { @@ -284,7 +286,7 @@ func (s *WidgetService) SendMessage(ctx context.Context, req WidgetSendMessageRe if req.WidgetToken == "" { return nil, errors.New("widget_token is required") } - if req.Content == "" { + if req.Content == "" && len(req.AttachmentIDs) == 0 { return nil, errors.New("content is required") } @@ -335,12 +337,18 @@ func (s *WidgetService) SendMessage(ctx context.Context, req WidgetSendMessageRe return nil, fmt.Errorf("failed to create message: %w", err) } + attachments, err := s.attachWidgetUploads(ctx, &msg, req.AttachmentIDs) + if err != nil { + return nil, err + } + applogger.L().Infof("Widget message: contact=%d conversation=%d message=%d", contactInbox.ContactID, conversation.ID, msg.ID) return &WidgetSendMessageResponse{ ConversationID: conversation.ID, Message: msg, + Attachments: attachments, }, nil } @@ -827,6 +835,12 @@ func (s *WidgetService) GetMessages(ctx context.Context, widgetToken string, con return s.messageRepo.FindByConversation(ctx, conversationID, offset, limit) } +func (s *WidgetService) GetMessageAttachments(ctx context.Context, messageID uint) ([]model.Attachment, error) { + var attachments []model.Attachment + err := s.messageRepo.DB().WithContext(ctx).Where("message_id = ?", messageID).Order("id ASC").Find(&attachments).Error + return attachments, err +} + // GetCableToken returns the pubsub_token for WebSocket connection. // Reference: Chatwoot widget SDK — fetches token for ActionCable subscription // The contact connects to /cable with pubsub_token to receive real-time events. @@ -1530,6 +1544,57 @@ func (s *WidgetService) findPublicContact(ctx context.Context, accountID uint, r return contact, nil } +func (s *WidgetService) attachWidgetUploads(ctx context.Context, message *model.Message, signedIDs []string) ([]model.Attachment, error) { + if len(signedIDs) == 0 { + return nil, nil + } + attachments := make([]model.Attachment, 0, len(signedIDs)) + for _, signedID := range signedIDs { + signedID = strings.TrimSpace(signedID) + if signedID == "" { + continue + } + var upload model.DirectUpload + if err := s.messageRepo.DB().WithContext(ctx).Where("upload_uuid = ?", signedID).First(&upload).Error; err != nil { + return nil, fmt.Errorf("direct upload not found: %w", err) + } + if upload.Source != model.DirectUploadSourceWidget { + return nil, errors.New("direct upload source mismatch") + } + if upload.Status != model.DirectUploadStatusPending { + return nil, errors.New("direct upload is not pending") + } + if upload.AccountID != 0 && upload.AccountID != message.AccountID { + return nil, errors.New("direct upload account mismatch") + } + if time.Now().After(upload.ExpiresAt) { + upload.Status = model.DirectUploadStatusExpired + _ = s.messageRepo.DB().WithContext(ctx).Save(&upload).Error + return nil, errors.New("direct upload has expired") + } + attachment := model.Attachment{ + MessageID: message.ID, + AccountID: message.AccountID, + FileType: upload.FileType, + FileURL: upload.FileURL, + ThumbURL: upload.ThumbURL, + FileSize: int(upload.FileSize), + FileName: upload.OriginalName, + Metadata: string(upload.Metadata), + } + if err := s.messageRepo.DB().WithContext(ctx).Create(&attachment).Error; err != nil { + return nil, fmt.Errorf("failed to create attachment: %w", err) + } + upload.Status = model.DirectUploadStatusCompleted + upload.AccountID = message.AccountID + if err := s.messageRepo.DB().WithContext(ctx).Save(&upload).Error; err != nil { + return nil, fmt.Errorf("failed to mark direct upload completed: %w", err) + } + attachments = append(attachments, attachment) + } + return attachments, nil +} + func (s *WidgetService) resolvePublicInbox(ctx context.Context, inboxIdentifier string) (*model.Inbox, *channelmodel.ChannelAPI, error) { if inboxIdentifier == "" { return nil, nil, errors.New("inbox identifier is required")