Files
gochat/internal/service/contact_service.go
T

609 lines
20 KiB
Go

package service
import (
"context"
"encoding/csv"
"encoding/json"
"errors"
"fmt"
"io"
"strconv"
"strings"
"gorm.io/datatypes"
"gorm.io/gorm"
"github.com/gochat/gochat/internal/model"
"github.com/gochat/gochat/internal/repository"
"github.com/gochat/gochat/internal/search"
applogger "github.com/gochat/gochat/pkg/logger"
pkgvalidator "github.com/gochat/gochat/pkg/validator"
)
// ContactService implements business logic for Contact operations.
// Reference: Chatwoot app/controllers/api/v1/contacts_controller.rb
type ContactService struct {
repo *repository.ContactRepo
contactInboxSvc *ContactInboxService
noteRepo *repository.NoteRepo
searchIndexer SearchIndexer
}
// NewContactService creates a new Contact service.
func NewContactService(repo *repository.ContactRepo, contactInboxSvc *ContactInboxService, noteRepo *repository.NoteRepo) *ContactService {
return &ContactService{repo: repo, contactInboxSvc: contactInboxSvc, noteRepo: noteRepo}
}
func (s *ContactService) SetSearchIndexer(indexer SearchIndexer) {
s.searchIndexer = indexer
}
func (s *ContactService) indexContact(ctx context.Context, contact *model.Contact) {
if s.searchIndexer != nil {
logSearchIndexError("contact", contact.ID, s.searchIndexer.IndexContact(ctx, contact))
}
}
func (s *ContactService) deleteContactIndex(ctx context.Context, accountID uint, id uint) {
if s.searchIndexer != nil {
logSearchIndexError("contact", id, s.searchIndexer.DeleteContact(ctx, accountID, id))
}
}
// Ready reports whether the service has the dependencies required for DB-backed operations.
func (s *ContactService) Ready() bool {
return s != nil && s.repo != nil
}
func (s *ContactService) DB() *gorm.DB {
if s == nil || s.repo == nil {
return nil
}
return s.repo.DB()
}
// ListByAccount retrieves all contacts for an account with optional sort.
func (s *ContactService) ListByAccount(ctx context.Context, accountID uint, offset, limit int, sort string, labels ...[]string) ([]model.Contact, int64, error) {
return s.repo.FindByAccount(ctx, accountID, offset, limit, sort, labels...)
}
// Search searches contacts by name, email, phone, or identifier with optional sort.
func (s *ContactService) Search(ctx context.Context, accountID uint, query string, offset, limit int, sort string, searchMode search.SearchMode, labels ...[]string) ([]model.Contact, int64, error) {
if query == "" {
return s.repo.FindByAccount(ctx, accountID, offset, limit, sort, labels...)
}
return s.repo.Search(ctx, accountID, query, offset, limit, sort, searchMode, labels...)
}
// GetByID retrieves a single contact.
func (s *ContactService) GetByID(ctx context.Context, id uint) (*model.Contact, error) {
return s.repo.FindByID(ctx, id)
}
// GetByAccountAndID retrieves a contact scoped to an account.
func (s *ContactService) GetByAccountAndID(ctx context.Context, accountID, id uint) (*model.Contact, error) {
return s.repo.FindByAccountAndID(ctx, accountID, id)
}
// ListContactInboxes retrieves all contact_inboxes for a contact.
func (s *ContactService) ListContactInboxes(ctx context.Context, contactID uint) ([]model.ContactInbox, error) {
return s.contactInboxSvc.ListByContact(ctx, contactID)
}
// CreateContactRequest is the DTO for creating a contact.
// Reference: Chatwoot app/controllers/api/v1/contacts_controller.rb#create
// When inbox_id is provided, a ContactInbox record is auto-created (Chatwoot pattern).
type CreateContactRequest struct {
Name string `json:"name" validate:"required,min=1"`
Email string `json:"email,omitempty" validate:"omitempty,email"`
Phone string `json:"phone,omitempty"`
Identifier string `json:"identifier,omitempty"`
AvatarURL string `json:"avatar_url,omitempty"`
InboxID *uint `json:"inbox_id,omitempty"`
SourceID string `json:"source_id,omitempty"`
AdditionalAttributes *model.JSONMap `json:"additional_attributes,omitempty"`
CustomAttributes *model.JSONMap `json:"custom_attributes,omitempty"`
ContactType string `json:"contact_type,omitempty"`
MiddleName string `json:"middle_name,omitempty"`
LastName string `json:"last_name,omitempty"`
CountryCode string `json:"country_code,omitempty"`
Location string `json:"location,omitempty"`
CompanyID *uint `json:"company_id,omitempty"`
}
// Create creates a new contact and optionally auto-creates a ContactInbox when inbox_id is provided.
// Reference: Chatwoot contacts_controller#create — auto-creates ContactInbox for channel source.
func (s *ContactService) Create(ctx context.Context, accountID uint, req CreateContactRequest) (*model.Contact, error) {
if err := pkgvalidator.ValidateStruct(req); err != nil {
return nil, err
}
contact := &model.Contact{
AccountID: accountID,
Name: req.Name,
Email: req.Email,
PhoneNumber: req.Phone,
Identifier: req.Identifier,
AvatarURL: req.AvatarURL,
MiddleName: req.MiddleName,
LastName: req.LastName,
CountryCode: req.CountryCode,
Location: req.Location,
ContactType: req.ContactType,
SourceID: req.SourceID,
CompanyID: req.CompanyID,
}
if req.AdditionalAttributes != nil {
contact.AdditionalAttributes = mergeContactJSON(contact.AdditionalAttributes, model.ToDatatypesJSON(req.AdditionalAttributes))
}
if req.CustomAttributes != nil {
contact.CustomAttributes = mergeContactJSON(contact.CustomAttributes, model.ToDatatypesJSON(req.CustomAttributes))
}
if err := s.repo.Create(ctx, contact); err != nil {
applogger.L().Errorf("Failed to create contact: %v", err)
return nil, err
}
s.indexContact(ctx, contact)
// Auto-create ContactInbox when inbox_id is provided (Chatwoot pattern)
if req.InboxID != nil && *req.InboxID > 0 {
sourceID := req.SourceID
if sourceID == "" {
sourceID = contact.Email // fallback: use email as source_id
}
ciReq := CreateContactInboxRequest{
ContactID: contact.ID,
InboxID: *req.InboxID,
SourceID: sourceID,
}
if _, err := s.contactInboxSvc.Create(ctx, ciReq); err != nil {
applogger.L().Errorf("Failed to auto-create ContactInbox for contact %d: %v", contact.ID, err)
// Non-blocking: contact is created, but ContactInbox creation failed
}
}
return contact, nil
}
// UpdateContactRequest is the DTO for updating a contact.
type UpdateContactRequest struct {
Name string `json:"name,omitempty" validate:"omitempty,min=1"`
Email string `json:"email,omitempty" validate:"omitempty,email"`
Phone string `json:"phone,omitempty"`
Identifier string `json:"identifier,omitempty"`
AvatarURL string `json:"avatar_url,omitempty"`
MiddleName string `json:"middle_name,omitempty"`
LastName string `json:"last_name,omitempty"`
CountryCode string `json:"country_code,omitempty"`
Location string `json:"location,omitempty"`
ContactType string `json:"contact_type,omitempty"`
AdditionalAttributes *model.JSONMap `json:"additional_attributes,omitempty"`
CustomAttributes *model.JSONMap `json:"custom_attributes,omitempty"`
CompanyID *uint `json:"company_id,omitempty"`
}
// Update modifies an existing contact.
func (s *ContactService) Update(ctx context.Context, accountID, id uint, req UpdateContactRequest) (*model.Contact, error) {
if err := pkgvalidator.ValidateStruct(req); err != nil {
return nil, err
}
contact, err := s.repo.FindByAccountAndID(ctx, accountID, id)
if err != nil {
return nil, err
}
if req.Name != "" {
contact.Name = req.Name
}
if req.Email != "" {
contact.Email = req.Email
}
if req.Phone != "" {
contact.PhoneNumber = req.Phone
}
if req.Identifier != "" {
contact.Identifier = req.Identifier
}
if req.AvatarURL != "" {
contact.AvatarURL = req.AvatarURL
}
if req.MiddleName != "" {
contact.MiddleName = req.MiddleName
}
if req.LastName != "" {
contact.LastName = req.LastName
}
if req.CountryCode != "" {
contact.CountryCode = req.CountryCode
}
if req.Location != "" {
contact.Location = req.Location
}
if req.ContactType != "" {
contact.ContactType = req.ContactType
}
if req.AdditionalAttributes != nil {
contact.AdditionalAttributes = model.ToDatatypesJSON(req.AdditionalAttributes)
}
if req.CustomAttributes != nil {
contact.CustomAttributes = model.ToDatatypesJSON(req.CustomAttributes)
}
if req.CompanyID != nil {
contact.CompanyID = req.CompanyID
}
if err := s.repo.Update(ctx, contact); err != nil {
return nil, err
}
s.indexContact(ctx, contact)
return contact, nil
}
// Delete soft-deletes a contact.
func (s *ContactService) Delete(ctx context.Context, accountID, id uint) error {
contact, err := s.repo.FindByAccountAndID(ctx, accountID, id)
if err != nil {
return err
}
if err := s.repo.Delete(ctx, contact.ID); err != nil {
return err
}
s.deleteContactIndex(ctx, accountID, contact.ID)
return nil
}
// CreateNoteRequest is the DTO for creating a contact note.
type CreateNoteRequest struct {
Content string `json:"content" validate:"required,min=1"`
}
// ListNotes retrieves notes for a contact.
// Reference: Chatwoot app/controllers/api/v1/contacts/notes_controller.rb #index
func (s *ContactService) ListNotes(ctx context.Context, accountID, contactID uint) ([]model.Note, error) {
// Verify contact belongs to account
_, err := s.repo.FindByAccountAndID(ctx, accountID, contactID)
if err != nil {
return nil, errors.New("contact not found")
}
return s.noteRepo.FindByContact(ctx, accountID, contactID)
}
// CreateNote creates a note for a contact.
// Reference: Chatwoot app/controllers/api/v1/contacts/notes_controller.rb #create
func (s *ContactService) CreateNote(ctx context.Context, accountID, contactID, userID uint, req CreateNoteRequest) (*model.Note, error) {
if err := pkgvalidator.ValidateStruct(req); err != nil {
return nil, err
}
// Verify contact belongs to account
_, err := s.repo.FindByAccountAndID(ctx, accountID, contactID)
if err != nil {
return nil, errors.New("contact not found")
}
note := &model.Note{
Content: req.Content,
AccountID: accountID,
ContactID: contactID,
UserID: &userID,
}
if err := s.noteRepo.Create(ctx, note); err != nil {
return nil, err
}
return note, nil
}
// ListActive retrieves contacts with recent activity for an account.
// GET /api/v1/accounts/:id/contacts/active
// Reference: Chatwoot contacts#active
func (s *ContactService) ListActive(ctx context.Context, accountID uint, offset, limit int, sort string) ([]model.Contact, int64, error) {
return s.repo.FindActive(ctx, accountID, offset, limit, sort)
}
// ExportCSV writes all contacts for an account as CSV.
// 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)
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 nil
}
// ImportCSVResult holds the result of a CSV import operation.
type ImportCSVResult struct {
Imported int `json:"imported"`
Skipped int `json:"skipped"`
Failed int `json:"failed"`
}
// ImportCSV reads contacts from a CSV reader and creates them.
// POST /api/v1/accounts/:id/contacts/import
// Reference: Chatwoot contacts#import
func (s *ContactService) ImportCSV(ctx context.Context, accountID uint, r io.Reader) (*ImportCSVResult, error) {
csvReader := csv.NewReader(r)
// Read header row
header, err := csvReader.Read()
if err != nil {
return nil, fmt.Errorf("failed to read CSV header: %w", err)
}
// Build column index map
colIndex := make(map[string]int)
for i, col := range header {
colIndex[col] = i
}
result := &ImportCSVResult{}
for {
row, err := csvReader.Read()
if err == io.EOF {
break
}
if err != nil {
result.Failed++
continue
}
contact := &model.Contact{
AccountID: accountID,
}
if idx, ok := colIndex["name"]; ok && idx < len(row) {
contact.Name = row[idx]
}
if idx, ok := colIndex["email"]; ok && idx < len(row) {
contact.Email = row[idx]
}
if idx, ok := colIndex["phone_number"]; ok && idx < len(row) {
contact.PhoneNumber = row[idx]
}
if idx, ok := colIndex["identifier"]; ok && idx < len(row) {
contact.Identifier = row[idx]
}
if idx, ok := colIndex["country_code"]; ok && idx < len(row) {
contact.CountryCode = row[idx]
}
if idx, ok := colIndex["location"]; ok && idx < len(row) {
contact.Location = row[idx]
}
if idx, ok := colIndex["contact_type"]; ok && idx < len(row) {
contact.ContactType = row[idx]
}
// Skip rows without name (required field)
if contact.Name == "" {
result.Skipped++
continue
}
// Check for duplicate by email
if contact.Email != "" {
existing, _ := s.repo.FindByEmail(ctx, accountID, contact.Email)
if existing != nil {
result.Skipped++
continue
}
}
contact.CustomAttributes = datatypes.JSON("{}")
contact.AdditionalAttributes = datatypes.JSON("{}")
if err := s.repo.Create(ctx, contact); err != nil {
applogger.L().Errorf("Failed to import contact row: %v", err)
result.Failed++
continue
}
result.Imported++
}
return result, nil
}
// Filter retrieves contacts matching advanced filter criteria.
// POST /api/v1/accounts/:id/contacts/filter
// Reference: Chatwoot ContactFilterService#perform — filters by contact_type, source,
// assignee, inbox, labels, status, and applies sort order with pagination.
func (s *ContactService) Filter(ctx context.Context, accountID uint, params repository.ContactFilterParams, offset, limit int) ([]model.Contact, int64, error) {
return s.repo.Filter(ctx, accountID, params, offset, limit)
}
// DeleteCustomAttributes removes all custom attributes from a contact.
// DELETE /api/v1/accounts/:id/contacts/:contact_id/custom_attributes
// Reference: Chatwoot contacts#destroy_custom_attributes
func (s *ContactService) DeleteCustomAttributes(ctx context.Context, accountID, contactID uint) error {
// Verify contact belongs to account
contact, err := s.repo.FindByAccountAndID(ctx, accountID, contactID)
if err != nil {
return errors.New("contact not found")
}
return s.repo.DeleteCustomAttributes(ctx, contact.ID)
}
func (s *ContactService) DestroyCustomAttributes(ctx context.Context, accountID, contactID uint, keys []string) (*model.Contact, error) {
contact, err := s.repo.FindByAccountAndID(ctx, accountID, contactID)
if err != nil {
return nil, errors.New("contact not found")
}
attrs := map[string]any{}
if len(contact.CustomAttributes) > 0 {
_ = json.Unmarshal(contact.CustomAttributes, &attrs)
}
for _, key := range keys {
delete(attrs, key)
}
bytes, _ := json.Marshal(attrs)
contact.CustomAttributes = datatypes.JSON(bytes)
if err := s.repo.Update(ctx, contact); err != nil {
return nil, err
}
s.indexContact(ctx, contact)
return contact, nil
}
func (s *ContactService) DeleteAvatar(ctx context.Context, accountID, contactID uint) (*model.Contact, error) {
contact, err := s.repo.FindByAccountAndID(ctx, accountID, contactID)
if err != nil {
return nil, errors.New("contact not found")
}
contact.AvatarURL = ""
if err := s.repo.Update(ctx, contact); err != nil {
return nil, err
}
s.indexContact(ctx, contact)
return contact, nil
}
func (s *ContactService) GetLabels(ctx context.Context, accountID, contactID uint) ([]string, error) {
if _, err := s.repo.FindByAccountAndID(ctx, accountID, contactID); err != nil {
return nil, errors.New("contact not found")
}
var rows []struct{ Name string }
err := s.DB().WithContext(ctx).Table("contact_labels").
Select("tags.name").
Joins("JOIN tags ON tags.id = contact_labels.tag_id").
Where("contact_labels.account_id = ? AND contact_labels.contact_id = ? AND tags.deleted_at IS NULL", accountID, contactID).
Order("contact_labels.created_at ASC, tags.name ASC").
Scan(&rows).Error
if err != nil {
return nil, err
}
labels := make([]string, 0, len(rows))
for _, row := range rows {
labels = append(labels, row.Name)
}
return labels, nil
}
func (s *ContactService) UpdateLabels(ctx context.Context, accountID, contactID uint, labels []string) ([]string, error) {
if _, err := s.repo.FindByAccountAndID(ctx, accountID, contactID); err != nil {
return nil, errors.New("contact not found")
}
normalized := normalizeContactServiceLabels(labels)
err := s.DB().WithContext(ctx).Transaction(func(tx *gorm.DB) error {
if err := tx.Where("account_id = ? AND contact_id = ?", accountID, contactID).Delete(&model.ContactLabel{}).Error; err != nil {
return err
}
for _, label := range normalized {
tag := model.Tag{AccountID: accountID, Name: label}
if err := tx.Where("account_id = ? AND name = ?", accountID, label).FirstOrCreate(&tag).Error; err != nil {
return err
}
contactLabel := model.ContactLabel{AccountID: accountID, ContactID: contactID, TagID: tag.ID}
if err := tx.Create(&contactLabel).Error; err != nil {
return err
}
}
return nil
})
if err != nil {
return nil, err
}
return normalized, nil
}
func normalizeContactServiceLabels(labels []string) []string {
seen := map[string]struct{}{}
result := make([]string, 0, len(labels))
for _, label := range labels {
label = strings.TrimSpace(label)
if label == "" {
continue
}
if _, ok := seen[label]; ok {
continue
}
seen[label] = struct{}{}
result = append(result, label)
}
return result
}
func mergeContactJSON(current datatypes.JSON, incoming datatypes.JSON) datatypes.JSON {
if len(incoming) == 0 {
return current
}
merged := map[string]any{}
if len(current) > 0 {
_ = json.Unmarshal(current, &merged)
}
incomingMap := map[string]any{}
_ = json.Unmarshal(incoming, &incomingMap)
for key, value := range incomingMap {
merged[key] = value
}
bytes, _ := json.Marshal(merged)
return datatypes.JSON(bytes)
}
// ContactableInbox represents an inbox that a contact can be associated with.
// Reference: Chatwoot contacts#contactable_inboxes
type ContactableInbox struct {
Inbox model.Inbox `json:"inbox"`
ContactInbox *model.ContactInbox `json:"contact_inbox,omitempty"`
SourceID string `json:"source_id,omitempty"`
}
// GetContactableInboxes returns inboxes that a contact can be added to.
// GET /api/v1/accounts/:id/contacts/:contact_id/contactable_inboxes
// Reference: Chatwoot contacts#contactable_inboxes — returns inboxes in the account
// that the contact is either already in or can be added to.
func (s *ContactService) GetContactableInboxes(ctx context.Context, accountID, contactID uint) ([]ContactableInbox, error) {
// Verify contact belongs to account
_, err := s.repo.FindByAccountAndID(ctx, accountID, contactID)
if err != nil {
return nil, errors.New("contact not found")
}
// Get existing contact_inboxes for this contact
existingInboxes, err := s.contactInboxSvc.ListByContact(ctx, contactID)
if err != nil {
return nil, fmt.Errorf("failed to list contact inboxes: %w", err)
}
result := make([]ContactableInbox, 0, len(existingInboxes))
for _, ci := range existingInboxes {
result = append(result, ContactableInbox{
Inbox: ci.Inbox,
ContactInbox: &ci,
SourceID: ci.SourceID,
})
}
return result, nil
}