feat(crm): persist contact export artifacts
This commit is contained in:
@@ -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
|
||||
|
||||
+7
-6
@@ -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{},
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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"}})
|
||||
|
||||
|
||||
@@ -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" }
|
||||
@@ -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 {
|
||||
|
||||
@@ -152,6 +152,7 @@ func defaultTestModels() []interface{} {
|
||||
&model.Tag{},
|
||||
&model.ConversationLabel{},
|
||||
&model.ContactLabel{},
|
||||
&model.ContactExport{},
|
||||
&model.DataImport{},
|
||||
&model.PlatformApp{},
|
||||
&model.Permissible{},
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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")
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
|
||||
@@ -0,0 +1 @@
|
||||
DROP TABLE IF EXISTS contact_exports;
|
||||
@@ -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);
|
||||
Reference in New Issue
Block a user