feat(crm): queue contact imports

This commit is contained in:
2026-06-07 02:17:50 +08:00
parent b16f922daa
commit 3a71001cb9
6 changed files with 186 additions and 20 deletions
+1 -1
View File
@@ -932,7 +932,7 @@ func (h *ContactHandler) Import(c *gin.Context) {
file, _, fileErr = c.Request.FormFile("file")
}
if fileErr != nil {
c.JSON(http.StatusUnprocessableEntity, gin.H{"error": "failed to import contacts"})
c.JSON(http.StatusUnprocessableEntity, gin.H{"error": "File is blank"})
return
}
defer file.Close()
@@ -93,7 +93,7 @@ func TestContactImportMissingFile(t *testing.T) {
assert.Equal(t, http.StatusUnprocessableEntity, w.Code)
var resp map[string]interface{}
json.Unmarshal(w.Body.Bytes(), &resp)
assert.Contains(t, resp, "error")
assert.Equal(t, "File is blank", resp["error"])
}
func TestContactableInboxesBadAccountID(t *testing.T) {
+47
View File
@@ -0,0 +1,47 @@
package service
import (
"context"
"encoding/json"
"fmt"
"sync"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/worker"
)
const TaskTypeContactImport = "contact:import"
type contactImportConfig struct {
CSVBase64 string `json:"csv_base64"`
}
type contactImportJob struct {
ImportID uint `json:"import_id"`
}
var contactImportRegistrations sync.Map
// RegisterContactImportJobs wires Chatwoot DataImportJob contacts import into
// the durable worker. Handlers are registered once per WorkerPool instance.
func RegisterContactImportJobs(wp *worker.WorkerPool, svc *ContactService) {
if wp == nil || svc == nil {
return
}
if _, loaded := contactImportRegistrations.LoadOrStore(wp, struct{}{}); loaded {
return
}
wp.Register(TaskTypeContactImport, svc.performContactImportJob)
}
func (s *ContactService) performContactImportJob(ctx context.Context, job *model.BackgroundJob) error {
var payload contactImportJob
if err := json.Unmarshal(job.Payload, &payload); err != nil {
return fmt.Errorf("unmarshal contact import job: %w", err)
}
if payload.ImportID == 0 {
return fmt.Errorf("invalid contact import job payload: %#v", payload)
}
_, err := s.performContactImport(ctx, payload.ImportID)
return err
}
+48 -8
View File
@@ -3,6 +3,7 @@ package service
import (
"bytes"
"context"
"encoding/base64"
"encoding/csv"
"encoding/json"
"errors"
@@ -56,6 +57,7 @@ func (s *ContactService) SetContactExportMailer(mailer ContactExportMailer) {
func (s *ContactService) SetWorkerPool(wp *worker.WorkerPool) {
s.worker = wp
RegisterContactExportJobs(wp, s)
RegisterContactImportJobs(wp, s)
}
func (s *ContactService) indexContact(ctx context.Context, contact *model.Contact) {
@@ -995,25 +997,63 @@ func (s *ContactService) ImportContacts(ctx context.Context, accountID, userID u
if !s.Ready() {
return nil, errors.New("contact service not ready")
}
csvData, err := io.ReadAll(r)
if err != nil {
return nil, fmt.Errorf("read import file: %w", err)
}
var userIDPtr *uint
if userID != 0 {
userIDPtr = &userID
}
dataImport := &model.DataImport{AccountID: accountID, UserID: userIDPtr, DataType: "contacts", Status: string(model.DataImportStatusPending)}
config, _ := json.Marshal(contactImportConfig{CSVBase64: base64.StdEncoding.EncodeToString(csvData)})
dataImport := &model.DataImport{AccountID: accountID, UserID: userIDPtr, DataType: "contacts", Status: string(model.DataImportStatusPending), ImportConfig: config}
if err := s.repo.DB().WithContext(ctx).Create(dataImport).Error; err != nil {
return nil, err
}
if err := s.repo.DB().WithContext(ctx).Model(dataImport).Update("status", string(model.DataImportStatusProcessing)).Error; err != nil {
if s.worker != nil {
_, err := s.worker.Enqueue(ctx, TaskTypeContactImport, contactImportJob{ImportID: dataImport.ID}, worker.WithQueue("low"), worker.WithMaxAttempts(3), worker.WithIdempotencyKey(fmt.Sprintf("contact-import:%d", dataImport.ID)))
if err != nil {
s.repo.DB().WithContext(ctx).Model(dataImport).Updates(map[string]any{
"status": string(model.DataImportStatusFailed),
"processing_errors": err.Error(),
})
return dataImport, err
}
return dataImport, nil
}
return s.performContactImport(ctx, dataImport.ID)
}
func (s *ContactService) performContactImport(ctx context.Context, importID uint) (*model.DataImport, error) {
if !s.Ready() {
return nil, errors.New("contact service not ready")
}
var dataImport model.DataImport
if err := s.repo.DB().WithContext(ctx).First(&dataImport, importID).Error; err != nil {
return nil, err
}
if dataImport.Status == string(model.DataImportStatusCompleted) {
return &dataImport, nil
}
var config contactImportConfig
if err := json.Unmarshal(dataImport.ImportConfig, &config); err != nil {
return nil, fmt.Errorf("unmarshal import config: %w", err)
}
csvData, err := base64.StdEncoding.DecodeString(config.CSVBase64)
if err != nil {
return nil, fmt.Errorf("decode import csv: %w", err)
}
if err := s.repo.DB().WithContext(ctx).Model(&dataImport).Update("status", string(model.DataImportStatusProcessing)).Error; err != nil {
return nil, err
}
result, err := s.ImportCSV(ctx, accountID, r)
result, err := s.ImportCSV(ctx, dataImport.AccountID, bytes.NewReader(csvData))
if err != nil {
s.repo.DB().WithContext(ctx).Model(dataImport).Updates(map[string]any{
s.repo.DB().WithContext(ctx).Model(&dataImport).Updates(map[string]any{
"status": string(model.DataImportStatusFailed),
"processing_errors": err.Error(),
})
return dataImport, err
return &dataImport, err
}
updates := map[string]any{
"status": string(model.DataImportStatusCompleted),
@@ -1021,13 +1061,13 @@ func (s *ContactService) ImportContacts(ctx context.Context, accountID, userID u
"failed_records": result.Failed,
"total_records": result.Imported + result.Skipped + result.Failed,
}
if err := s.repo.DB().WithContext(ctx).Model(dataImport).Updates(updates).Error; err != nil {
if err := s.repo.DB().WithContext(ctx).Model(&dataImport).Updates(updates).Error; err != nil {
return nil, err
}
if err := s.repo.DB().WithContext(ctx).First(dataImport, dataImport.ID).Error; err != nil {
if err := s.repo.DB().WithContext(ctx).First(&dataImport, dataImport.ID).Error; err != nil {
return nil, err
}
return dataImport, nil
return &dataImport, nil
}
// ImportCSV reads contacts from a CSV reader and creates them.
@@ -336,6 +336,62 @@ func TestContactService_ContactExportJobRetriesMissingExport(t *testing.T) {
assert.NotEmpty(t, job.LastError)
}
func TestContactService_ImportContacts_QueuesDurableDataImport(t *testing.T) {
db, _, svc := setupContactService(t)
account := createTestAccount(t, db)
userID := uint(42)
require.NoError(t, db.Create(&model.Tag{AccountID: account.ID, Name: "vip"}).Error)
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return time.Date(2026, 6, 7, 9, 0, 0, 0, time.UTC) }))
svc.SetWorkerPool(wp)
dataImport, err := svc.ImportContacts(context.Background(), account.ID, userID, strings.NewReader("name,email,labels\nQueued,queued@test.com,vip\n"))
require.NoError(t, err)
assert.Equal(t, "contacts", dataImport.DataType)
assert.Equal(t, string(model.DataImportStatusPending), dataImport.Status)
assert.NotEmpty(t, dataImport.ImportConfig)
require.NotNil(t, dataImport.UserID)
assert.Equal(t, userID, *dataImport.UserID)
var contactCount int64
require.NoError(t, db.Model(&model.Contact{}).Where("account_id = ? AND email = ?", account.ID, "queued@test.com").Count(&contactCount).Error)
assert.Equal(t, int64(0), contactCount)
var jobCount int64
require.NoError(t, db.Model(&model.BackgroundJob{}).Where("job_type = ? AND queue = ? AND status = ?", TaskTypeContactImport, "low", model.BackgroundJobStatusQueued).Count(&jobCount).Error)
assert.Equal(t, int64(1), jobCount)
processed, err := wp.ProcessOne(context.Background())
require.NoError(t, err)
assert.True(t, processed)
var completed model.DataImport
require.NoError(t, db.First(&completed, dataImport.ID).Error)
assert.Equal(t, string(model.DataImportStatusCompleted), completed.Status)
assert.Equal(t, 1, completed.TotalRecords)
assert.Equal(t, 1, completed.ProcessedRecords)
require.NoError(t, db.Model(&model.Contact{}).Where("account_id = ? AND email = ?", account.ID, "queued@test.com").Count(&contactCount).Error)
assert.Equal(t, int64(1), contactCount)
}
func TestContactService_ContactImportJobRetriesMissingImport(t *testing.T) {
db, _, svc := setupContactService(t)
wp := worker.NewWorkerPoolWithOptions(db, worker.WithNow(func() time.Time { return time.Date(2026, 6, 7, 9, 15, 0, 0, time.UTC) }), worker.WithBackoff(func(attempt int) time.Duration { return time.Minute }))
svc.SetWorkerPool(wp)
_, err := wp.Enqueue(context.Background(), TaskTypeContactImport, contactImportJob{ImportID: 9999}, worker.WithQueue("low"), worker.WithMaxAttempts(3))
require.NoError(t, err)
processed, err := wp.ProcessOne(context.Background())
if err == nil || !processed {
t.Fatalf("expected missing contact import to retry, processed=%v err=%v", processed, err)
}
var job model.BackgroundJob
require.NoError(t, db.Where("job_type = ?", TaskTypeContactImport).First(&job).Error)
assert.Equal(t, model.BackgroundJobStatusRetrying, job.Status)
assert.NotEmpty(t, job.LastError)
}
func TestContactService_ExportContacts_FiltersByLabelAndColumns(t *testing.T) {
db, _, svc := setupContactService(t)
account := createTestAccount(t, db)