diff --git a/docs/parity/gochat_routes.txt b/docs/parity/gochat_routes.txt index 3c8b735f..07803cb8 100644 --- a/docs/parity/gochat_routes.txt +++ b/docs/parity/gochat_routes.txt @@ -191,6 +191,7 @@ GET /api/v1/accounts/:account_id/contacts/:contact_id/notes GET /api/v1/accounts/:account_id/contacts/:contact_id/notes/:note_id GET /api/v1/accounts/:account_id/contacts/active GET /api/v1/accounts/:account_id/contacts/export +GET /api/v1/accounts/:account_id/contacts/export/:export_id/download GET /api/v1/accounts/:account_id/contacts/search GET /api/v1/accounts/:account_id/conversations/ GET /api/v1/accounts/:account_id/conversations/:conversation_id @@ -816,4 +817,4 @@ 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 /widget/direct_uploads/:upload_uuid -TOTAL: 818 +TOTAL: 819 diff --git a/internal/app/app.go b/internal/app/app.go index c4b6b43e..fc34c695 100644 --- a/internal/app/app.go +++ b/internal/app/app.go @@ -24,12 +24,12 @@ import ( // App is the central application container holding all shared services. type App struct { - config *config.Config - reloader *config.ConfigReloader - db *gorm.DB - pubsub pubsub.PubSub - engine *gin.Engine - wsHub *ws.Hub + config *config.Config + reloader *config.ConfigReloader + db *gorm.DB + pubsub pubsub.PubSub + engine *gin.Engine + wsHub *ws.Hub notificationDeliverySvc *service.NotificationDeliveryService } @@ -237,6 +237,7 @@ func autoMigrate(db *gorm.DB) error { &model.Call{}, &model.CaptainAssistantInbox{}, &model.Contactable{}, + &model.ContactExport{}, &model.DataImport{}, &model.EmailTemplate{}, &model.MessageReaction{}, diff --git a/internal/handler/api/v1/contact_handler.go b/internal/handler/api/v1/contact_handler.go index 2bbf609d..11a784d1 100644 --- a/internal/handler/api/v1/contact_handler.go +++ b/internal/handler/api/v1/contact_handler.go @@ -1,6 +1,7 @@ package v1 import ( + "io" "net/http" "strconv" @@ -706,9 +707,41 @@ func (h *ContactHandler) ExportRequest(c *gin.Context) { return } + var req service.ContactExportRequest + if err := c.ShouldBindJSON(&req); err != nil && err != io.EOF { + c.JSON(http.StatusUnprocessableEntity, gin.H{"error": "failed to export contacts"}) + return + } + + if _, svcErr := h.svc.ExportContacts(c.Request.Context(), accountID, getUserID(c), req); svcErr != nil { + c.JSON(http.StatusUnprocessableEntity, gin.H{"error": "failed to export contacts"}) + return + } + c.Status(http.StatusOK) } +// DownloadExport streams a previously generated contacts export artifact. +func (h *ContactHandler) DownloadExport(c *gin.Context) { + accountID := parseAccountIDParam(c) + if accountID == 0 { + c.JSON(http.StatusBadRequest, gin.H{"error": "invalid account id"}) + return + } + exportID, err := parseUintParam(c, "export_id") + if err != nil || exportID == 0 { + c.JSON(http.StatusBadRequest, gin.H{"error": "invalid export id"}) + return + } + export, svcErr := h.svc.GetContactExport(c.Request.Context(), accountID, exportID) + if svcErr != nil || len(export.CSVData) == 0 { + c.JSON(http.StatusNotFound, gin.H{"error": "export not found"}) + return + } + c.Header("Content-Disposition", "attachment; filename="+export.FileName) + c.Data(http.StatusOK, export.ContentType, export.CSVData) +} + // Import uploads contacts from a CSV file. // POST /api/v1/accounts/:id/contacts/import // Reference: Chatwoot contacts#import diff --git a/internal/handler/api/v1/contact_handler_crud_test.go b/internal/handler/api/v1/contact_handler_crud_test.go index 5f432879..ff109433 100644 --- a/internal/handler/api/v1/contact_handler_crud_test.go +++ b/internal/handler/api/v1/contact_handler_crud_test.go @@ -52,7 +52,9 @@ func (s *ContactHandlerCRUDTestSuite) SetupSuite() { &model.Contact{}, &model.Tag{}, &model.ContactLabel{}, + &model.ContactExport{}, &model.DataImport{}, + &model.Notification{}, &model.Conversation{}, &model.Message{}, &model.ContactInbox{}, @@ -84,6 +86,8 @@ func (s *ContactHandlerCRUDTestSuite) SetupSuite() { s.router.Use(gin.Recovery(), s.mockAuthMiddleware()) s.router.GET("/api/v1/accounts/:id/contacts", s.handler.List) s.router.GET("/api/v1/accounts/:id/contacts/search", s.handler.Search) + s.router.POST("/api/v1/accounts/:id/contacts/export", s.handler.ExportRequest) + s.router.GET("/api/v1/accounts/:id/contacts/export/:export_id/download", s.handler.DownloadExport) s.router.POST("/api/v1/accounts/:id/contacts/import", s.handler.Import) s.router.GET("/api/v1/accounts/:id/contacts/:contact_id", s.handler.Get) s.router.POST("/api/v1/accounts/:id/contacts", s.handler.Create) @@ -123,7 +127,9 @@ func (s *ContactHandlerCRUDTestSuite) SetupTest() { s.db.Exec("DELETE FROM contact_notes") s.db.Exec("DELETE FROM contact_labels") s.db.Exec("DELETE FROM tags") + s.db.Exec("DELETE FROM contact_exports") s.db.Exec("DELETE FROM data_imports") + s.db.Exec("DELETE FROM notifications") s.db.Exec("DELETE FROM messages") s.db.Exec("DELETE FROM notes") s.db.Exec("DELETE FROM contact_inboxes") @@ -635,6 +641,63 @@ func (s *ContactHandlerCRUDTestSuite) TestImport_CreatesDataImportAndReturnsOK() s.Equal(int64(1), labelCount) } +func (s *ContactHandlerCRUDTestSuite) TestExportRequest_CreatesArtifactNotificationAndReturnsOK() { + contact := &model.Contact{AccountID: s.account.ID, Name: "Exported", Email: "exported@example.com"} + s.Require().NoError(s.db.Create(contact).Error) + tag := &model.Tag{AccountID: s.account.ID, Name: "vip"} + s.Require().NoError(s.db.Create(tag).Error) + s.Require().NoError(s.db.Create(&model.ContactLabel{AccountID: s.account.ID, ContactID: contact.ID, TagID: tag.ID}).Error) + + bodyBytes, _ := json.Marshal(map[string]any{ + "column_names": []string{"email", "labels"}, + "label": "vip", + }) + w := httptest.NewRecorder() + req, _ := http.NewRequest("POST", + fmt.Sprintf("/api/v1/accounts/%d/contacts/export", s.account.ID), + bytes.NewReader(bodyBytes)) + req.Header.Set("Content-Type", "application/json") + s.router.ServeHTTP(w, req) + + s.Equal(http.StatusOK, w.Code) + s.Empty(w.Body.String()) + + var export model.ContactExport + s.Require().NoError(s.db.Where("account_id = ?", s.account.ID).First(&export).Error) + s.Equal(string(model.DataImportStatusCompleted), export.Status) + s.Equal(1, export.RowCount) + s.Contains(string(export.CSVData), "email,labels") + s.Contains(string(export.CSVData), "exported@example.com,vip") + s.Contains(export.FileURL, fmt.Sprintf("/contacts/export/%d/download", export.ID)) + + var notification model.Notification + s.Require().NoError(s.db.Where("user_id = ? AND notification_type = ?", s.user.ID, "contacts_export_complete").First(¬ification).Error) + s.Equal("ContactExport", notification.PrimaryActorType) + s.Equal(export.ID, notification.PrimaryActorID) +} + +func (s *ContactHandlerCRUDTestSuite) TestDownloadExport_ReturnsPersistedCSVArtifact() { + export := &model.ContactExport{ + AccountID: s.account.ID, + UserID: &s.user.ID, + Status: string(model.DataImportStatusCompleted), + FileName: "contacts.csv", + ContentType: "text/csv", + CSVData: []byte("\ufeffemail\nexported@example.com\n"), + } + s.Require().NoError(s.db.Create(export).Error) + + w := httptest.NewRecorder() + req, _ := http.NewRequest("GET", + fmt.Sprintf("/api/v1/accounts/%d/contacts/export/%d/download", s.account.ID, export.ID), nil) + s.router.ServeHTTP(w, req) + + s.Equal(http.StatusOK, w.Code) + s.Equal("text/csv", w.Header().Get("Content-Type")) + s.Contains(w.Header().Get("Content-Disposition"), "contacts.csv") + s.Contains(w.Body.String(), "exported@example.com") +} + func (s *ContactHandlerCRUDTestSuite) TestLabels_UpdateListAndFilter() { bodyBytes, _ := json.Marshal(map[string]interface{}{"labels": []string{"vip", "trial"}}) diff --git a/internal/model/contact_export.go b/internal/model/contact_export.go new file mode 100644 index 00000000..b9ee33c7 --- /dev/null +++ b/internal/model/contact_export.go @@ -0,0 +1,34 @@ +package model + +import ( + "encoding/json" + "time" + + "gorm.io/gorm" +) + +// ContactExport stores the completed CSV artifact generated by contacts#export. +// Reference: Chatwoot Account::ContactsExportJob attaches contacts_export to the account. +type ContactExport struct { + ID uint `gorm:"primaryKey" json:"id"` + AccountID uint `gorm:"not null;index" json:"account_id"` + UserID *uint `gorm:"index" json:"user_id,omitempty"` + Status string `gorm:"size:50;not null;default:'pending'" json:"status"` + FileName string `gorm:"size:255;not null" json:"file_name"` + ContentType string `gorm:"size:100;not null;default:'text/csv'" json:"content_type"` + FileURL string `gorm:"size:512" json:"file_url"` + CSVData []byte `gorm:"type:bytea" json:"-"` + RowCount int `json:"row_count"` + ColumnNames json.RawMessage `gorm:"type:jsonb" json:"column_names"` + FilterParams json.RawMessage `gorm:"type:jsonb" json:"filter_params"` + Error string `gorm:"type:text" json:"error,omitempty"` + CompletedAt *time.Time `json:"completed_at,omitempty"` + CreatedAt time.Time `gorm:"autoCreateTime" json:"created_at"` + UpdatedAt time.Time `gorm:"autoUpdateTime" json:"updated_at"` + DeletedAt gorm.DeletedAt `gorm:"index" json:"deleted_at,omitempty"` + + Account Account `gorm:"foreignKey:AccountID" json:"account,omitempty"` + User User `gorm:"foreignKey:UserID" json:"user,omitempty"` +} + +func (ContactExport) TableName() string { return "contact_exports" } diff --git a/internal/repository/contact_repo.go b/internal/repository/contact_repo.go index ba0bb57d..1b94e743 100644 --- a/internal/repository/contact_repo.go +++ b/internal/repository/contact_repo.go @@ -233,14 +233,71 @@ func (r *ContactRepo) FindActive(ctx context.Context, accountID uint, offset, li // FindAllForExport retrieves all contacts for an account (no pagination, for CSV export). // Reference: Chatwoot contacts#export func (r *ContactRepo) FindAllForExport(ctx context.Context, accountID uint) ([]model.Contact, error) { + return r.FindForExport(ctx, accountID, ContactFilterParams{}) +} + +// FindForExport retrieves all contacts matching export filters without pagination. +// Reference: Chatwoot Account::ContactsExportJob#contacts. +func (r *ContactRepo) FindForExport(ctx context.Context, accountID uint, params ContactFilterParams) ([]model.Contact, error) { var contacts []model.Contact - err := r.db.WithContext(ctx). - Where("account_id = ?", accountID). - Order("id ASC"). - Find(&contacts).Error + q := r.exportQuery(ctx, accountID, params) + err := q.Distinct("contacts.*").Order("contacts.id ASC").Find(&contacts).Error return contacts, err } +func (r *ContactRepo) exportQuery(ctx context.Context, accountID uint, params ContactFilterParams) *gorm.DB { + q := r.db.WithContext(ctx).Model(&model.Contact{}).Where("contacts.account_id = ?", accountID) + if params.ContactType != "" { + q = q.Where("contacts.contact_type = ?", params.ContactType) + } + if params.ContactSource != "" { + q = q.Where("contacts.source_id = ?", params.ContactSource) + } + if params.InboxID != nil { + q = q.Joins("JOIN contact_inboxes ON contact_inboxes.contact_id = contacts.id AND contact_inboxes.inbox_id = ?", *params.InboxID) + } + if params.Labels != "" { + q = applyContactLabelFilter(q, accountID, strings.Split(params.Labels, ",")) + } + if params.Status == "active" { + q = q.Where("contacts.last_activity_at IS NOT NULL") + } else if params.Status == "inactive" { + q = q.Where("contacts.last_activity_at IS NULL") + } + if params.UpdatedWithin != nil { + threshold := time.Now().Add(-time.Duration(*params.UpdatedWithin) * time.Second) + q = q.Where("contacts.updated_at >= ?", threshold) + } + return q +} + +// ContactLabelsByContactIDs returns approved contact label names grouped by contact ID. +func (r *ContactRepo) ContactLabelsByContactIDs(ctx context.Context, accountID uint, contactIDs []uint) (map[uint][]string, error) { + result := make(map[uint][]string, len(contactIDs)) + if len(contactIDs) == 0 { + return result, nil + } + type row struct { + ContactID uint + Name string + } + var rows []row + err := r.db.WithContext(ctx). + Table("contact_labels"). + Select("contact_labels.contact_id, tags.name"). + Joins("JOIN tags ON tags.id = contact_labels.tag_id AND tags.account_id = contact_labels.account_id"). + Where("contact_labels.account_id = ? AND contact_labels.contact_id IN ?", accountID, contactIDs). + Order("tags.name ASC"). + Scan(&rows).Error + if err != nil { + return nil, err + } + for _, r := range rows { + result[r.ContactID] = append(result[r.ContactID], r.Name) + } + return result, nil +} + // UpdateCustomAttributes updates only the custom_attributes JSON field of a contact. // Reference: Chatwoot contacts#update_custom_attributes func (r *ContactRepo) UpdateCustomAttributes(ctx context.Context, id uint, customAttrs datatypes.JSON) error { diff --git a/internal/repository/testdb_helper.go b/internal/repository/testdb_helper.go index 5f008834..9f9cd91b 100644 --- a/internal/repository/testdb_helper.go +++ b/internal/repository/testdb_helper.go @@ -152,6 +152,7 @@ func defaultTestModels() []interface{} { &model.Tag{}, &model.ConversationLabel{}, &model.ContactLabel{}, + &model.ContactExport{}, &model.DataImport{}, &model.PlatformApp{}, &model.Permissible{}, diff --git a/internal/router/router.go b/internal/router/router.go index 01edddad..53da2e4a 100644 --- a/internal/router/router.go +++ b/internal/router/router.go @@ -981,6 +981,7 @@ func registerV1Routes(g *gin.RouterGroup, h *Handlers) { contacts.GET("/active", h.Contact.Active) contacts.GET("/export", h.Contact.Export) contacts.POST("/export", h.Contact.ExportRequest) + contacts.GET("/export/:export_id/download", h.Contact.DownloadExport) contacts.POST("/import", h.Contact.Import) contacts.GET("/:contact_id", h.Contact.Get) contacts.PUT("/:contact_id", h.Contact.Update) diff --git a/internal/service/contact_service.go b/internal/service/contact_service.go index 18b49c9a..03ce787c 100644 --- a/internal/service/contact_service.go +++ b/internal/service/contact_service.go @@ -1,6 +1,7 @@ package service import ( + "bytes" "context" "encoding/csv" "encoding/json" @@ -9,6 +10,7 @@ import ( "io" "strconv" "strings" + "time" "gorm.io/datatypes" "gorm.io/gorm" @@ -305,41 +307,11 @@ func (s *ContactService) ListActive(ctx context.Context, accountID uint, offset, // GET /api/v1/accounts/:id/contacts/export // Reference: Chatwoot contacts#export func (s *ContactService) ExportCSV(ctx context.Context, accountID uint, w io.Writer) error { - contacts, err := s.repo.FindAllForExport(ctx, accountID) + csvData, _, err := s.GenerateContactExportCSV(ctx, accountID, ContactExportRequest{}) if err != nil { - return fmt.Errorf("failed to fetch contacts for export: %w", err) - } - - csvWriter := csv.NewWriter(w) - // Write header - header := []string{"id", "name", "email", "phone_number", "identifier", "country_code", "location", "contact_type", "blocked", "created_at"} - if err := csvWriter.Write(header); err != nil { - return fmt.Errorf("failed to write CSV header: %w", err) - } - - // Write rows - for _, c := range contacts { - row := []string{ - strconv.FormatUint(uint64(c.ID), 10), - c.Name, - c.Email, - c.PhoneNumber, - c.Identifier, - c.CountryCode, - c.Location, - c.ContactType, - strconv.FormatBool(c.Blocked), - c.CreatedAt.Format("2006-01-02T15:04:05Z07:00"), - } - if err := csvWriter.Write(row); err != nil { - return fmt.Errorf("failed to write CSV row: %w", err) - } - } - - csvWriter.Flush() - if err := csvWriter.Error(); err != nil { - return fmt.Errorf("CSV flush error: %w", err) + return err } + _, err = w.Write(csvData) return nil } @@ -350,6 +322,308 @@ type ImportCSVResult struct { Failed int `json:"failed"` } +type ContactExportRequest struct { + ColumnNames []string `json:"column_names"` + Payload []ContactExportFilterCondition `json:"payload"` + Label string `json:"label"` +} + +type ContactExportFilterCondition struct { + AttributeKey string `json:"attribute_key"` + FilterType string `json:"filter_type"` + Operator string `json:"operator"` + Values []any `json:"values"` +} + +func (s *ContactService) ExportContacts(ctx context.Context, accountID, userID uint, req ContactExportRequest) (*model.ContactExport, error) { + if !s.Ready() { + return nil, errors.New("contact service not ready") + } + + var account model.Account + if err := s.repo.DB().WithContext(ctx).First(&account, accountID).Error; err != nil { + return nil, err + } + + var userIDPtr *uint + if userID != 0 { + userIDPtr = &userID + } + columnsJSON, _ := json.Marshal(req.ColumnNames) + filterJSON, _ := json.Marshal(map[string]any{"payload": req.Payload, "label": req.Label}) + export := &model.ContactExport{ + AccountID: accountID, + UserID: userIDPtr, + Status: string(model.DataImportStatusPending), + FileName: contactExportFilename(account), + ContentType: "text/csv", + ColumnNames: columnsJSON, + FilterParams: filterJSON, + } + if err := s.repo.DB().WithContext(ctx).Create(export).Error; err != nil { + return nil, err + } + if err := s.repo.DB().WithContext(ctx).Model(export).Update("status", string(model.DataImportStatusProcessing)).Error; err != nil { + return nil, err + } + + csvData, rowCount, err := s.GenerateContactExportCSV(ctx, accountID, req) + if err != nil { + s.repo.DB().WithContext(ctx).Model(export).Updates(map[string]any{ + "status": string(model.DataImportStatusFailed), + "error": err.Error(), + }) + return export, err + } + + completedAt := time.Now() + export.FileURL = fmt.Sprintf("/api/v1/accounts/%d/contacts/export/%d/download", accountID, export.ID) + updates := map[string]any{ + "status": string(model.DataImportStatusCompleted), + "csv_data": csvData, + "row_count": rowCount, + "file_url": export.FileURL, + "completed_at": completedAt, + } + if err := s.repo.DB().WithContext(ctx).Model(export).Updates(updates).Error; err != nil { + return nil, err + } + if err := s.repo.DB().WithContext(ctx).First(export, export.ID).Error; err != nil { + return nil, err + } + if err := s.createContactExportNotification(ctx, export); err != nil { + applogger.L().Warnf("contact export notification failed: %v", err) + } + return export, nil +} + +func (s *ContactService) GenerateContactExportCSV(ctx context.Context, accountID uint, req ContactExportRequest) ([]byte, int, error) { + params := contactExportFilterParams(req) + contacts, err := s.repo.FindForExport(ctx, accountID, params) + if err != nil { + return nil, 0, fmt.Errorf("failed to fetch contacts for export: %w", err) + } + + headers := validContactExportHeaders(req.ColumnNames) + labelsByContactID := map[uint][]string{} + if containsString(headers, "labels") { + ids := make([]uint, 0, len(contacts)) + for _, contact := range contacts { + ids = append(ids, contact.ID) + } + labelsByContactID, err = s.repo.ContactLabelsByContactIDs(ctx, accountID, ids) + if err != nil { + return nil, 0, fmt.Errorf("failed to fetch contact labels for export: %w", err) + } + } + + var body bytes.Buffer + body.Write([]byte{0xEF, 0xBB, 0xBF}) + csvWriter := csv.NewWriter(&body) + if err := csvWriter.Write(headers); err != nil { + return nil, 0, fmt.Errorf("failed to write CSV header: %w", err) + } + for _, contact := range contacts { + row := make([]string, 0, len(headers)) + for _, header := range headers { + row = append(row, contactExportValue(contact, header, labelsByContactID[contact.ID])) + } + if err := csvWriter.Write(row); err != nil { + return nil, 0, fmt.Errorf("failed to write CSV row: %w", err) + } + } + csvWriter.Flush() + if err := csvWriter.Error(); err != nil { + return nil, 0, fmt.Errorf("CSV flush error: %w", err) + } + return body.Bytes(), len(contacts), nil +} + +func (s *ContactService) GetContactExport(ctx context.Context, accountID, exportID uint) (*model.ContactExport, error) { + var export model.ContactExport + err := s.repo.DB().WithContext(ctx).Where("account_id = ? AND id = ?", accountID, exportID).First(&export).Error + if err != nil { + return nil, err + } + return &export, nil +} + +func contactExportFilename(account model.Account) string { + name := strings.TrimSpace(account.Name) + if name == "" { + name = "account" + } + name = strings.NewReplacer("/", "_", "\\", "_", " ", "_").Replace(name) + return fmt.Sprintf("%s_%d_contacts.csv", name, account.ID) +} + +func contactExportFilterParams(req ContactExportRequest) repository.ContactFilterParams { + params := repository.ContactFilterParams{} + if strings.TrimSpace(req.Label) != "" { + params.Labels = strings.TrimSpace(req.Label) + } + for _, condition := range req.Payload { + if len(condition.Values) == 0 { + continue + } + value := strings.TrimSpace(contactExportFilterValue(condition.Values[0])) + switch strings.TrimSpace(condition.AttributeKey) { + case "contact_type": + params.ContactType = value + case "source_id", "contact_source": + params.ContactSource = value + case "status": + params.Status = value + case "labels", "label_list": + params.Labels = strings.Join(contactExportFilterValues(condition.Values), ",") + case "inbox_id": + if n, err := strconv.ParseUint(value, 10, 32); err == nil && n != 0 { + inboxID := uint(n) + params.InboxID = &inboxID + } + case "updated_within": + if n, err := strconv.Atoi(value); err == nil && n > 0 { + params.UpdatedWithin = &n + } + } + } + return params +} + +func contactExportFilterValue(value any) string { + switch v := value.(type) { + case string: + return v + case float64: + return strconv.FormatFloat(v, 'f', -1, 64) + case int: + return strconv.Itoa(v) + case uint: + return strconv.FormatUint(uint64(v), 10) + case json.Number: + return v.String() + default: + return fmt.Sprint(v) + } +} + +func contactExportFilterValues(values []any) []string { + result := make([]string, 0, len(values)) + for _, value := range values { + text := strings.TrimSpace(contactExportFilterValue(value)) + if text != "" { + result = append(result, text) + } + } + return result +} + +func validContactExportHeaders(columnNames []string) []string { + requested := columnNames + if len(requested) == 0 { + requested = []string{"id", "name", "email", "phone_number", "labels"} + } + allowed := map[string]struct{}{ + "id": {}, "name": {}, "middle_name": {}, "last_name": {}, "email": {}, "phone_number": {}, + "identifier": {}, "country_code": {}, "location": {}, "contact_type": {}, "blocked": {}, + "source_id": {}, "company_id": {}, "last_activity_at": {}, "created_at": {}, "updated_at": {}, "labels": {}, + } + seen := map[string]struct{}{} + headers := make([]string, 0, len(requested)) + for _, header := range requested { + header = strings.TrimSpace(header) + if header == "" { + continue + } + if _, ok := allowed[header]; !ok { + continue + } + if _, ok := seen[header]; ok { + continue + } + seen[header] = struct{}{} + headers = append(headers, header) + } + return headers +} + +func contactExportValue(contact model.Contact, header string, labels []string) string { + switch header { + case "id": + return strconv.FormatUint(uint64(contact.ID), 10) + case "name": + return contact.Name + case "middle_name": + return contact.MiddleName + case "last_name": + return contact.LastName + case "email": + return contact.Email + case "phone_number": + return contact.PhoneNumber + case "identifier": + return contact.Identifier + case "country_code": + return contact.CountryCode + case "location": + return contact.Location + case "contact_type": + return contact.ContactType + case "blocked": + return strconv.FormatBool(contact.Blocked) + case "source_id": + return contact.SourceID + case "company_id": + if contact.CompanyID == nil { + return "" + } + return strconv.FormatUint(uint64(*contact.CompanyID), 10) + case "last_activity_at": + if contact.LastActivityAt == nil { + return "" + } + return strconv.FormatInt(*contact.LastActivityAt, 10) + case "created_at": + return contact.CreatedAt.Format(time.RFC3339) + case "updated_at": + return contact.UpdatedAt.Format(time.RFC3339) + case "labels": + return strings.Join(labels, ",") + default: + return "" + } +} + +func containsString(values []string, needle string) bool { + for _, value := range values { + if value == needle { + return true + } + } + return false +} + +func (s *ContactService) createContactExportNotification(ctx context.Context, export *model.ContactExport) error { + if export.UserID == nil || *export.UserID == 0 { + return nil + } + attrs, _ := json.Marshal(map[string]any{ + "file_url": export.FileURL, + "file_name": export.FileName, + "row_count": export.RowCount, + }) + notification := &model.Notification{ + AccountID: &export.AccountID, + UserID: *export.UserID, + NotificationType: "contacts_export_complete", + PrimaryActorType: "ContactExport", + PrimaryActorID: export.ID, + EmailEnabled: true, + AdditionalAttributes: attrs, + } + return s.repo.DB().WithContext(ctx).Create(notification).Error +} + func (s *ContactService) ImportContacts(ctx context.Context, accountID, userID uint, r io.Reader) (*model.DataImport, error) { if !s.Ready() { return nil, errors.New("contact service not ready") diff --git a/internal/service/contact_service_g3_test.go b/internal/service/contact_service_g3_test.go index 4f5a54f7..42d888a7 100644 --- a/internal/service/contact_service_g3_test.go +++ b/internal/service/contact_service_g3_test.go @@ -126,22 +126,17 @@ func TestContactService_ExportCSV_WritesCorrectHeadersAndRows(t *testing.T) { err := svc.ExportCSV(context.Background(), account.ID, &buf) require.NoError(t, err) - csvOutput := buf.String() + csvOutput := strings.TrimPrefix(buf.String(), "\ufeff") lines := strings.Split(csvOutput, "\n") // Verify header - assert.Equal(t, "id,name,email,phone_number,identifier,country_code,location,contact_type,blocked,created_at", lines[0]) + assert.Equal(t, "id,name,email,phone_number,labels", lines[0]) // Verify data row contains the contact info dataLine := lines[1] assert.Contains(t, dataLine, "Export Contact") assert.Contains(t, dataLine, "export@test.com") assert.Contains(t, dataLine, "+1234567890") - assert.Contains(t, dataLine, "id123") - assert.Contains(t, dataLine, "US") - assert.Contains(t, dataLine, "New York") - assert.Contains(t, dataLine, "lead") - assert.Contains(t, dataLine, "false") } func TestContactService_ExportCSV_EmptyAccount(t *testing.T) { @@ -152,11 +147,11 @@ func TestContactService_ExportCSV_EmptyAccount(t *testing.T) { err := svc.ExportCSV(context.Background(), account.ID, &buf) require.NoError(t, err) - csvOutput := buf.String() + csvOutput := strings.TrimPrefix(buf.String(), "\ufeff") lines := strings.Split(csvOutput, "\n") // Should have header only (plus trailing empty line from csv writer) - assert.Equal(t, "id,name,email,phone_number,identifier,country_code,location,contact_type,blocked,created_at", lines[0]) + assert.Equal(t, "id,name,email,phone_number,labels", lines[0]) // No data rows assert.Empty(t, strings.TrimSpace(lines[1])) } @@ -185,6 +180,57 @@ func TestContactService_ExportCSV_MultipleContacts(t *testing.T) { assert.Len(t, lines, 4) } +func TestContactService_ExportContacts_PersistsArtifactAndNotification(t *testing.T) { + db, _, svc := setupContactService(t) + account := createTestAccount(t, db) + user := createTestUser(t, db, account.ID) + contact := &model.Contact{AccountID: account.ID, Name: "Export Alice", Email: "alice@example.com", PhoneNumber: "+111"} + require.NoError(t, db.Create(contact).Error) + tag := &model.Tag{AccountID: account.ID, Name: "vip"} + require.NoError(t, db.Create(tag).Error) + require.NoError(t, db.Create(&model.ContactLabel{AccountID: account.ID, ContactID: contact.ID, TagID: tag.ID}).Error) + + export, err := svc.ExportContacts(context.Background(), account.ID, user.ID, ContactExportRequest{}) + require.NoError(t, err) + assert.Equal(t, string(model.DataImportStatusCompleted), export.Status) + assert.Equal(t, 1, export.RowCount) + assert.Contains(t, export.FileName, "contacts.csv") + assert.Contains(t, export.FileURL, "/contacts/export/") + assert.Contains(t, string(export.CSVData), "id,name,email,phone_number,labels") + assert.Contains(t, string(export.CSVData), "alice@example.com") + assert.Contains(t, string(export.CSVData), "vip") + + var notification model.Notification + require.NoError(t, db.Where("user_id = ? AND notification_type = ?", user.ID, "contacts_export_complete").First(¬ification).Error) + assert.Equal(t, "ContactExport", notification.PrimaryActorType) + assert.Equal(t, export.ID, notification.PrimaryActorID) + assert.True(t, notification.EmailEnabled) +} + +func TestContactService_ExportContacts_FiltersByLabelAndColumns(t *testing.T) { + db, _, svc := setupContactService(t) + account := createTestAccount(t, db) + keep := &model.Contact{AccountID: account.ID, Name: "Keep", Email: "keep@example.com"} + drop := &model.Contact{AccountID: account.ID, Name: "Drop", Email: "drop@example.com"} + require.NoError(t, db.Create(keep).Error) + require.NoError(t, db.Create(drop).Error) + tag := &model.Tag{AccountID: account.ID, Name: "vip"} + require.NoError(t, db.Create(tag).Error) + require.NoError(t, db.Create(&model.ContactLabel{AccountID: account.ID, ContactID: keep.ID, TagID: tag.ID}).Error) + + csvData, rowCount, err := svc.GenerateContactExportCSV(context.Background(), account.ID, ContactExportRequest{ + ColumnNames: []string{"email", "labels", "bogus", "email"}, + Label: "vip", + }) + require.NoError(t, err) + assert.Equal(t, 1, rowCount) + output := strings.TrimPrefix(string(csvData), "\ufeff") + assert.Contains(t, output, "email,labels") + assert.Contains(t, output, "keep@example.com,vip") + assert.NotContains(t, output, "drop@example.com") + assert.NotContains(t, output, "bogus") +} + // ========== ImportCSV ========== func TestContactService_ImportCSV_ImportsValidRows(t *testing.T) { diff --git a/internal/service/service_test_helper.go b/internal/service/service_test_helper.go index 8452ac81..0d52314c 100644 --- a/internal/service/service_test_helper.go +++ b/internal/service/service_test_helper.go @@ -68,7 +68,9 @@ func setupServiceTestDB(t *testing.T) *gorm.DB { &model.CompanyNote{}, &model.Tag{}, &model.ContactLabel{}, + &model.ContactExport{}, &model.DataImport{}, + &model.Notification{}, ); err != nil { t.Fatalf("failed to auto-migrate models: %v", err) } @@ -273,7 +275,9 @@ func setupContactServiceTestDB(t *testing.T) *gorm.DB { &model.ContactInbox{}, &model.Tag{}, &model.ContactLabel{}, + &model.ContactExport{}, &model.DataImport{}, + &model.Notification{}, ); err != nil { t.Fatalf("failed to auto-migrate contact models: %v", err) } diff --git a/migrations/000021_add_contact_exports.down.sql b/migrations/000021_add_contact_exports.down.sql new file mode 100644 index 00000000..d101f806 --- /dev/null +++ b/migrations/000021_add_contact_exports.down.sql @@ -0,0 +1 @@ +DROP TABLE IF EXISTS contact_exports; diff --git a/migrations/000021_add_contact_exports.up.sql b/migrations/000021_add_contact_exports.up.sql new file mode 100644 index 00000000..50133777 --- /dev/null +++ b/migrations/000021_add_contact_exports.up.sql @@ -0,0 +1,22 @@ +CREATE TABLE IF NOT EXISTS contact_exports ( + id BIGSERIAL PRIMARY KEY, + account_id BIGINT NOT NULL, + user_id BIGINT, + status VARCHAR(50) NOT NULL DEFAULT 'pending', + file_name VARCHAR(255) NOT NULL, + content_type VARCHAR(100) NOT NULL DEFAULT 'text/csv', + file_url VARCHAR(512), + csv_data BYTEA, + row_count INTEGER DEFAULT 0, + column_names JSONB, + filter_params JSONB, + error TEXT, + completed_at TIMESTAMP WITH TIME ZONE, + created_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(), + updated_at TIMESTAMP WITH TIME ZONE DEFAULT NOW(), + deleted_at TIMESTAMP WITH TIME ZONE +); + +CREATE INDEX IF NOT EXISTS idx_contact_exports_account_id ON contact_exports(account_id); +CREATE INDEX IF NOT EXISTS idx_contact_exports_user_id ON contact_exports(user_id); +CREATE INDEX IF NOT EXISTS idx_contact_exports_deleted_at ON contact_exports(deleted_at);