diff --git a/cmd/cli/specs.go b/cmd/cli/specs.go index df2fe7033..050712646 100644 --- a/cmd/cli/specs.go +++ b/cmd/cli/specs.go @@ -318,6 +318,46 @@ the hold. Resuming a lead that is not held succeeds and changes nothing.`, {Name: "contact", Help: "The contact's id"}, }, }, + { + Name: "lead-cc", Short: "Contacts copied on one lead's emails", + Method: http.MethodGet, Path: "/campaigns/{id}/leads/{contact}/cc", + Args: []argSpec{ + {Name: "id", Help: "The campaign's id"}, + {Name: "contact", Help: "The lead's contact id"}, + }, + Table: output.Table{Root: "cc", Columns: []output.Column{ + col("CONTACT", "contact_id"), col("EMAIL", "email"), col("STATUS", "status"), + }, Empty: "This lead copies nobody."}, + }, + { + Name: "set-lead-cc", Short: "Replace the contacts copied on one lead's emails", + Long: `Copy up to two contacts on every email this campaign sends one lead, follow-ups +included, so colleagues at one company share a single thread. The list replaces +the current one. A copied contact who is also a lead of the campaign has their +own sequence held while any lead copies them, so they never get two threads.`, + Example: " $ warmbly campaign set-lead-cc CAMPAIGN_ID CONTACT_ID --cc COLLEAGUE_ID\n" + + " $ warmbly campaign set-lead-cc CAMPAIGN_ID CONTACT_ID --input '{\"contact_ids\":[]}' # copy nobody", + Method: http.MethodPut, Path: "/campaigns/{id}/leads/{contact}/cc", Body: bodyRequired, + Args: []argSpec{ + {Name: "id", Help: "The campaign's id"}, + {Name: "contact", Help: "The lead's contact id"}, + }, + Flag: []flagSpec{ + {Name: "cc", Help: "A contact id to copy (repeatable, at most 2)", Kind: flagStrings, Key: "contact_ids"}, + }, + Success: "Lead CC replaced.", + }, + { + Name: "lead-cc-suggestions", Short: "The lead's likely colleagues to copy", + Method: http.MethodGet, Path: "/campaigns/{id}/leads/{contact}/cc/suggestions", + Args: []argSpec{ + {Name: "id", Help: "The campaign's id"}, + {Name: "contact", Help: "The lead's contact id"}, + }, + Table: output.Table{Root: "data", Columns: []output.Column{ + col("CONTACT", "contact_id"), col("EMAIL", "email"), col("COMPANY", "company"), col("MATCH", "reason"), + }, Empty: "No contacts share the lead's company or email domain."}, + }, { Name: "logs", Short: "The campaign's send log", Method: http.MethodGet, Path: "/campaigns/{id}/logs", Paginate: true, diff --git a/cmd/warmblyctl/api_resources.go b/cmd/warmblyctl/api_resources.go index 68e01dae8..c9bdc2eeb 100644 --- a/cmd/warmblyctl/api_resources.go +++ b/cmd/warmblyctl/api_resources.go @@ -72,6 +72,11 @@ var apiSpecs = []apiSpec{ {name: "campaign lead-hold", summary: "Whether one lead's flow is held", method: "GET", path: "/campaigns/{id}/leads/{child}/hold", child: "contact"}, {name: "campaign pause-lead", summary: "Hold one lead's flow until a date, or until resumed", method: "POST", path: "/campaigns/{id}/leads/{child}/pause", body: bodyOptional, child: "contact"}, {name: "campaign resume-lead", summary: "Lift one lead's hold now", method: "POST", path: "/campaigns/{id}/leads/{child}/resume", child: "contact"}, + // Contacts copied on every email to one lead. --data carries + // {"contact_ids": ["", ...]}, at most two; [] copies nobody. + {name: "campaign lead-cc", summary: "Contacts copied on one lead's emails", method: "GET", path: "/campaigns/{id}/leads/{child}/cc", child: "contact"}, + {name: "campaign set-lead-cc", summary: "Replace the contacts copied on one lead's emails", method: "PUT", path: "/campaigns/{id}/leads/{child}/cc", body: bodyRequired, child: "contact"}, + {name: "campaign lead-cc-suggestions", summary: "The lead's likely colleagues to copy", method: "GET", path: "/campaigns/{id}/leads/{child}/cc/suggestions", child: "contact"}, // Contacts. {name: "contact list", summary: "List or search contacts; --data carries the filter body", method: "POST", path: "/contacts/search", body: bodyOptional, query: []string{"limit", "cursor"}}, diff --git a/docs/content/docs/api/cli.mdx b/docs/content/docs/api/cli.mdx index 34b00110a..a1f22ab0b 100644 --- a/docs/content/docs/api/cli.mdx +++ b/docs/content/docs/api/cli.mdx @@ -238,7 +238,7 @@ Run `warmbly --help` for the flags, and `warmbly "}, bcc: []string{"crm@acme.test"}}, + } + for _, tc := range []struct { + address string + owner *uuid.UUID + isCopy bool + }{ + {"jonas@acme.test", &copied, true}, + {"BOSS@acme.test", nil, true}, + {"crm@acme.test", nil, true}, + {"forwarded@elsewhere.test", &lead, false}, + } { + owner, isCopy := s.copyBounceOwner(context.Background(), uuid.New(), lead, tc.address) + if isCopy != tc.isCopy || (owner == nil) != (tc.owner == nil) || (owner != nil && *owner != *tc.owner) { + t.Errorf("%s: owner %v copy %v, want %v %v", tc.address, owner, isCopy, tc.owner, tc.isCopy) + } + } +} diff --git a/internal/app/advanced/incoming_reply_test.go b/internal/app/advanced/incoming_reply_test.go index 0e6ee36a0..02b520c12 100644 --- a/internal/app/advanced/incoming_reply_test.go +++ b/internal/app/advanced/incoming_reply_test.go @@ -84,6 +84,17 @@ type incomingReplyProgressRepo struct { receivingSent bool completeErr error advanced *incomingReplyAdvancedRepo + copies []models.CampaignLeadCC + copiedLead *repository.CopiedLeadRef + copiesErr error +} + +func (r *incomingReplyProgressRepo) ListLeadCC(context.Context, uuid.UUID, uuid.UUID) ([]models.CampaignLeadCC, error) { + return r.copies, r.copiesErr +} + +func (r *incomingReplyProgressRepo) LeadForCopiedReply(context.Context, uuid.UUID, uuid.UUID) (*repository.CopiedLeadRef, error) { + return r.copiedLead, nil } func (r *incomingReplyProgressRepo) IsInboundReplySource(context.Context, uuid.UUID, uuid.UUID) (bool, error) { diff --git a/internal/app/advanced/lead_copy_reply_test.go b/internal/app/advanced/lead_copy_reply_test.go new file mode 100644 index 000000000..4415bce26 --- /dev/null +++ b/internal/app/advanced/lead_copy_reply_test.go @@ -0,0 +1,185 @@ +package advanced + +import ( + "context" + "errors" + "strings" + "testing" + "time" + + "github.com/google/uuid" + "github.com/warmbly/warmbly/internal/errx" + "github.com/warmbly/warmbly/internal/models" + "github.com/warmbly/warmbly/internal/repository" +) + +type copyReplyAdvancedRepo struct { + *incomingReplyAdvancedRepo + suppressed []string +} + +func (r *copyReplyAdvancedRepo) UpsertSuppressedRecipient(_ context.Context, s *models.SuppressedRecipient) error { + r.suppressed = append(r.suppressed, strings.ToLower(s.Email)) + return nil +} + +type copyReplyContactRepo struct { + incomingReplyContactRepo +} + +func (copyReplyContactRepo) SetSubscribedByEmail(context.Context, uuid.UUID, string, bool) error { + return nil +} + +type copyReplyProgressRepo struct { + *incomingReplyProgressRepo + heldEverywhere int +} + +func (r *copyReplyProgressRepo) HoldLeadEverywhere(context.Context, uuid.UUID, *time.Time, string, string) ([]uuid.UUID, error) { + r.heldEverywhere++ + return nil, nil +} + +// newCopyReplyService is the incoming-reply harness with the lead copying +// jonas@acme.test, answering in the lead's thread. +func newCopyReplyService(t *testing.T) (*service, *copyReplyProgressRepo, *copyReplyAdvancedRepo, uuid.UUID) { + t.Helper() + orgID, accountID, leadID := uuid.New(), uuid.New(), uuid.New() + account := &models.Email{ID: accountID, OrganizationID: &orgID, Email: "sender@example.test"} + svc, progress := newIncomingReplyService(account, nil, leadID) + progress.copies = []models.CampaignLeadCC{{ContactID: uuid.New(), Email: "jonas@acme.test", Status: models.LeadCCStatusActive}} + adv := ©ReplyAdvancedRepo{incomingReplyAdvancedRepo: progress.advanced} + wrapped := ©ReplyProgressRepo{incomingReplyProgressRepo: progress} + svc.repo = adv + svc.campaignProgressRepo = wrapped + svc.contactRepo = copyReplyContactRepo{svc.contactRepo.(incomingReplyContactRepo)} + return svc, wrapped, adv, accountID +} + +func copyReply(accountID uuid.UUID, subject, body string, inReplyTo []string) *models.EmailMessageStoreData { + return &models.EmailMessageStoreData{ + ID: uuid.New(), EmailID: accountID, Folder: models.FolderInbox, + FromAddr: []string{"Jonas "}, + ToAddr: []string{"sender@example.test"}, + InReplyTo: inReplyTo, + Subject: subject, + Snippet: body, + BodyText: body, + } +} + +// A copy asking to stop is about the copy: the lead they were copied on is +// not suppressed with them. +func TestCopyOptOutSuppressesOnlyTheCopy(t *testing.T) { + svc, _, adv, accountID := newCopyReplyService(t) + if xerr := svc.ProcessIncomingReply(context.Background(), accountID, + copyReply(accountID, "Re: Hello", "Please remove me from your list.", []string{""})); xerr != nil { + t.Fatal(xerr) + } + if len(adv.suppressed) != 1 || adv.suppressed[0] != "jonas@acme.test" { + t.Fatalf("suppressed %v, want only the copy", adv.suppressed) + } +} + +// A copy's away message says nothing about the lead's desk. +func TestCopyOutOfOfficeDoesNotHoldTheLead(t *testing.T) { + svc, progress, _, accountID := newCopyReplyService(t) + if xerr := svc.ProcessIncomingReply(context.Background(), accountID, + copyReply(accountID, "Automatic reply: Hello", "I am out of the office until Monday.", []string{""})); xerr != nil { + t.Fatal(xerr) + } + if progress.heldEverywhere != 0 { + t.Fatalf("the lead was held %d times for a copy's away message", progress.heldEverywhere) + } +} + +// A copy answering further down the thread names no message of ours; the +// mailbox that wrote to the lead ties it back, and it is the lead's reply. +func TestCopyReplyDownThreadCountsForTheLead(t *testing.T) { + svc, progress, _, accountID := newCopyReplyService(t) + copyContact := &models.Contact{ID: uuid.New(), Email: "jonas@acme.test"} + svc.taskRepo = incomingReplyTaskRepo{} + cr := svc.contactRepo.(copyReplyContactRepo) + cr.senderContact = copyContact + svc.contactRepo = cr + progress.copiedLead = &repository.CopiedLeadRef{CampaignID: uuid.New(), ContactID: uuid.New(), SequenceID: uuid.New()} + + if xerr := svc.ProcessIncomingReply(context.Background(), accountID, + copyReply(accountID, "Re: Hello", "Sounds good, let's talk Tuesday.", []string{""})); xerr != nil { + t.Fatal(xerr) + } + if progress.replied != 1 { + t.Fatalf("RecordEmailReplied calls = %d, want the copy's reply counted for the lead", progress.replied) + } +} + +// The link in a copied message cannot say who used it, so it opts out the +// lead and everyone copied; a sequence action is about the lead alone. +func TestUnsubscribeLinkOptsOutTheLeadsCopies(t *testing.T) { + for _, tc := range []struct { + name string + run func(s *service, org, campaign, lead uuid.UUID) *errx.Error + want []string + }{ + {"link", func(s *service, org, campaign, lead uuid.UUID) *errx.Error { + return s.UnsubscribeFromLink(context.Background(), org, campaign, lead, "one_click") + }, []string{"task-contact@example.test", "jonas@acme.test"}}, + {"sequence action", func(s *service, _, campaign, lead uuid.UUID) *errx.Error { + return s.Unsubscribe(context.Background(), campaign, lead) + }, []string{"task-contact@example.test"}}, + } { + svc, _, adv, _ := newCopyReplyService(t) + campaign := svc.campaignRepo.(incomingReplyCampaignRepo).campaign + lead := svc.contactRepo.(copyReplyContactRepo).taskContact.ID + if xerr := tc.run(svc, *campaign.OrganizationID, campaign.ID, lead); xerr != nil { + t.Fatalf("%s: %v", tc.name, xerr) + } + if strings.Join(adv.suppressed, ",") != strings.Join(tc.want, ",") { + t.Fatalf("%s: suppressed %v, want %v", tc.name, adv.suppressed, tc.want) + } + } +} + +// An unreadable copy list is an error, never "not a copy": guessing would +// suppress or hold the lead for what a copy did. +func TestCopyLookupFailureIsNotReadAsTheLead(t *testing.T) { + svc, progress, adv, accountID := newCopyReplyService(t) + progress.copiesErr = errors.New("database unavailable") + if xerr := svc.ProcessIncomingReply(context.Background(), accountID, + copyReply(accountID, "Re: Hello", "Please remove me from your list.", []string{""})); xerr == nil { + t.Fatal("a failed copy lookup was processed as if the sender were not a copy") + } + if len(adv.suppressed) != 0 { + t.Fatalf("suppressed %v on a failed lookup", adv.suppressed) + } +} + +// A fresh message from a copy is not a reply to anything, so it credits no lead. +func TestCopyFreshMessageCreditsNoLead(t *testing.T) { + svc, progress, _, accountID := newCopyReplyService(t) + svc.taskRepo = incomingReplyTaskRepo{} + cr := svc.contactRepo.(copyReplyContactRepo) + cr.senderContact = &models.Contact{ID: uuid.New(), Email: "jonas@acme.test"} + svc.contactRepo = cr + progress.copiedLead = &repository.CopiedLeadRef{CampaignID: uuid.New(), ContactID: uuid.New(), SequenceID: uuid.New()} + + if xerr := svc.ProcessIncomingReply(context.Background(), accountID, + copyReply(accountID, "Quick question", "Unrelated: are you at the fair next week?", nil)); xerr != nil { + t.Fatal(xerr) + } + if progress.replied != 0 { + t.Fatalf("RecordEmailReplied calls = %d, want none for a message that replies to nothing", progress.replied) + } +} + +// The opt-out is not acknowledged while a copy is still sendable. +func TestUnsubscribeFailsWhenCopiesCannotBeRead(t *testing.T) { + svc, progress, _, _ := newCopyReplyService(t) + progress.copiesErr = errors.New("database unavailable") + campaign := svc.campaignRepo.(incomingReplyCampaignRepo).campaign + lead := svc.contactRepo.(copyReplyContactRepo).taskContact.ID + if xerr := svc.UnsubscribeFromLink(context.Background(), *campaign.OrganizationID, campaign.ID, lead, "link"); xerr == nil { + t.Fatal("the unsubscribe succeeded without reaching the lead's copies") + } +} diff --git a/internal/app/advanced/service.go b/internal/app/advanced/service.go index 927299464..d56c599eb 100644 --- a/internal/app/advanced/service.go +++ b/internal/app/advanced/service.go @@ -641,6 +641,72 @@ func (s *service) unsubscribe(ctx context.Context, expectOrg *uuid.UUID, campaig "contact_email": contact.Email, "source": via, }) + + // The link in a message is the same for everyone it copied and cannot say + // who used it, so it opts all of them out. A sequence action is about the + // lead alone. + if via != "action" { + return s.unsubscribeLeadCopies(ctx, *campaign.OrganizationID, campaignID, contactID, contact.Email, via, reason) + } + return nil +} + +// isLeadCopy reports whether sender is one of the contacts copied on the +// lead's emails rather than the lead answering from another address. A failed +// read is an error, never "not a copy", which would charge the lead. +func (s *service) isLeadCopy(ctx context.Context, campaignID, contactID uuid.UUID, leadEmail, sender string) (bool, error) { + if s.campaignProgressRepo == nil || sender == "" || strings.EqualFold(leadEmail, sender) { + return false, nil + } + copies, err := s.campaignProgressRepo.ListLeadCC(ctx, campaignID, contactID) + if err != nil { + return false, err + } + for _, cp := range copies { + if strings.EqualFold(strings.TrimSpace(cp.Email), sender) { + return true, nil + } + } + return false, nil +} + +// unsubscribeLeadCopies suppresses every contact copied on one lead's emails. +// A failure fails the request, so the opt-out is retried rather than +// acknowledged with a copy still sendable; the upserts are idempotent. +func (s *service) unsubscribeLeadCopies(ctx context.Context, orgID, campaignID, contactID uuid.UUID, leadEmail, via, reason string) *errx.Error { + if s.campaignProgressRepo == nil { + return nil + } + copies, err := s.campaignProgressRepo.ListLeadCC(ctx, campaignID, contactID) + if err != nil { + return toErrx(err) + } + for _, cp := range copies { + addr := strings.ToLower(strings.TrimSpace(cp.Email)) + if addr == "" { + continue + } + if err := s.repo.UpsertSuppressedRecipient(ctx, &models.SuppressedRecipient{ + OrganizationID: orgID, + Email: addr, + Kind: models.SuppressionKindEmail, + Reason: reason + " on an email copied to them", + Source: models.DeliverabilityEventUnsubscribe, + CampaignID: &campaignID, + Metadata: map[string]interface{}{"via": via, "copied_on": leadEmail}, + }); err != nil { + return toErrx(err) + } + if err := s.contactRepo.SetSubscribedByEmail(ctx, orgID, addr, false); err != nil { + log.Warn().Err(err).Str("contact_id", cp.ContactID.String()).Msg("unsubscribe: could not clear a copied contact's subscription flag") + } + s.emit(ctx, orgID, models.WebhookEventCampaignUnsubscribed, map[string]any{ + "campaign_id": campaignID.String(), + "contact_id": cp.ContactID.String(), + "contact_email": cp.Email, + "source": via, + }) + } return nil } @@ -1173,6 +1239,10 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid. // contactEmail is the address we mailed, which is not always the one // that answered; an opt-out has to reach both. var contactEmail string + // senderIsCopy is a reply from a contact copied on the lead's emails. It + // counts as the lead's reply, but the copy's own away message or opt-out + // is about the copy, not the lead. + var senderIsCopy bool // First, try exact message threading via In-Reply-To. for _, mid := range msg.InReplyTo { @@ -1223,6 +1293,11 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid. campaignID = ct.CampaignID contactID = ct.ContactID sequenceID = ct.SequenceID + isCopy, cerr := s.isLeadCopy(ctx, *ct.CampaignID, *ct.ContactID, contactEmail, sender) + if cerr != nil { + return toErrx(cerr) + } + senderIsCopy = isCopy break } @@ -1245,6 +1320,29 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid. } } + // A copied contact answering further down the thread (to the lead's own + // reply, say) names no message of ours; the mailbox that wrote to the + // lead is the evidence, and the reply is the lead's. A fresh message with + // no parent is not a reply to anything and credits nobody. + if campaignID == nil && contactID != nil && !referencesCampaignThread && len(msg.InReplyTo) > 0 { + ref, err := s.campaignProgressRepo.LeadForCopiedReply(ctx, *contactID, emailAccountID) + if err != nil { + return toErrx(err) + } + if ref != nil { + lead, lerr := s.contactRepo.GetByID(ctx, ref.ContactID) + if lerr != nil { + return lerr + } + if lead != nil { + campaignID, sequenceID = &ref.CampaignID, &ref.SequenceID + contactID = &ref.ContactID + contactEmail = strings.TrimSpace(lead.Email) + senderIsCopy = true + } + } + } + if campaignID != nil && contactID != nil && sequenceID != nil { cID, ctID, sID := *campaignID, *contactID, *sequenceID @@ -1416,7 +1514,7 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid. } var held *time.Time - if campaignID != nil && contactID != nil && verdict.Class == replyclassify.ClassOutOfOffice && settings.ReplyIntent.HoldOnOutOfOffice { + if campaignID != nil && contactID != nil && !senderIsCopy && verdict.Class == replyclassify.ClassOutOfOffice && settings.ReplyIntent.HoldOnOutOfOffice { held = s.holdForOutOfOffice(ctx, *account.OrganizationID, *contactID, settings.ReplyIntent, msg) } @@ -1446,8 +1544,9 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid. if err := s.contactRepo.SetSubscribedByEmail(ctx, *account.OrganizationID, sender, false); err != nil { log.Warn().Err(err).Msg("reply opt-out: could not clear the contact's subscription flag") } - // Answered from another address: the one we mailed asked to stop too. - if contactEmail != "" && !strings.EqualFold(contactEmail, sender) { + // Answered from another address: the one we mailed asked to stop too, + // unless it was a copy asking for themselves. + if contactEmail != "" && !senderIsCopy && !strings.EqualFold(contactEmail, sender) { _ = s.repo.UpsertSuppressedRecipient(ctx, &models.SuppressedRecipient{ OrganizationID: *account.OrganizationID, Email: strings.ToLower(contactEmail), diff --git a/internal/app/campaign/lead_cc.go b/internal/app/campaign/lead_cc.go new file mode 100644 index 000000000..fabc74e17 --- /dev/null +++ b/internal/app/campaign/lead_cc.go @@ -0,0 +1,111 @@ +package campaign + +import ( + "context" + "errors" + "fmt" + + "github.com/google/uuid" + "github.com/warmbly/warmbly/internal/config" + "github.com/warmbly/warmbly/internal/errx" + "github.com/warmbly/warmbly/internal/models" + "github.com/warmbly/warmbly/internal/repository" +) + +// leadCCSuggestionLimit bounds the colleagues offered when picking copies. +const leadCCSuggestionLimit = 8 + +func (s *campaignService) ListLeadCC(ctx context.Context, orgID, campaignID, contactID uuid.UUID) ([]models.CampaignLeadCC, *errx.Error) { + if s.campaignProgressRepo == nil { + return nil, errx.InternalError() + } + if xerr := s.ownedCampaign(ctx, orgID, campaignID); xerr != nil { + return nil, xerr + } + if _, err := s.campaignProgressRepo.GetLeadHold(ctx, campaignID, contactID); err != nil { + if errors.Is(err, repository.ErrLeadNotInCampaign) { + return nil, errx.New(errx.NotFound, "contact is not a lead of this campaign") + } + return nil, errx.InternalError() + } + cc, err := s.campaignProgressRepo.ListLeadCC(ctx, campaignID, contactID) + if err != nil { + return nil, errx.InternalError() + } + return cc, nil +} + +func (s *campaignService) SetLeadCC(ctx context.Context, orgID, campaignID, contactID uuid.UUID, contactIDs []string) ([]models.CampaignLeadCC, *errx.Error) { + if s.campaignProgressRepo == nil { + return nil, errx.InternalError() + } + if xerr := s.ownedCampaign(ctx, orgID, campaignID); xerr != nil { + return nil, xerr + } + ids := make([]uuid.UUID, 0, len(contactIDs)) + seen := map[uuid.UUID]bool{} + for _, raw := range contactIDs { + id, err := uuid.Parse(raw) + if err != nil { + return nil, errx.New(errx.BadRequest, "contact_ids must be contact ids") + } + if !seen[id] { + seen[id] = true + ids = append(ids, id) + } + } + if len(ids) > config.CampaignLeadMaxCC { + return nil, errx.NewWithIdentifier(errx.BadRequest, "lead_cc_limit", + fmt.Sprintf("A lead can have at most %d contacts copied on their emails", config.CampaignLeadMaxCC)) + } + + before, err := s.campaignProgressRepo.ListLeadCC(ctx, campaignID, contactID) + if err != nil { + return nil, errx.InternalError() + } + switch err := s.campaignProgressRepo.SetLeadCC(ctx, orgID, campaignID, contactID, ids); { + case err == nil: + case errors.Is(err, repository.ErrLeadNotInCampaign): + return nil, errx.New(errx.NotFound, "contact is not a lead of this campaign") + case errors.Is(err, repository.ErrLeadCCSelf): + return nil, errx.NewWithIdentifier(errx.BadRequest, "lead_cc_self", "A lead cannot be copied on their own emails") + case errors.Is(err, repository.ErrLeadCCContactNotFound): + return nil, errx.NewWithIdentifier(errx.NotFound, "lead_cc_contact_not_found", "A contact to copy was not found in this workspace") + case errors.Is(err, repository.ErrLeadCCLeadIsCopied): + return nil, errx.NewWithIdentifier(errx.Conflict, "lead_cc_lead_is_copied", + "This lead is copied on another lead's emails in this campaign, so it sends none of its own to copy anyone on") + case errors.Is(err, repository.ErrLeadCCHasCopies): + return nil, errx.NewWithIdentifier(errx.Conflict, "lead_cc_has_copies", + "A contact to copy has contacts copied on their own emails in this campaign; remove those first") + default: + return nil, errx.InternalError() + } + + // A removed copy's own lead is released and may be due now. + for _, b := range before { + if !seen[b.ContactID] { + s.WakeCampaigns(ctx, orgID, []string{campaignID.String()}) + break + } + } + + cc, err := s.campaignProgressRepo.ListLeadCC(ctx, campaignID, contactID) + if err != nil { + return nil, errx.InternalError() + } + return cc, nil +} + +func (s *campaignService) SuggestLeadCC(ctx context.Context, orgID, campaignID, contactID uuid.UUID) ([]models.CampaignLeadCCSuggestion, *errx.Error) { + if s.campaignProgressRepo == nil { + return nil, errx.InternalError() + } + if xerr := s.ownedCampaign(ctx, orgID, campaignID); xerr != nil { + return nil, xerr + } + out, err := s.campaignProgressRepo.SuggestLeadCC(ctx, orgID, campaignID, contactID, leadCCSuggestionLimit) + if err != nil { + return nil, errx.InternalError() + } + return out, nil +} diff --git a/internal/app/campaign/service.go b/internal/app/campaign/service.go index 72722845a..2b6b6cdbe 100644 --- a/internal/app/campaign/service.go +++ b/internal/app/campaign/service.go @@ -89,6 +89,14 @@ type CampaignService interface { ResumeLead(ctx context.Context, orgID, campaignID, contactID uuid.UUID) *errx.Error // GetLeadHold reads the live hold on one lead (nil when it is not held). GetLeadHold(ctx context.Context, orgID, campaignID, contactID uuid.UUID) (*models.LeadHold, *errx.Error) + + // ListLeadCC reads the contacts copied on every email to one lead. + ListLeadCC(ctx context.Context, orgID, campaignID, contactID uuid.UUID) ([]models.CampaignLeadCC, *errx.Error) + // SetLeadCC replaces the contacts copied on one lead and returns the new + // list. A copied contact's own lead in the campaign is held meanwhile. + SetLeadCC(ctx context.Context, orgID, campaignID, contactID uuid.UUID, contactIDs []string) ([]models.CampaignLeadCC, *errx.Error) + // SuggestLeadCC offers the lead's likely colleagues to copy. + SuggestLeadCC(ctx context.Context, orgID, campaignID, contactID uuid.UUID) ([]models.CampaignLeadCCSuggestion, *errx.Error) } // Bounds on a manual lead hold. A hold in the past would lift the moment it diff --git a/internal/app/consumer/event_send_result.go b/internal/app/consumer/event_send_result.go index 0d6eb3ae2..fa137dc59 100644 --- a/internal/app/consumer/event_send_result.go +++ b/internal/app/consumer/event_send_result.go @@ -15,6 +15,7 @@ import ( "github.com/warmbly/warmbly/internal/errx" "github.com/warmbly/warmbly/internal/infrastructure/pubsub" "github.com/warmbly/warmbly/internal/models" + "github.com/warmbly/warmbly/internal/pkg/mailhdr" "github.com/warmbly/warmbly/internal/repository" ) @@ -167,7 +168,7 @@ func (s *JobsService) HandleEmailFailed(ctx context.Context, result models.SendE switch task.TaskType { case "campaign": - return s.failCampaignSend(ctx, task, reason, code, nil) + return s.failCampaignSend(ctx, task, reason, code, refusedRecipient(result), nil) case "email": s.notifyUserSendFailed(ctx, task, reason) case "placement": @@ -194,7 +195,7 @@ func (s *JobsService) failWarmupSend(ctx context.Context, task *repository.Task, // day the send was counted against, for giving the daily counters back; nil // reads it off the task, which is right for a worker result but not for the // reclaimer, whose sends can be counted on an earlier day. -func (s *JobsService) failCampaignSend(ctx context.Context, task *repository.Task, reason, code string, countedOn *time.Time) error { +func (s *JobsService) failCampaignSend(ctx context.Context, task *repository.Task, reason, code, refused string, countedOn *time.Time) error { ct, err := s.TaskRepo.GetCampaignTask(ctx, task.ID) if err != nil { return err @@ -216,6 +217,11 @@ func (s *JobsService) failCampaignSend(ctx context.Context, task *repository.Tas } } + // A copy the server refused is not the lead's bounce: it is dropped from + // later emails and the lead's step is retried without it. + copyRefused := code == string(errx.MailErrorCodeRecipientRejected) && refused != "" && + recipient != "" && !strings.EqualFold(mailhdr.Bare(refused), recipient) + // A refusal on the SENDING DOMAIN's authentication is not about this lead: // the recipient received nothing, and every retry from that domain fails // identically until its DNS is fixed. Give the reservation back without @@ -234,14 +240,22 @@ func (s *JobsService) failCampaignSend(ctx context.Context, task *repository.Tas } // Only permanent recipient refusals are bounce evidence; deferrals may use the same wording. - if s.Evidence != nil && ct.ContactID != nil && ct.SequenceID != nil && + if s.Evidence != nil && ct.ContactID != nil && ct.SequenceID != nil && !copyRefused && code != string(errx.MailErrorCodeServerUnreachable) && emailverify.NamesRecipient(reason) { s.Evidence.RecordEvidence(ctx, *ct.ContactID, models.Step(&campaignID, ct.SequenceID), "bounced_recipient", "send:"+ct.SequenceID.String(), reason) } + // A refused copy that the retry will leave off costs the lead no attempt; + // one that could not be recorded is counted, so it cannot loop forever. + copyExcluded := copyRefused && s.recordRefusedCopy(ctx, task, ct, campaign, refused, reason) + attempts, exhausted, rolledBack := 0, false, false if ct.ContactID != nil && ct.SequenceID != nil && s.CampaignProgressRepo != nil { - attempts, exhausted, rolledBack, err = s.CampaignProgressRepo.RecordSendFailure(ctx, campaignID, *ct.ContactID, *ct.SequenceID, reason) + if copyExcluded { + attempts, exhausted, rolledBack, err = s.CampaignProgressRepo.WalkBackSend(ctx, campaignID, *ct.ContactID, *ct.SequenceID, reason, false) + } else { + attempts, exhausted, rolledBack, err = s.CampaignProgressRepo.RecordSendFailure(ctx, campaignID, *ct.ContactID, *ct.SequenceID, reason) + } if err != nil { return err } @@ -271,7 +285,7 @@ func (s *JobsService) failCampaignSend(ctx context.Context, task *repository.Tas // reputation, so it goes through the bounce pipeline (progress, optional // suppression, guardrails, warmup health, webhooks) and the lead is // dropped as bounced instead of being offered again. - if rolledBack && code == string(errx.MailErrorCodeRecipientRejected) { + if rolledBack && !copyRefused && code == string(errx.MailErrorCodeRecipientRejected) { if s.recordSynchronousBounce(ctx, task, ct, campaign, recipient, reason) { s.logCampaignSendFailure(ctx, campaignID, ct, recipient, reason, code, attempts, false, false, false) s.publishCampaignUpdated(ctx, campaign, campaignID, "") @@ -335,6 +349,52 @@ func (s *JobsService) recordSynchronousBounce(ctx context.Context, task *reposit return true } +// recordRefusedCopy feeds a copied address the server refused at RCPT into +// the bounce pipeline under its own name, and reports whether the next send +// is sure to leave it off: a lead's copy through its bounced mark, a +// campaign-wide one through the recorded bounce the send path reads. +func (s *JobsService) recordRefusedCopy(ctx context.Context, task *repository.Task, ct *repository.CampaignTask, campaign *models.Campaign, refused, reason string) bool { + if campaign == nil || campaign.OrganizationID == nil || ct.CampaignID == nil || ct.ContactID == nil { + return false + } + address := strings.ToLower(mailhdr.Bare(refused)) + var owner *uuid.UUID + if s.CampaignProgressRepo != nil { + id, err := s.CampaignProgressRepo.MarkLeadCCBounced(ctx, *ct.CampaignID, *ct.ContactID, address) + if err != nil { + log.Warn().Err(err).Str("task_id", task.ID.String()).Msg("could not mark a refused copy bounced") + } + owner = id + } + if s.AdvancedService == nil { + return owner != nil + } + taskID := task.ID + req := &models.IngestDeliverabilityEventRequest{ + EventType: models.DeliverabilityEventBounce, + Provider: "smtp_reject", + TaskID: &taskID, + CampaignID: ct.CampaignID, + ContactID: owner, + RecipientEmail: address, + Reason: reason, + IdempotencyKey: "reject:" + taskID.String() + ":" + address, + } + if xerr := s.AdvancedService.IngestDeliverabilityEvent(ctx, *campaign.OrganizationID, req); xerr != nil { + log.Warn().Str("task_id", taskID.String()).Str("error", xerr.Message).Msg("could not record a refused copy as a bounce") + return owner != nil + } + return true +} + +// refusedRecipient is the address the server refused, when the worker knew. +func refusedRecipient(result models.SendEmailResult) string { + if result.Error == nil { + return "" + } + return result.Error.Recipient +} + // publishCampaignUpdated pulses the campaign for every teammate (status "" keeps // the dashboard's status as is). func (s *JobsService) publishCampaignUpdated(ctx context.Context, campaign *models.Campaign, campaignID uuid.UUID, status string) { diff --git a/internal/app/consumer/refused_copy_test.go b/internal/app/consumer/refused_copy_test.go new file mode 100644 index 000000000..c4f8f085b --- /dev/null +++ b/internal/app/consumer/refused_copy_test.go @@ -0,0 +1,138 @@ +package jobs + +import ( + "context" + "testing" + "time" + + "github.com/google/uuid" + + "github.com/warmbly/warmbly/internal/app/advanced" + "github.com/warmbly/warmbly/internal/errx" + "github.com/warmbly/warmbly/internal/models" + "github.com/warmbly/warmbly/internal/repository" +) + +type copyContactRepo struct { + repository.ContactRepository + email string +} + +func (r copyContactRepo) GetByID(_ context.Context, id uuid.UUID) (*models.Contact, *errx.Error) { + return &models.Contact{ID: id, Email: r.email}, nil +} + +type copyCampaignRepo struct { + repository.CampaignRepository + org uuid.UUID +} + +func (r copyCampaignRepo) GetByID(_ context.Context, id uuid.UUID) (*models.Campaign, error) { + return &models.Campaign{ID: id, OrganizationID: &r.org, Status: "active"}, nil +} + +func (copyCampaignRepo) DecrementCampaignDailySend(context.Context, uuid.UUID, time.Time, bool) error { + return nil +} + +type copyProgressRepo struct { + repository.CampaignProgressRepository + copyID uuid.UUID + bounced []string + counted []bool +} + +func (r *copyProgressRepo) RecordSendFailure(context.Context, uuid.UUID, uuid.UUID, uuid.UUID, string) (int, bool, bool, error) { + r.counted = append(r.counted, true) + return 1, false, true, nil +} + +func (r *copyProgressRepo) WalkBackSend(_ context.Context, _, _, _ uuid.UUID, _ string, count bool) (int, bool, bool, error) { + r.counted = append(r.counted, count) + return 0, false, true, nil +} + +func (*copyProgressRepo) HasSentSteps(context.Context, uuid.UUID, uuid.UUID) (bool, error) { + return false, nil +} + +func (r *copyProgressRepo) MarkLeadCCBounced(_ context.Context, _, _ uuid.UUID, address string) (*uuid.UUID, error) { + r.bounced = append(r.bounced, address) + return &r.copyID, nil +} + +type copyAdvanced struct { + advanced.Service + events []*models.IngestDeliverabilityEventRequest +} + +func (a *copyAdvanced) IngestDeliverabilityEvent(_ context.Context, _ uuid.UUID, req *models.IngestDeliverabilityEventRequest) *errx.Error { + a.events = append(a.events, req) + return nil +} + +// A copy the server refused at RCPT bounces the copy, never the lead: no +// address evidence against the lead, and the bounce event names the copy. +func TestRefusedCopyIsNotTheLeadsBounce(t *testing.T) { + for _, tc := range []struct { + name string + refused string + wantCopy bool + }{ + {"a copy", "Jonas ", true}, + {"the lead", "ana@acme.test", false}, + } { + campaign, lead, step, taskID := uuid.New(), uuid.New(), uuid.New(), uuid.New() + progress := ©ProgressRepo{copyID: uuid.New()} + adv := ©Advanced{} + ev := &recordingEvidence{} + s := &JobsService{ + TaskRepo: &evidenceTaskRepo{ + task: &repository.Task{ID: taskID, TaskType: "campaign", EmailAccountID: uuid.New(), Status: "completed"}, + ct: &repository.CampaignTask{TaskID: taskID, CampaignID: &campaign, ContactID: &lead, SequenceID: &step}, + }, + CampaignRepo: copyCampaignRepo{org: uuid.New()}, + CampaignProgressRepo: progress, + ContactRepo: copyContactRepo{email: "ana@acme.test"}, + AdvancedService: adv, + Evidence: ev, + } + err := s.HandleEmailFailed(context.Background(), models.SendEmailResult{ + TaskID: taskID, + Error: &models.EmailSendError{ + Code: string(errx.MailErrorCodeRecipientRejected), + Message: `The mail server rejected the recipient: 550 "5.1.1 no such user"`, + Recipient: tc.refused, + }, + }) + if err != nil { + t.Fatalf("%s: %v", tc.name, err) + } + if len(adv.events) != 1 { + t.Fatalf("%s: %d bounce events, want 1", tc.name, len(adv.events)) + } + got := adv.events[0] + if tc.wantCopy { + if len(ev.kinds) != 0 { + t.Fatalf("%s: evidence %v recorded against the lead", tc.name, ev.kinds) + } + if got.RecipientEmail != "jonas@acme.test" || got.ContactID == nil || *got.ContactID != progress.copyID { + t.Fatalf("%s: bounce event %+v, want it on the copy", tc.name, got) + } + if len(progress.bounced) != 1 { + t.Fatalf("%s: copy marked bounced %v times, want once", tc.name, progress.bounced) + } + // The retry leaves the copy off, so the lead is not charged an attempt. + if len(progress.counted) != 1 || progress.counted[0] { + t.Fatalf("%s: walk-backs %v, want one that counts no attempt", tc.name, progress.counted) + } + continue + } + if len(progress.counted) != 1 || !progress.counted[0] { + t.Fatalf("%s: walk-backs %v, want the lead's attempt counted", tc.name, progress.counted) + } + if got.RecipientEmail != "ana@acme.test" || got.ContactID == nil || *got.ContactID != lead || len(progress.bounced) != 0 { + t.Fatalf("%s: bounce event %+v (copies marked %v), want the lead's own bounce", tc.name, got, progress.bounced) + } + } +} diff --git a/internal/app/consumer/stuck_send_reclaimer.go b/internal/app/consumer/stuck_send_reclaimer.go index d2c66cfbd..084525709 100644 --- a/internal/app/consumer/stuck_send_reclaimer.go +++ b/internal/app/consumer/stuck_send_reclaimer.go @@ -109,7 +109,7 @@ func (s *JobsService) reclaimStuckSend(ctx context.Context, d repository.StuckDi return "", err } dispatchedAt := d.DispatchedAt - if err := s.failCampaignSend(ctx, task, reason, "SEND_OUTCOME_LOST", &dispatchedAt); err != nil { + if err := s.failCampaignSend(ctx, task, reason, "SEND_OUTCOME_LOST", "", &dispatchedAt); err != nil { return "", err } return "reclaimed", nil diff --git a/internal/app/contact/campaign_state.go b/internal/app/contact/campaign_state.go index 12e273796..54b9b1600 100644 --- a/internal/app/contact/campaign_state.go +++ b/internal/app/contact/campaign_state.go @@ -114,8 +114,15 @@ func holdCopy(hold *models.LeadHold) string { return "Paused for this contact" } what := "Paused for this contact" - if hold.Source == models.LeadHoldSourceOutOfOffice { + switch hold.Source { + case models.LeadHoldSourceOutOfOffice: what = "Out of office" + case models.LeadHoldSourceCC: + // The reason is the lead they are copied on; no end is expected. + if r := strings.TrimSpace(hold.Reason); r != "" { + return "Copied on the emails to " + r + ", so none of their own are sent" + } + return "Copied on another lead's emails, so none of their own are sent" } if r := strings.TrimSpace(hold.Reason); r != "" { what += ": " + r diff --git a/internal/app/orgtransfer/spec.go b/internal/app/orgtransfer/spec.go index cba564291..e06c45aec 100644 --- a/internal/app/orgtransfer/spec.go +++ b/internal/app/orgtransfer/spec.go @@ -462,6 +462,12 @@ var Tables = []Table{ Name: "campaign_segments", Group: models.OrgDataGroupCampaigns, Scope: `campaign_id IN ` + orgCampaigns, }, + { + // Both contacts are in the contacts group, which campaigns require. + Name: "campaign_lead_cc", Group: models.OrgDataGroupCampaigns, + Scope: `campaign_id IN ` + orgCampaigns, + Note: "Must travel with the leads, or a copied contact held on their own lead is released into a second sequence.", + }, { Name: "campaign_lead_removals", Group: models.OrgDataGroupCampaigns, Scope: `campaign_id IN ` + orgCampaigns, diff --git a/internal/app/worker/wmail/send.go b/internal/app/worker/wmail/send.go index 0b2b0b1ad..e74552912 100644 --- a/internal/app/worker/wmail/send.go +++ b/internal/app/worker/wmail/send.go @@ -563,5 +563,6 @@ func MailErrorToSendError(err *errx.MailError) *models.EmailSendError { UserTitle: userInfo.Title, UserMessage: userInfo.Message, ActionRequired: userInfo.ActionRequired, + Recipient: err.Recipient, } } diff --git a/internal/client/smtpimap/smtp/client.go b/internal/client/smtpimap/smtp/client.go index 1feca07fd..035093b4b 100644 --- a/internal/client/smtpimap/smtp/client.go +++ b/internal/client/smtpimap/smtp/client.go @@ -475,7 +475,9 @@ func (c *Client) sendRaw(ctx context.Context, from string, to []string, data []b if !permanentReply(err) { return errx.ErrMailServerUnreachableAt("rcpt to", err) } - return errx.ErrMailRecipientRejected(err.Error()) + refused := errx.ErrMailRecipientRejected(err.Error()) + refused.Recipient = r + return refused } } w, err := client.Data() diff --git a/internal/client/smtpimap/smtp/server_test.go b/internal/client/smtpimap/smtp/server_test.go index 8ec23ba3c..29b51f4b8 100644 --- a/internal/client/smtpimap/smtp/server_test.go +++ b/internal/client/smtpimap/smtp/server_test.go @@ -412,4 +412,8 @@ func TestRecipientRejectionCarriesTheServersReason(t *testing.T) { if !strings.Contains(err.Message, "no such user here") { t.Errorf("message = %q, want the server's own reason in it", err.Message) } + // Named apart from the message, so a refused copy is not read as the lead. + if err.Recipient != "to@example.test" { + t.Errorf("recipient = %q, want the refused address", err.Recipient) + } } diff --git a/internal/config/constants.go b/internal/config/constants.go index 53f02cc57..bb2ed1a49 100644 --- a/internal/config/constants.go +++ b/internal/config/constants.go @@ -159,6 +159,10 @@ const ( // tick retries it; this bounds that loop for a mailbox that can never send. CampaignSendMaxAttempts = 5 + // CampaignLeadMaxCC caps the contacts copied on one lead's emails. Every + // copy is one more recipient who did not ask for the email. + CampaignLeadMaxCC = 2 + // CampaignNotDueGraceSeconds is how far in the future a step's hard // constraints (wait_after, start date, sending window, mailbox min-gap) // may sit while a firing task still sends it. Beyond this the scheduler diff --git a/internal/errx/email.go b/internal/errx/email.go index ed067ae34..bdbcf43af 100644 --- a/internal/errx/email.go +++ b/internal/errx/email.go @@ -109,6 +109,8 @@ type MailError struct { ResolvedAt *time.Time `json:"resolved_at"` Message string `json:"message"` + // Recipient is the one address a per-recipient refusal was about. + Recipient string `json:"recipient,omitempty"` // RetryAfter is provider guidance for transient throttles. It stays local // to the worker; persisted error records should not depend on a stale delay. diff --git a/internal/infrastructure/db/migrations/000230_campaign_lead_cc.down.sql b/internal/infrastructure/db/migrations/000230_campaign_lead_cc.down.sql new file mode 100644 index 000000000..619770f83 --- /dev/null +++ b/internal/infrastructure/db/migrations/000230_campaign_lead_cc.down.sql @@ -0,0 +1,16 @@ +DROP TRIGGER IF EXISTS campaign_lead_enrol_cc_hold ON campaign_leads; +DROP FUNCTION IF EXISTS campaign_lead_enrol_cc_hold(); +DROP TRIGGER IF EXISTS campaign_lead_cc_hold ON campaign_lead_cc; +DROP FUNCTION IF EXISTS campaign_lead_cc_hold(); + +DROP TABLE IF EXISTS campaign_lead_cc; + +-- The source is gone, so its holds go with it rather than failing the check. +UPDATE campaign_leads +SET paused_at = NULL, paused_until = NULL, pause_reason = NULL, pause_source = NULL +WHERE pause_source = 'cc'; + +ALTER TABLE public.campaign_leads DROP CONSTRAINT IF EXISTS campaign_leads_pause_source_check; +ALTER TABLE public.campaign_leads + ADD CONSTRAINT campaign_leads_pause_source_check + CHECK (pause_source IS NULL OR pause_source IN ('manual', 'out_of_office', 'inbox_tagging')) NOT VALID; diff --git a/internal/infrastructure/db/migrations/000230_campaign_lead_cc.up.sql b/internal/infrastructure/db/migrations/000230_campaign_lead_cc.up.sql new file mode 100644 index 000000000..6bb5aceb5 --- /dev/null +++ b/internal/infrastructure/db/migrations/000230_campaign_lead_cc.up.sql @@ -0,0 +1,93 @@ +-- Extra recipients copied on every email one campaign sends one lead, so two +-- people at the same company can be reached in one thread instead of two +-- parallel sequences (issue #731). The copies are contacts, not free text, so +-- suppression, bounces and verification apply to them exactly as to a lead. +CREATE TABLE campaign_lead_cc ( + campaign_id uuid NOT NULL, + contact_id uuid NOT NULL, + cc_contact_id uuid NOT NULL REFERENCES contacts (id) ON DELETE CASCADE, + position smallint NOT NULL DEFAULT 0, + -- A bounce attributed to this copy on this lead's thread. It is dropped + -- from later emails whatever the workspace's auto-suppress setting says. + bounced_at timestamptz, + created_at timestamptz NOT NULL DEFAULT now(), + PRIMARY KEY (campaign_id, contact_id, cc_contact_id), + FOREIGN KEY (campaign_id, contact_id) REFERENCES campaign_leads (campaign_id, contact_id) ON DELETE CASCADE, + CONSTRAINT campaign_lead_cc_not_self CHECK (cc_contact_id <> contact_id) +); + +-- Serves the contact FK's cascade and "is this contact copied on a lead here". +CREATE INDEX idx_campaign_lead_cc_cc ON campaign_lead_cc (cc_contact_id, campaign_id); + +-- A contact copied on another lead's thread is reached there, so their own +-- lead in the same campaign is held with source 'cc' instead of starting a +-- second sequence. +ALTER TABLE public.campaign_leads DROP CONSTRAINT IF EXISTS campaign_leads_pause_source_check; +-- NOT VALID: 000231 validates it in its own transaction. +ALTER TABLE public.campaign_leads + ADD CONSTRAINT campaign_leads_pause_source_check + CHECK (pause_source IS NULL OR pause_source IN ('manual', 'out_of_office', 'inbox_tagging', 'cc')) NOT VALID; + +-- Triggers rather than callers, so every path that enrols a lead (segments, +-- imports, contact edits, org transfer) holds a copied contact the same way. +CREATE FUNCTION campaign_lead_cc_hold() RETURNS trigger +LANGUAGE plpgsql AS $$ +BEGIN + IF TG_OP = 'INSERT' THEN + -- A member's own live pause outranks this; anything else is replaced. + UPDATE campaign_leads cl + SET paused_at = CASE + WHEN cl.paused_at IS NOT NULL AND (cl.paused_until IS NULL OR cl.paused_until > NOW()) + THEN cl.paused_at ELSE NOW() END, + paused_until = NULL, + pause_reason = (SELECT c.email FROM contacts c WHERE c.id = NEW.contact_id), + pause_source = 'cc' + WHERE cl.campaign_id = NEW.campaign_id + AND cl.contact_id = NEW.cc_contact_id + AND NOT (cl.pause_source = 'manual' AND cl.paused_at IS NOT NULL + AND (cl.paused_until IS NULL OR cl.paused_until > NOW())); + RETURN NEW; + END IF; + -- Released once no lead in the campaign copies them any more. + UPDATE campaign_leads cl + SET paused_at = NULL, paused_until = NULL, pause_reason = NULL, pause_source = NULL + WHERE cl.campaign_id = OLD.campaign_id + AND cl.contact_id = OLD.cc_contact_id + AND cl.pause_source = 'cc' + AND NOT EXISTS ( + SELECT 1 FROM campaign_lead_cc x + WHERE x.campaign_id = OLD.campaign_id AND x.cc_contact_id = OLD.cc_contact_id + ); + RETURN OLD; +END; +$$; + +CREATE TRIGGER campaign_lead_cc_hold + AFTER INSERT OR DELETE ON campaign_lead_cc + FOR EACH ROW EXECUTE FUNCTION campaign_lead_cc_hold(); + +-- The other order: a contact already copied on a lead is enrolled later. +CREATE FUNCTION campaign_lead_enrol_cc_hold() RETURNS trigger +LANGUAGE plpgsql AS $$ +DECLARE + lead_email text; +BEGIN + SELECT c.email INTO lead_email + FROM campaign_lead_cc x + JOIN contacts c ON c.id = x.contact_id + WHERE x.campaign_id = NEW.campaign_id AND x.cc_contact_id = NEW.contact_id + ORDER BY x.created_at + LIMIT 1; + IF FOUND AND NEW.paused_at IS NULL THEN + NEW.paused_at := NOW(); + NEW.paused_until := NULL; + NEW.pause_reason := lead_email; + NEW.pause_source := 'cc'; + END IF; + RETURN NEW; +END; +$$; + +CREATE TRIGGER campaign_lead_enrol_cc_hold + BEFORE INSERT ON campaign_leads + FOR EACH ROW EXECUTE FUNCTION campaign_lead_enrol_cc_hold(); diff --git a/internal/infrastructure/db/migrations/000231_validate_lead_cc_hold_source.down.sql b/internal/infrastructure/db/migrations/000231_validate_lead_cc_hold_source.down.sql new file mode 100644 index 000000000..cbdfb637a --- /dev/null +++ b/internal/infrastructure/db/migrations/000231_validate_lead_cc_hold_source.down.sql @@ -0,0 +1,3 @@ +-- A validated constraint has no unvalidated form to return to; 000230's down +-- migration replaces it. +SELECT 1; diff --git a/internal/infrastructure/db/migrations/000231_validate_lead_cc_hold_source.up.sql b/internal/infrastructure/db/migrations/000231_validate_lead_cc_hold_source.up.sql new file mode 100644 index 000000000..4662b98e1 --- /dev/null +++ b/internal/infrastructure/db/migrations/000231_validate_lead_cc_hold_source.up.sql @@ -0,0 +1,3 @@ +-- Validates the pause_source CHECK 000230 re-added NOT VALID. A VALIDATE only +-- takes a SHARE UPDATE EXCLUSIVE lock, so writes continue while it scans. +ALTER TABLE public.campaign_leads VALIDATE CONSTRAINT campaign_leads_pause_source_check; diff --git a/internal/models/audit.go b/internal/models/audit.go index 3335b22e8..af75393fc 100644 --- a/internal/models/audit.go +++ b/internal/models/audit.go @@ -53,7 +53,7 @@ const ( AuditEntityCampaign AuditEntityType = "campaign" // AuditEntityCampaignLead is ONE contact inside ONE campaign: the entity id // is the contact and metadata carries the campaign. Written when a member - // pauses or resumes that lead's flow. + // pauses or resumes that lead's flow, or changes who it copies. AuditEntityCampaignLead AuditEntityType = "campaign_lead" AuditEntityContact AuditEntityType = "contact" AuditEntityEmailAccount AuditEntityType = "email_account" diff --git a/internal/models/campaign_lead_cc.go b/internal/models/campaign_lead_cc.go new file mode 100644 index 000000000..ebdc76219 --- /dev/null +++ b/internal/models/campaign_lead_cc.go @@ -0,0 +1,60 @@ +package models + +import ( + "time" + + "github.com/google/uuid" +) + +// CampaignLeadCC is a contact copied on every email one campaign sends one +// lead, so several people at one company share a single thread (issue #731). +type CampaignLeadCC struct { + ContactID uuid.UUID `json:"contact_id"` + Email string `json:"email"` + FirstName string `json:"first_name"` + LastName string `json:"last_name"` + Company string `json:"company,omitempty"` + // Status says whether the next email copies them: one of the + // LeadCCStatus constants. Anything but "active" is left off. + Status string `json:"status"` + BouncedAt *time.Time `json:"bounced_at,omitempty"` +} + +// Copied reports whether the next email to the lead carries this address. +func (c CampaignLeadCC) Copied() bool { return c.Status == LeadCCStatusActive } + +// Why a copy is or is not on the next email, in the order they are decided. +const ( + LeadCCStatusActive = "active" + // LeadCCStatusUnsubscribed is an opted-out or suppressed address. + LeadCCStatusUnsubscribed = "unsubscribed" + // LeadCCStatusBounced is an address that bounced on this thread or on any + // campaign email of its own. + LeadCCStatusBounced = "bounced" + // LeadCCStatusUndeliverable is an address verification refused, under the + // same rule the campaign applies to its leads. + LeadCCStatusUndeliverable = "undeliverable" +) + +// LeadHoldSourceCC holds a contact's own lead while they are copied on another +// lead's thread in the same campaign, so they never get two sequences. Written +// only by the campaign_lead_cc triggers (migration 000230). +const LeadHoldSourceCC = "cc" + +// SetCampaignLeadCC replaces the contacts copied on one lead. An empty list +// removes them all. +type SetCampaignLeadCC struct { + ContactIDs []string `json:"contact_ids"` +} + +// CampaignLeadCCSuggestion is a contact who looks like a colleague of the +// lead, offered first when picking who to copy. Reason is "company" when the +// company names match and "domain" when only the email domain does. +type CampaignLeadCCSuggestion struct { + ContactID uuid.UUID `json:"contact_id"` + Email string `json:"email"` + FirstName string `json:"first_name"` + LastName string `json:"last_name"` + Company string `json:"company,omitempty"` + Reason string `json:"reason"` +} diff --git a/internal/models/contact.go b/internal/models/contact.go index 802d85cd7..d5f0be571 100644 --- a/internal/models/contact.go +++ b/internal/models/contact.go @@ -111,6 +111,8 @@ type ContactCampaignProgress struct { // Hold is the per-lead pause, set only while it is live. Present on any // status: a held lead that has also replied still reads "replied". Hold *LeadHold `json:"hold,omitempty"` + // CC is the contacts copied on every email to this lead in this campaign. + CC []CampaignLeadCC `json:"cc,omitempty"` } // LeadHold is one contact's flow parked inside one campaign. Source is diff --git a/internal/models/contact_campaign_state.go b/internal/models/contact_campaign_state.go index b9c2a2aa5..da025d76d 100644 --- a/internal/models/contact_campaign_state.go +++ b/internal/models/contact_campaign_state.go @@ -37,6 +37,9 @@ type ContactCampaignState struct { // hand. The drawer renders it with "resume now" and "stop" next to it. Hold *LeadHold `json:"hold,omitempty"` + // CC is the contacts copied on every email to this lead in this campaign. + CC []CampaignLeadCC `json:"cc"` + // Next is nil once the flow has ended for the contact; EndedReason says why. Next *ContactNextAction `json:"next,omitempty"` EndedReason string `json:"ended_reason,omitempty"` diff --git a/internal/models/worker.go b/internal/models/worker.go index a28aad697..c16e51810 100644 --- a/internal/models/worker.go +++ b/internal/models/worker.go @@ -111,6 +111,9 @@ type EmailSendError struct { UserTitle string `json:"user_title,omitempty" avro:"user_title"` UserMessage string `json:"user_message,omitempty" avro:"user_message"` ActionRequired string `json:"action_required,omitempty" avro:"action_required"` + // Recipient is the address a refusal named, when the server refused one + // recipient rather than the message. + Recipient string `json:"recipient,omitempty" avro:"recipient"` } // SendEmailResult is the result from worker after sending email diff --git a/internal/repository/campaign_lead_cc_live_test.go b/internal/repository/campaign_lead_cc_live_test.go new file mode 100644 index 000000000..b184bf60b --- /dev/null +++ b/internal/repository/campaign_lead_cc_live_test.go @@ -0,0 +1,264 @@ +package repository + +import ( + "context" + "errors" + "strings" + "testing" + + "github.com/google/uuid" + + "github.com/warmbly/warmbly/internal/models" +) + +// Contacts copied on one lead's emails (issue #731). The hold that keeps a +// copied contact from getting a second sequence lives in triggers, so it is +// asserted against the real routing query. +// +// WARMBLY_TEST_DB=postgres://warmbly:warmbly@localhost:15432/?sslmode=disable \ +// go test ./internal/repository/ -run LiveLeadCC -v + +func routedIDs(pairs []ContactSequencePair) map[uuid.UUID]bool { + out := map[uuid.UUID]bool{} + for _, p := range pairs { + out[p.ContactID] = true + } + return out +} + +func TestLiveLeadCCHoldsTheCopiedLeadUntilReleased(t *testing.T) { + _, pool := liveContactDB(t) + f := newRoutedPairsFixture(t, pool, 3) + repo := NewCampaignProgressRepository(pool) + ctx := context.Background() + a, b, c := f.leads[0], f.leads[1], f.leads[2] + + if err := repo.SetLeadCC(ctx, f.org, f.campaign, a, []uuid.UUID{b}); err != nil { + t.Fatalf("SetLeadCC: %v", err) + } + pairs, _, _ := f.find(t, nil, 25) + got := routedIDs(pairs) + if !got[a] || got[b] || !got[c] { + t.Fatalf("routed %v; want the lead and the uncopied lead, not the copied one", got) + } + hold, err := repo.GetLeadHold(ctx, f.campaign, b) + if err != nil || hold == nil || hold.Source != models.LeadHoldSourceCC { + t.Fatalf("copied lead's hold = %+v, %v; want a cc hold", hold, err) + } + // Nothing is left to wait for, so a campaign of only copied leads can finish. + if n, err := repo.CountHeldLeads(ctx, f.campaign); err != nil || n != 0 { + t.Fatalf("CountHeldLeads = %d, %v; want 0 for a cc hold", n, err) + } + + cc, err := repo.ListLeadCC(ctx, f.campaign, a) + if err != nil || len(cc) != 1 || cc[0].ContactID != b || !cc[0].Copied() { + t.Fatalf("ListLeadCC = %+v, %v; want the copy, active", cc, err) + } + + if err := repo.SetLeadCC(ctx, f.org, f.campaign, a, nil); err != nil { + t.Fatalf("SetLeadCC to nobody: %v", err) + } + if hold, err := repo.GetLeadHold(ctx, f.campaign, b); err != nil || hold != nil { + t.Fatalf("hold after release = %+v, %v; want none", hold, err) + } + pairs, _, _ = f.find(t, nil, 25) + if !routedIDs(pairs)[b] { + t.Fatal("a released copy was not routed again") + } +} + +// A contact copied first and enrolled afterwards is held on arrival, whichever +// path enrols them. +func TestLiveLeadCCHoldsAContactEnrolledAfterBeingCopied(t *testing.T) { + _, pool := liveContactDB(t) + f := newRoutedPairsFixture(t, pool, 2) + repo := NewCampaignProgressRepository(pool) + ctx := context.Background() + a, b := f.leads[0], f.leads[1] + + if _, err := pool.Exec(ctx, `DELETE FROM campaign_leads WHERE campaign_id = $1 AND contact_id = $2`, f.campaign, b); err != nil { + t.Fatalf("drop lead: %v", err) + } + if err := repo.SetLeadCC(ctx, f.org, f.campaign, a, []uuid.UUID{b}); err != nil { + t.Fatalf("SetLeadCC: %v", err) + } + if _, err := pool.Exec(ctx, `INSERT INTO campaign_leads (campaign_id, contact_id) VALUES ($1, $2)`, f.campaign, b); err != nil { + t.Fatalf("enrol: %v", err) + } + if hold, err := repo.GetLeadHold(ctx, f.campaign, b); err != nil || hold == nil || hold.Source != models.LeadHoldSourceCC { + t.Fatalf("hold on enrol = %+v, %v; want a cc hold", hold, err) + } + // Removing the lead that copies them releases them too. + if _, err := pool.Exec(ctx, `DELETE FROM campaign_leads WHERE campaign_id = $1 AND contact_id = $2`, f.campaign, a); err != nil { + t.Fatalf("drop copying lead: %v", err) + } + if hold, err := repo.GetLeadHold(ctx, f.campaign, b); err != nil || hold != nil { + t.Fatalf("hold after the copying lead left = %+v, %v; want none", hold, err) + } +} + +func TestLiveLeadCCRefusals(t *testing.T) { + _, pool := liveContactDB(t) + f := newRoutedPairsFixture(t, pool, 3) + other := newRoutedPairsFixture(t, pool, 1) + repo := NewCampaignProgressRepository(pool) + ctx := context.Background() + a, b, c := f.leads[0], f.leads[1], f.leads[2] + + cases := []struct { + name string + lead uuid.UUID + cc []uuid.UUID + want error + }{ + {"self", a, []uuid.UUID{a}, ErrLeadCCSelf}, + {"another workspace's contact", a, []uuid.UUID{other.leads[0]}, ErrLeadCCContactNotFound}, + {"not a lead", other.leads[0], nil, ErrLeadNotInCampaign}, + } + for _, tc := range cases { + if err := repo.SetLeadCC(ctx, f.org, f.campaign, tc.lead, tc.cc); !errors.Is(err, tc.want) { + t.Fatalf("%s: err = %v, want %v", tc.name, err, tc.want) + } + } + + if err := repo.SetLeadCC(ctx, f.org, f.campaign, a, []uuid.UUID{b}); err != nil { + t.Fatalf("SetLeadCC: %v", err) + } + if err := repo.SetLeadCC(ctx, f.org, f.campaign, b, []uuid.UUID{c}); !errors.Is(err, ErrLeadCCLeadIsCopied) { + t.Fatalf("copies on a copied lead: err = %v, want ErrLeadCCLeadIsCopied", err) + } + if err := repo.SetLeadCC(ctx, f.org, f.campaign, c, []uuid.UUID{a}); !errors.Is(err, ErrLeadCCHasCopies) { + t.Fatalf("copying a lead with copies: err = %v, want ErrLeadCCHasCopies", err) + } + // A refused write leaves the list as it was. + if cc, _ := repo.ListLeadCC(ctx, f.campaign, a); len(cc) != 1 || cc[0].ContactID != b { + t.Fatalf("list after refusals = %+v; want it unchanged", cc) + } +} + +func TestLiveLeadCCStatusAndBounces(t *testing.T) { + _, pool := liveContactDB(t) + f := newRoutedPairsFixture(t, pool, 3) + repo := NewCampaignProgressRepository(pool) + ctx := context.Background() + a, b, c := f.leads[0], f.leads[1], f.leads[2] + + if err := repo.SetLeadCC(ctx, f.org, f.campaign, a, []uuid.UUID{b, c}); err != nil { + t.Fatalf("SetLeadCC: %v", err) + } + if _, err := pool.Exec(ctx, `UPDATE contacts SET subscribed = false WHERE id = $1`, b); err != nil { + t.Fatalf("unsubscribe: %v", err) + } + var cEmail string + if err := pool.QueryRow(ctx, `SELECT email FROM contacts WHERE id = $1`, c).Scan(&cEmail); err != nil { + t.Fatalf("read email: %v", err) + } + id, err := repo.MarkLeadCCBounced(ctx, f.campaign, a, strings.ToUpper(cEmail)) + if err != nil || id == nil || *id != c { + t.Fatalf("MarkLeadCCBounced = %v, %v; want the copy", id, err) + } + if id, err := repo.MarkLeadCCBounced(ctx, f.campaign, a, "nobody@test.local"); err != nil || id != nil { + t.Fatalf("MarkLeadCCBounced on a stranger = %v, %v; want nil", id, err) + } + + cc, err := repo.ListLeadCC(ctx, f.campaign, a) + if err != nil || len(cc) != 2 { + t.Fatalf("ListLeadCC = %+v, %v", cc, err) + } + want := map[uuid.UUID]string{b: models.LeadCCStatusUnsubscribed, c: models.LeadCCStatusBounced} + for _, x := range cc { + if x.Status != want[x.ContactID] || x.Copied() { + t.Fatalf("copy %s status %q; want %q and not copied", x.ContactID, x.Status, want[x.ContactID]) + } + } +} + +func TestLiveLeadForCopiedReplyFindsTheLeadsThread(t *testing.T) { + _, pool := liveContactDB(t) + f := newRoutedPairsFixture(t, pool, 2) + repo := NewCampaignProgressRepository(pool) + ctx := context.Background() + a, b := f.leads[0], f.leads[1] + + if err := repo.SetLeadCC(ctx, f.org, f.campaign, a, []uuid.UUID{b}); err != nil { + t.Fatalf("SetLeadCC: %v", err) + } + if _, err := pool.Exec(ctx, `UPDATE campaign_leads SET email_account_id = $3 WHERE campaign_id = $1 AND contact_id = $2`, + f.campaign, a, f.mailbox); err != nil { + t.Fatalf("bind sender: %v", err) + } + if _, err := pool.Exec(ctx, `INSERT INTO campaign_contact_progress (campaign_id, contact_id, sequence_id, sent_at) + VALUES ($1, $2, $3, NOW())`, f.campaign, a, f.step); err != nil { + t.Fatalf("progress: %v", err) + } + + ref, err := repo.LeadForCopiedReply(ctx, b, f.mailbox) + if err != nil || ref == nil || ref.CampaignID != f.campaign || ref.ContactID != a || ref.SequenceID != f.step { + t.Fatalf("LeadForCopiedReply = %+v, %v; want the lead's step", ref, err) + } + // Another mailbox never wrote to the lead, so it is no evidence. + if ref, err := repo.LeadForCopiedReply(ctx, b, uuid.New()); err != nil || ref != nil { + t.Fatalf("LeadForCopiedReply from another mailbox = %+v, %v; want nil", ref, err) + } +} + +// The Leads list carries each lead's copies, and a copied lead reads paused. +func TestLiveLeadCCShowsInTheLeadsList(t *testing.T) { + handle, pool := liveContactDB(t) + f := newRoutedPairsFixture(t, pool, 2) + repo := NewCampaignProgressRepository(pool) + contacts := NewContactRepostory(handle) + ctx := context.Background() + a, b := f.leads[0], f.leads[1] + + if err := repo.SetLeadCC(ctx, f.org, f.campaign, a, []uuid.UUID{b}); err != nil { + t.Fatalf("SetLeadCC: %v", err) + } + res, xerr := contacts.Search(ctx, f.org.String(), nil, nil, models.SearchContacts{ + CampaignIDs: []string{f.campaign.String()}, + }, 25) + if xerr != nil { + t.Fatalf("search: %v", xerr) + } + byID := map[uuid.UUID]*models.ContactCampaignProgress{} + for i := range res.Data { + byID[res.Data[i].ID] = res.Data[i].CampaignLead + } + if lead := byID[a]; lead == nil || len(lead.CC) != 1 || lead.CC[0].ContactID != b || lead.CC[0].Status != models.LeadCCStatusActive { + t.Fatalf("the copying lead reads %+v; want its one active copy", lead) + } + if lead := byID[b]; lead == nil || lead.Status != models.LeadStatusPaused || lead.Hold == nil || lead.Hold.Source != models.LeadHoldSourceCC { + t.Fatalf("the copied lead reads %+v; want paused with a cc hold", lead) + } +} + +// A refused copy walks the step back without spending the lead's attempt, and +// a campaign-wide copy that bounced here is reported so the send leaves it off. +func TestLiveLeadCCRefusedCopyCostsNoAttempt(t *testing.T) { + _, pool := liveContactDB(t) + f := newRoutedPairsFixture(t, pool, 1) + repo := NewCampaignProgressRepository(pool) + ctx := context.Background() + lead := f.leads[0] + + if _, err := pool.Exec(ctx, `INSERT INTO campaign_contact_progress (campaign_id, contact_id, sequence_id, sent_at, dispatched_at) + VALUES ($1, $2, $3, NOW(), NOW())`, f.campaign, lead, f.step); err != nil { + t.Fatalf("progress: %v", err) + } + attempts, _, rolled, err := repo.WalkBackSend(ctx, f.campaign, lead, f.step, "copy refused", false) + if err != nil || !rolled || attempts != 0 { + t.Fatalf("WalkBackSend = %d, %v, %v; want rolled back with no attempt", attempts, rolled, err) + } + + if _, err := pool.Exec(ctx, `INSERT INTO deliverability_events (organization_id, campaign_id, event_type, recipient_email, idempotency_key) + VALUES ($1, $2, 'bounce', 'Crm@Acme.test', $3)`, f.org, f.campaign, "test:"+uuid.NewString()); err != nil { + t.Fatalf("event: %v", err) + } + t.Cleanup(func() { + _, _ = pool.Exec(context.Background(), `DELETE FROM deliverability_events WHERE organization_id = $1`, f.org) + }) + got, err := repo.BouncedCopyAddresses(ctx, f.campaign, []string{"crm@acme.test", "boss@acme.test"}) + if err != nil || !got["crm@acme.test"] || got["boss@acme.test"] { + t.Fatalf("BouncedCopyAddresses = %v, %v; want only the bounced address", got, err) + } +} diff --git a/internal/repository/pg_campaign_lead_cc.go b/internal/repository/pg_campaign_lead_cc.go new file mode 100644 index 000000000..b40f33dfd --- /dev/null +++ b/internal/repository/pg_campaign_lead_cc.go @@ -0,0 +1,285 @@ +package repository + +import ( + "context" + "errors" + "fmt" + "strings" + + "github.com/google/uuid" + "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgxpool" + "github.com/warmbly/warmbly/internal/models" +) + +// What SetLeadCC refuses. The service turns each into its own error code. +var ( + // ErrLeadCCContactNotFound is a copy that is not a contact of the workspace. + ErrLeadCCContactNotFound = errors.New("a contact to copy was not found") + // ErrLeadCCSelf is the lead copied on their own emails. + ErrLeadCCSelf = errors.New("a lead cannot be copied on their own emails") + // ErrLeadCCLeadIsCopied is a lead already copied on another lead's thread + // in the campaign; their own emails are held, so copies would reach nobody. + ErrLeadCCLeadIsCopied = errors.New("the lead is copied on another lead in this campaign") + // ErrLeadCCHasCopies is a contact whose own lead copies others: holding it + // would silently strand the people it copies. + ErrLeadCCHasCopies = errors.New("a contact to copy has copies of their own in this campaign") +) + +// CopiedLeadRef names the lead, and the step, a copied contact's reply answers. +type CopiedLeadRef struct { + CampaignID uuid.UUID + ContactID uuid.UUID + SequenceID uuid.UUID +} + +// personalMailDomainsSQL never identify a company, so a shared one is not a +// reason to suggest two contacts are colleagues. +const personalMailDomainsSQL = `'gmail.com','googlemail.com','yahoo.com','yahoo.de','hotmail.com','hotmail.de','outlook.com','outlook.de','live.com','live.de','msn.com','aol.com','icloud.com','me.com','gmx.com','gmx.de','gmx.net','web.de','t-online.de','freenet.de','posteo.de','mailbox.org','proton.me','protonmail.com','mail.com','yandex.com'` + +// leadCCStatusSQL derives models.LeadCCStatus* for a copied contact aliased c +// on the campaign_lead_cc row aliased x. cp is the bound campaign id. +func leadCCStatusSQL(cp string) string { + return `CASE + WHEN c.subscribed IS FALSE + OR recipient_suppressed((SELECT organization_id FROM campaigns WHERE id = ` + cp + `), c.email) + THEN '` + models.LeadCCStatusUnsubscribed + `' + WHEN x.bounced_at IS NOT NULL + OR EXISTS (SELECT 1 FROM campaign_contact_progress b WHERE b.contact_id = c.id AND b.bounced_at IS NOT NULL) + THEN '` + models.LeadCCStatusBounced + `' + WHEN ` + undeliverableClause(cp) + ` THEN '` + models.LeadCCStatusUndeliverable + `' + ELSE '` + models.LeadCCStatusActive + `' + END` +} + +// leadCCSelectSQL lists one lead's copies in the order they were chosen. $1 is +// the campaign, $2 the lead. +func leadCCSelectSQL() string { + return ` + SELECT c.id, c.email, c.first_name, c.last_name, c.company, ` + leadCCStatusSQL("$1") + `, x.bounced_at + FROM campaign_lead_cc x + JOIN contacts c ON c.id = x.cc_contact_id + WHERE x.campaign_id = $1 AND x.contact_id = $2 + ORDER BY x.position, x.created_at` +} + +// listLeadCC is shared by the campaign and contact repositories, so the send +// path and the drawer read one definition. +func listLeadCC(ctx context.Context, pool *pgxpool.Pool, campaignID, contactID uuid.UUID) ([]models.CampaignLeadCC, error) { + rows, err := pool.Query(ctx, leadCCSelectSQL(), campaignID, contactID) + if err != nil { + return nil, err + } + defer rows.Close() + out := []models.CampaignLeadCC{} + for rows.Next() { + var cc models.CampaignLeadCC + if err := rows.Scan(&cc.ContactID, &cc.Email, &cc.FirstName, &cc.LastName, &cc.Company, &cc.Status, &cc.BouncedAt); err != nil { + return nil, err + } + out = append(out, cc) + } + return out, rows.Err() +} + +func (r *campaignProgressRepository) ListLeadCC(ctx context.Context, campaignID, contactID uuid.UUID) ([]models.CampaignLeadCC, error) { + return listLeadCC(ctx, r.db, campaignID, contactID) +} + +func (r *campaignProgressRepository) SetLeadCC(ctx context.Context, orgID, campaignID, contactID uuid.UUID, ccIDs []uuid.UUID) error { + tx, err := r.db.Begin(ctx) + if err != nil { + return err + } + defer tx.Rollback(ctx) + + // One writer per campaign, so two edits cannot build a chain of copies + // that each checked against the other's snapshot. + if _, err := tx.Exec(ctx, `SELECT pg_advisory_xact_lock(hashtextextended('campaign_lead_cc:' || $1::text, 0))`, campaignID); err != nil { + return err + } + + var isLead bool + if err := tx.QueryRow(ctx, ` + SELECT EXISTS ( + SELECT 1 FROM campaign_leads cl + JOIN campaigns cam ON cam.id = cl.campaign_id AND cam.organization_id = $1 + WHERE cl.campaign_id = $2 AND cl.contact_id = $3 + )`, orgID, campaignID, contactID).Scan(&isLead); err != nil { + return err + } + if !isLead { + return ErrLeadNotInCampaign + } + + if len(ccIDs) > 0 { + for _, id := range ccIDs { + if id == contactID { + return ErrLeadCCSelf + } + } + var found int + if err := tx.QueryRow(ctx, + `SELECT COUNT(*) FROM contacts WHERE organization_id = $1 AND id = ANY($2::uuid[])`, + orgID, ccIDs).Scan(&found); err != nil { + return err + } + if found != len(ccIDs) { + return ErrLeadCCContactNotFound + } + var leadIsCopied, ccHasCopies bool + if err := tx.QueryRow(ctx, ` + SELECT + EXISTS (SELECT 1 FROM campaign_lead_cc WHERE campaign_id = $1 AND cc_contact_id = $2), + EXISTS (SELECT 1 FROM campaign_lead_cc WHERE campaign_id = $1 AND contact_id = ANY($3::uuid[]))`, + campaignID, contactID, ccIDs).Scan(&leadIsCopied, &ccHasCopies); err != nil { + return err + } + if leadIsCopied { + return ErrLeadCCLeadIsCopied + } + if ccHasCopies { + return ErrLeadCCHasCopies + } + } + + if _, err := tx.Exec(ctx, ` + DELETE FROM campaign_lead_cc + WHERE campaign_id = $1 AND contact_id = $2 + AND NOT (cc_contact_id = ANY(COALESCE($3::uuid[], '{}')))`, + campaignID, contactID, ccIDs); err != nil { + return err + } + if len(ccIDs) > 0 { + if _, err := tx.Exec(ctx, ` + INSERT INTO campaign_lead_cc (campaign_id, contact_id, cc_contact_id, position) + SELECT $1, $2, u.id, u.ord - 1 + FROM unnest($3::uuid[]) WITH ORDINALITY AS u(id, ord) + ON CONFLICT (campaign_id, contact_id, cc_contact_id) DO UPDATE SET position = EXCLUDED.position`, + campaignID, contactID, ccIDs); err != nil { + return err + } + } + return tx.Commit(ctx) +} + +func (r *campaignProgressRepository) MarkLeadCCBounced(ctx context.Context, campaignID, contactID uuid.UUID, address string) (*uuid.UUID, error) { + address = strings.TrimSpace(address) + if address == "" { + return nil, nil + } + var id uuid.UUID + err := r.db.QueryRow(ctx, ` + UPDATE campaign_lead_cc x + SET bounced_at = COALESCE(x.bounced_at, NOW()) + FROM contacts c + WHERE c.id = x.cc_contact_id + AND x.campaign_id = $1 AND x.contact_id = $2 + AND lower(c.email) = lower($3) + RETURNING x.cc_contact_id`, campaignID, contactID, address).Scan(&id) + if errors.Is(err, pgx.ErrNoRows) { + return nil, nil + } + if err != nil { + return nil, err + } + return &id, nil +} + +func (r *campaignProgressRepository) LeadForCopiedReply(ctx context.Context, ccContactID, emailAccountID uuid.UUID) (*CopiedLeadRef, error) { + var ref CopiedLeadRef + err := r.db.QueryRow(ctx, ` + SELECT p.campaign_id, p.contact_id, p.sequence_id + FROM campaign_lead_cc x + JOIN campaign_leads cl ON cl.campaign_id = x.campaign_id AND cl.contact_id = x.contact_id + JOIN campaign_contact_progress p ON p.campaign_id = x.campaign_id AND p.contact_id = x.contact_id + WHERE x.cc_contact_id = $1 + AND cl.email_account_id = $2 + AND p.sent_at IS NOT NULL + AND `+progressIsEmailStep("p")+` + ORDER BY p.sent_at DESC + LIMIT 1`, ccContactID, emailAccountID).Scan(&ref.CampaignID, &ref.ContactID, &ref.SequenceID) + if errors.Is(err, pgx.ErrNoRows) { + return nil, nil + } + if err != nil { + return nil, err + } + return &ref, nil +} + +func (r *campaignProgressRepository) SuggestLeadCC(ctx context.Context, orgID, campaignID, contactID uuid.UUID, limit int) ([]models.CampaignLeadCCSuggestion, error) { + rows, err := r.db.Query(ctx, fmt.Sprintf(` + WITH lead AS ( + SELECT lower(btrim(company)) AS co, lower(split_part(email, '@', 2)) AS dom + FROM contacts WHERE id = $2 AND organization_id = $1 + ) + SELECT c.id, c.email, c.first_name, c.last_name, c.company, + CASE WHEN lead.co <> '' AND lower(btrim(c.company)) = lead.co THEN 'company' ELSE 'domain' END AS reason + FROM contacts c, lead + WHERE c.organization_id = $1 + AND c.id <> $2 + AND c.subscribed IS NOT FALSE + AND ( + (lead.co <> '' AND lower(btrim(c.company)) = lead.co) + OR (lead.dom <> '' AND lead.dom NOT IN (%s) AND lower(split_part(c.email, '@', 2)) = lead.dom) + ) + AND NOT EXISTS ( + SELECT 1 FROM campaign_lead_cc x + WHERE x.campaign_id = $3 AND x.contact_id = $2 AND x.cc_contact_id = c.id + ) + ORDER BY (lead.co <> '' AND lower(btrim(c.company)) = lead.co) DESC, c.first_name, c.last_name, c.email + LIMIT $4`, personalMailDomainsSQL), orgID, contactID, campaignID, limit) + if err != nil { + return nil, err + } + defer rows.Close() + out := []models.CampaignLeadCCSuggestion{} + for rows.Next() { + var s models.CampaignLeadCCSuggestion + if err := rows.Scan(&s.ContactID, &s.Email, &s.FirstName, &s.LastName, &s.Company, &s.Reason); err != nil { + return nil, err + } + out = append(out, s) + } + return out, rows.Err() +} + +func (r *campaignProgressRepository) BouncedCopyAddresses(ctx context.Context, campaignID uuid.UUID, addresses []string) (map[string]bool, error) { + out := map[string]bool{} + if len(addresses) == 0 { + return out, nil + } + rows, err := r.db.Query(ctx, ` + SELECT DISTINCT lower(recipient_email) + FROM deliverability_events + WHERE campaign_id = $1 AND event_type = 'bounce' + AND lower(recipient_email) = ANY($2::text[])`, campaignID, addresses) + if err != nil { + return nil, err + } + defer rows.Close() + for rows.Next() { + var a string + if err := rows.Scan(&a); err != nil { + return nil, err + } + out[a] = true + } + return out, rows.Err() +} + +// leadCCJSONSQL is one lead's copies as a JSON array for the Leads list, with +// the lead row aliased hl and the campaign bound at cp. +func leadCCJSONSQL(cp string) string { + return `( + SELECT COALESCE(json_agg(json_build_object( + 'contact_id', c.id, 'email', c.email, 'first_name', c.first_name, + 'last_name', c.last_name, 'company', c.company, + 'status', ` + leadCCStatusSQL(cp) + `, 'bounced_at', x.bounced_at + ) ORDER BY x.position, x.created_at), '[]'::json) + FROM campaign_lead_cc x + JOIN contacts c ON c.id = x.cc_contact_id + WHERE x.campaign_id = hl.campaign_id AND x.contact_id = hl.contact_id + )` +} diff --git a/internal/repository/pg_campaign_progress.go b/internal/repository/pg_campaign_progress.go index a35d822ab..235abc6ab 100644 --- a/internal/repository/pg_campaign_progress.go +++ b/internal/repository/pg_campaign_progress.go @@ -165,6 +165,12 @@ type CampaignProgressRepository interface { // A lead with nothing else delivered is unbound from the mailbox that // failed, so rotation can offer it a working one. RecordSendFailure(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID, reason string) (attempts int, exhausted bool, rolledBack bool, err error) + // WalkBackSend is RecordSendFailure with the attempt optionally left + // uncounted, for a failure the retry is known not to repeat. + WalkBackSend(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID, reason string, countAttempt bool) (attempts int, exhausted bool, rolledBack bool, err error) + // BouncedCopyAddresses reports which of the campaign's own CC/BCC + // addresses have bounced on a send of this campaign. + BouncedCopyAddresses(ctx context.Context, campaignID uuid.UUID, addresses []string) (map[string]bool, error) // LastSenderForLead is the mailbox a lead was LAST actually sent from, // read from the campaign tasks that dispatched its steps. It answers the // case campaign_leads.email_account_id cannot: a lead removed from the @@ -328,6 +334,24 @@ type CampaignProgressRepository interface { // including a dated hold that has since expired). Returns // ErrLeadNotInCampaign when the contact is not a lead of the campaign. GetLeadHold(ctx context.Context, campaignID, contactID uuid.UUID) (*models.LeadHold, error) + + // ListLeadCC reads the contacts copied on one lead, each with whether the + // next email carries them. + ListLeadCC(ctx context.Context, campaignID, contactID uuid.UUID) ([]models.CampaignLeadCC, error) + // SetLeadCC replaces the contacts copied on one lead. See the ErrLeadCC + // errors for what it refuses. + SetLeadCC(ctx context.Context, orgID, campaignID, contactID uuid.UUID, ccIDs []uuid.UUID) error + // MarkLeadCCBounced records a bounce on the copy of one lead's thread sent + // to address, and returns that copy's contact, or nil when address is not + // one of the lead's copies. + MarkLeadCCBounced(ctx context.Context, campaignID, contactID uuid.UUID, address string) (*uuid.UUID, error) + // LeadForCopiedReply finds the lead whose thread a copied contact is + // answering in: the latest email step sent from emailAccountID to a lead + // that copies them. Nil when there is none. + LeadForCopiedReply(ctx context.Context, ccContactID, emailAccountID uuid.UUID) (*CopiedLeadRef, error) + // SuggestLeadCC offers the lead's likely colleagues: same company name, or + // the same email domain when that domain is not a personal mail service. + SuggestLeadCC(ctx context.Context, orgID, campaignID, contactID uuid.UUID, limit int) ([]models.CampaignLeadCCSuggestion, error) } // ErrLeadNotInCampaign is returned when a hold is asked for on a contact that @@ -564,6 +588,14 @@ func (r *campaignProgressRepository) ListStuckDispatches(ctx context.Context, ol // so a duplicate worker result after the step was already walked back (or // re-sent) is a no-op. func (r *campaignProgressRepository) RecordSendFailure(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID, reason string) (int, bool, bool, error) { + return r.WalkBackSend(ctx, campaignID, contactID, sequenceID, reason, true) +} + +func (r *campaignProgressRepository) WalkBackSend(ctx context.Context, campaignID, contactID, sequenceID uuid.UUID, reason string, countAttempt bool) (int, bool, bool, error) { + inc := 0 + if countAttempt { + inc = 1 + } if len(reason) > 500 { reason = reason[:500] } @@ -578,7 +610,7 @@ func (r *campaignProgressRepository) RecordSendFailure(ctx context.Context, camp SET sent_at = NULL, dispatched_at = NULL, dispatch_task_id = NULL, - send_attempts = send_attempts + 1, + send_attempts = send_attempts + $5, failed_at = NOW(), failure_reason = $4 WHERE campaign_id = $1 AND contact_id = $2 AND sequence_id = $3 @@ -599,7 +631,7 @@ func (r *campaignProgressRepository) RecordSendFailure(ctx context.Context, camp SELECT send_attempts FROM walked ` var attempts int - err := r.db.QueryRow(ctx, query, campaignID, contactID, sequenceID, reason).Scan(&attempts) + err := r.db.QueryRow(ctx, query, campaignID, contactID, sequenceID, reason, inc).Scan(&attempts) if err != nil { if errors.Is(err, pgx.ErrNoRows) { return 0, false, false, nil @@ -2295,7 +2327,9 @@ func (r *campaignProgressRepository) CountHeldLeads(ctx context.Context, campaig var n int err := r.db.QueryRow(ctx, ` SELECT COUNT(*) FROM campaign_leads cl - WHERE cl.campaign_id = $1 AND `+liveHold("cl"), + WHERE cl.campaign_id = $1 AND `+liveHold("cl")+` + -- Reached in another lead's thread: nothing is left to wait for. + AND cl.pause_source IS DISTINCT FROM 'cc'`, campaignID).Scan(&n) return n, err } diff --git a/internal/repository/pg_contact.go b/internal/repository/pg_contact.go index 187289c6b..6a6bd57db 100644 --- a/internal/repository/pg_contact.go +++ b/internal/repository/pg_contact.go @@ -1513,6 +1513,14 @@ func (r *contactRepository) buildContactFilter(ctx context.Context, orgID string }, nil } +// leadRowJSON is the campaign_leads half of a Leads-list row: the fields read +// together because they come from one lead row. +type leadRowJSON struct { + Sender *string `json:"sender"` + CC []models.CampaignLeadCC `json:"cc"` + Hold *models.LeadHold `json:"hold"` +} + func (r *contactRepository) Search( ctx context.Context, orgID string, @@ -1650,6 +1658,8 @@ func (r *contactRepository) Search( 'lead', ( SELECT json_build_object( 'sender', (SELECT ea.email FROM email_accounts ea WHERE ea.id = hl.email_account_id), + -- Contacts copied on every email to this lead. + 'cc', %[5]s, 'hold', CASE WHEN %[4]s THEN json_build_object( 'since', hl.paused_at, 'until', hl.paused_until, 'reason', COALESCE(hl.pause_reason, ''), 'source', COALESCE(hl.pause_source, '') @@ -1689,7 +1699,7 @@ func (r *contactRepository) Search( ) FROM campaign_contact_progress p WHERE p.campaign_id = %[1]s AND p.contact_id = c.id - )`, singleCampaignPlaceholder, config.CampaignSendMaxAttempts, undeliverableClause(singleCampaignPlaceholder), liveHold("hl")) + )`, singleCampaignPlaceholder, config.CampaignSendMaxAttempts, undeliverableClause(singleCampaignPlaceholder), liveHold("hl"), leadCCJSONSQL(singleCampaignPlaceholder)) } // campaign_count is only ever read by the min/max filters and the @@ -1828,10 +1838,7 @@ func (r *contactRepository) Search( Step *string `json:"step"` // The lead row's own fields, read together because they come // from one campaign_leads row. - Lead *struct { - Sender *string `json:"sender"` - Hold *models.LeadHold `json:"hold"` - } `json:"lead"` + Lead *leadRowJSON `json:"lead"` Undeliverable bool `json:"undeliverable"` } @@ -1843,10 +1850,7 @@ func (r *contactRepository) Search( // which the outer query already excludes. lead := lp.Lead if lead == nil { - lead = &struct { - Sender *string `json:"sender"` - Hold *models.LeadHold `json:"hold"` - }{} + lead = &leadRowJSON{} } status := models.LeadStatusPending switch { @@ -1888,6 +1892,7 @@ func (r *contactRepository) Search( c.CampaignLead = &models.ContactCampaignProgress{ Status: status, Hold: lead.Hold, + CC: lead.CC, Sender: sender, Sent: lp.Sent, Opened: lp.Opened, diff --git a/internal/repository/pg_contact_campaign_state.go b/internal/repository/pg_contact_campaign_state.go index 20a2e5c32..1dfe0d91b 100644 --- a/internal/repository/pg_contact_campaign_state.go +++ b/internal/repository/pg_contact_campaign_state.go @@ -76,6 +76,12 @@ func (r *contactRepository) ListCampaignStates(ctx context.Context, orgID, conta } st.Steps = steps st.TotalSteps = len(steps) + cc, err := listLeadCC(ctx, r.DB.Pool, st.CampaignID, contactID) + if err != nil { + db.CaptureError(err, "", nil, "ListCampaignStates cc") + return nil, errx.InternalError() + } + st.CC = cc var sent, replied, bounced, failed bool var emailSteps, emailSent int diff --git a/internal/tasks/campaign_copies.go b/internal/tasks/campaign_copies.go new file mode 100644 index 000000000..f8a34f743 --- /dev/null +++ b/internal/tasks/campaign_copies.go @@ -0,0 +1,88 @@ +package tasks + +import ( + "context" + "strings" + + "github.com/google/uuid" + "github.com/warmbly/warmbly/internal/models" + "github.com/warmbly/warmbly/internal/pkg/mailhdr" +) + +// campaignCopies resolves who is copied on one send: the campaign's own CC and +// BCC, then the lead's copied contacts. An address that is suppressed, is the +// lead's own, or already appears earlier is left off, so a copy can never +// reach someone the lead's email would not have been allowed to. +func (s *tasksService) campaignCopies(ctx context.Context, orgID uuid.UUID, campaign *models.Campaign, contact *models.Contact) (cc, bcc []string, err error) { + seen := map[string]bool{strings.ToLower(strings.TrimSpace(contact.Email)): true} + keep := func(addr string) (bool, error) { + bare := strings.ToLower(mailhdr.Bare(addr)) + if bare == "" || seen[bare] { + return false, nil + } + if s.advanced != nil { + suppressed, _, xerr := s.advanced.ShouldSuppressRecipient(ctx, orgID, bare) + if xerr != nil { + return false, xerr + } + if suppressed { + return false, nil + } + } + seen[bare] = true + return true, nil + } + + // A campaign-wide copy that bounced on this campaign is dropped, so one bad + // address cannot keep failing every lead's send. + var wide []string + for _, a := range append(append([]string{}, campaign.CC...), campaign.BCC...) { + if bare := strings.ToLower(mailhdr.Bare(a)); bare != "" { + wide = append(wide, bare) + } + } + bounced, berr := s.campaignProgressRepo.BouncedCopyAddresses(ctx, campaign.ID, wide) + if berr != nil { + return nil, nil, berr + } + for a := range bounced { + seen[a] = true + } + + for _, a := range campaign.CC { + ok, kerr := keep(a) + if kerr != nil { + return nil, nil, kerr + } + if ok { + cc = append(cc, a) + } + } + for _, a := range campaign.BCC { + ok, kerr := keep(a) + if kerr != nil { + return nil, nil, kerr + } + if ok { + bcc = append(bcc, a) + } + } + + copies, lerr := s.campaignProgressRepo.ListLeadCC(ctx, campaign.ID, contact.ID) + if lerr != nil { + return nil, nil, lerr + } + for _, c := range copies { + // The status already applied suppression, bounces and verification. + if !c.Copied() { + continue + } + bare := strings.ToLower(strings.TrimSpace(c.Email)) + if bare == "" || seen[bare] { + continue + } + seen[bare] = true + cc = append(cc, c.Email) + } + return cc, bcc, nil +} diff --git a/internal/tasks/campaign_copies_test.go b/internal/tasks/campaign_copies_test.go new file mode 100644 index 000000000..69a6fed21 --- /dev/null +++ b/internal/tasks/campaign_copies_test.go @@ -0,0 +1,68 @@ +package tasks + +import ( + "context" + "reflect" + "strings" + "testing" + + "github.com/google/uuid" + "github.com/warmbly/warmbly/internal/app/advanced" + "github.com/warmbly/warmbly/internal/errx" + "github.com/warmbly/warmbly/internal/models" + "github.com/warmbly/warmbly/internal/repository" +) + +type copiesAdvanced struct { + advanced.Service + suppressed map[string]bool +} + +func (f copiesAdvanced) ShouldSuppressRecipient(_ context.Context, _ uuid.UUID, recipient string) (bool, string, *errx.Error) { + return f.suppressed[strings.ToLower(recipient)], "", nil +} + +type copiesProgress struct { + repository.CampaignProgressRepository + cc []models.CampaignLeadCC +} + +func (f copiesProgress) ListLeadCC(context.Context, uuid.UUID, uuid.UUID) ([]models.CampaignLeadCC, error) { + return f.cc, nil +} + +func (f copiesProgress) BouncedCopyAddresses(context.Context, uuid.UUID, []string) (map[string]bool, error) { + return map[string]bool{"refused@acme.test": true}, nil +} + +// A copy never reaches someone the lead's own email could not: suppressed +// campaign copies, the lead's own address, repeats and lead copies the status +// already refused are all left off. +func TestCampaignCopiesFiltersEveryCopy(t *testing.T) { + s := &tasksService{ + advanced: copiesAdvanced{suppressed: map[string]bool{"gone@acme.test": true}}, + campaignProgressRepo: copiesProgress{cc: []models.CampaignLeadCC{ + {Email: "jonas@acme.test", Status: models.LeadCCStatusActive}, + {Email: "bounced@acme.test", Status: models.LeadCCStatusBounced}, + {Email: "Boss@acme.test", Status: models.LeadCCStatusActive}, + {Email: "ana@acme.test", Status: models.LeadCCStatusActive}, + }}, + } + campaign := &models.Campaign{ + ID: uuid.New(), + CC: []string{"Boss ", "gone@acme.test", "ANA@acme.test"}, + BCC: []string{"crm@acme.test", "boss@acme.test", "Refused@acme.test"}, + } + contact := &models.Contact{ID: uuid.New(), Email: "ana@acme.test"} + + cc, bcc, err := s.campaignCopies(context.Background(), uuid.New(), campaign, contact) + if err != nil { + t.Fatalf("campaignCopies: %v", err) + } + if want := []string{"Boss ", "jonas@acme.test"}; !reflect.DeepEqual(cc, want) { + t.Fatalf("cc = %v, want %v", cc, want) + } + if want := []string{"crm@acme.test"}; !reflect.DeepEqual(bcc, want) { + t.Fatalf("bcc = %v, want %v", bcc, want) + } +} diff --git a/internal/tasks/campaign_task.go b/internal/tasks/campaign_task.go index a061af98c..3f63dc3df 100644 --- a/internal/tasks/campaign_task.go +++ b/internal/tasks/campaign_task.go @@ -540,6 +540,18 @@ func (s *tasksService) HandleCampaignTask(task *proto.ProcessTask) (result *errx taskRecord.EmailAccountID = account.ID } + // Who else the email copies, read for an email step before the send is + // reserved. Fail closed: copies that cannot be checked against suppression + // are not sent, and neither is the email without the copies chosen. + copyCC, copyBCC, cerr := s.campaignCopies(ctx, orgID, campaign, contact) + if cerr != nil { + errs.CaptureException(cerr) + s.taskRepo.RecordTaskFailure(ctx, taskID, "Could not read who the email copies", cerr.Error()) + s.retryCampaignTickLater(ctx, taskRecord) + executionStatus = "failed" + return errx.InternalError() + } + // STEP 9.4: The conversation this step joins. A follow-up is a nudge on the // email the contact already has, not a second cold email, so every step // after their first is threaded onto the last one they received: the @@ -815,8 +827,8 @@ func (s *tasksService) HandleCampaignTask(task *proto.ProcessTask) (result *errx emailMsg := EmailMessage{ From: account.Email, To: []string{contact.Email}, - CC: campaign.CC, - BCC: campaign.BCC, + CC: copyCC, + BCC: copyBCC, Subject: subject, BodyHTML: bodyHTML, BodyPlain: bodyPlain, diff --git a/skills/warmbly-api/SKILL.md b/skills/warmbly-api/SKILL.md index b55c0122b..e332010f7 100644 --- a/skills/warmbly-api/SKILL.md +++ b/skills/warmbly-api/SKILL.md @@ -36,7 +36,7 @@ Run `warmblyctl --help` for subcommands and `warmblyctl | Family | Covers | |---|---| | `me` | Identity and granted scopes | -| `campaign` | list, get, create, update, delete, steps, senders, preflight, start, stop, test-email, logs, plan, pause-lead / resume-lead | +| `campaign` | list, get, create, update, delete, steps, senders, preflight, start, stop, test-email, logs, plan, pause-lead / resume-lead, lead-cc / set-lead-cc / lead-cc-suggestions | | `contact` | list (search), get, lookup, create, update, delete, notes, timeline, import, imports, import-status, import-start, import-cancel, export | | `mailbox` | list, get, update, delete, auth-check, sync, skip-folders, identity, refresh-identity, behavior, verify, send, warmup-start/pause/resume/stop/status | | `inbox` | list, count, thread, seen, reply, compose, agent drafts, scheduled sends | @@ -99,6 +99,12 @@ These commands put real mail on the wire: `campaign start`, `campaign resume-lead` lifts it. An out-of-office auto-reply already writes the same hold by itself, across every campaign that contact is a lead of, so a lead reading `paused` for that reason needs nothing from you. +- To reach two people at one company in ONE thread, copy the second on the + first lead's emails: `campaign set-lead-cc --id --contact + --data '{"contact_ids":[""]}'` (at most two; + `campaign lead-cc-suggestions` lists likely colleagues). Do not enrol both as + leads of the same campaign for this: a copied contact's own lead is held + anyway, and resuming that hold sends them a second thread. - If deliverability analytics show rising bounces or complaints, stop the campaign first and report; do not push volume into a degrading mailbox. diff --git a/skills/warmbly-cli/SKILL.md b/skills/warmbly-cli/SKILL.md index 388600ca8..47c9661c7 100644 --- a/skills/warmbly-cli/SKILL.md +++ b/skills/warmbly-cli/SKILL.md @@ -72,7 +72,7 @@ gives the arguments and flags. Ids are positional, not flags. | Command | Covers | |---|---| | `status` | one call for "what is happening": mailboxes needing attention, what is sending, what is unread | -| `campaign` | list, view, create, edit, delete, steps, senders, segments, preflight, test, start, stop, logs, plan, pause-lead / resume-lead | +| `campaign` | list, view, create, edit, delete, steps, senders, segments, preflight, test, start, stop, logs, plan, pause-lead / resume-lead, lead-cc / set-lead-cc / lead-cc-suggestions | | `contact` | list, view, create, edit, delete, lookup, timeline, emails, notes, import, imports, import-status, import-start, import-cancel, export, verify | | `mailbox` | list, view, edit, check, sync, skip-folders, identity, refresh-identity, behavior, warmup, hold, release, send | | `inbox` | list, view, thread, read, reply, compose, drafts, scheduled, snooze | @@ -130,6 +130,13 @@ Everything else is safe to run freely. recipient answers with an out-of-office auto-reply, in every campaign that contact is a lead of, so do not also pause a lead that reads `paused` for that reason. +- To reach two people at one company in ONE thread, copy the second on the + first lead's emails: `warmbly campaign set-lead-cc CAMPAIGN_ID CONTACT_ID + --cc COLLEAGUE_ID` (at most two; `lead-cc-suggestions` lists likely + colleagues). Do not enrol both as leads of the same campaign for this: a + copied contact's own lead is held anyway, and `resume-lead` on that hold + sends them a second thread. Every copy is a recipient who did not ask for + the email, so keep it to small, personal campaigns. - If deliverability shows rising bounces or complaints, stop the campaign and report. Do not push volume into a degrading mailbox. - To check where copy lands before a launch, `warmbly placement test --mailbox diff --git a/web/src/components/app/campaigns/new/steps.tsx b/web/src/components/app/campaigns/new/steps.tsx index 313e6dc9c..ad2e78908 100644 --- a/web/src/components/app/campaigns/new/steps.tsx +++ b/web/src/components/app/campaigns/new/steps.tsx @@ -21,7 +21,7 @@ import type Sequence from "@/lib/api/models/app/campaigns/sequences/Sequence"; import type { DraftMeta } from "./serverDraft"; import EmailContentEditor from "@/components/app/campaigns/sequences/EmailContentEditor"; import { useSegments } from "@/lib/api/hooks/app/segments"; -import { CheckSquare } from "@/components/ui/check-square"; +import { Checkbox } from "@/components/ui/checkbox"; import TagSelector from "@/components/app/popup/select/TagSelector"; import ScrollStrip from "@/components/ui/scroll-strip"; import { DateTimePicker } from "@/components/ui/DateTimePicker"; @@ -148,20 +148,14 @@ export function LeadsStep({ {shown.map((l) => { const on = picked.has(l.id); return ( - + ); })} {!q && ( diff --git a/web/src/components/app/contacts/ContactsTable.tsx b/web/src/components/app/contacts/ContactsTable.tsx index cea0d5c7c..c818270bb 100644 --- a/web/src/components/app/contacts/ContactsTable.tsx +++ b/web/src/components/app/contacts/ContactsTable.tsx @@ -69,6 +69,7 @@ import type { CampaignLeadCounts } from "@/lib/api/models/app/contacts/SearchCon import ContactsEditBulk from "./ContactsEditBulk"; import PauseLeadDialog from "./PauseLeadDialog"; import { useResumeLead } from "@/lib/api/hooks/app/campaigns/useLeadHold"; +import { CC_RESUME_CONFIRM } from "@/lib/leadHold"; import { selectionOf } from "@/lib/api/models/app/contacts/ContactSelection"; import type ContactSelection from "@/lib/api/models/app/contacts/ContactSelection"; import * as rowSelection from "./selection"; @@ -549,22 +550,26 @@ export default function ContactsTable({ const [pauseTarget, setPauseTarget] = React.useState<{ id: string; name: string } | null>(null); const resumeLead = useResumeLead(); const resumeOne = React.useCallback( - async (contactId: string) => { + (contactId: string, copied?: boolean) => { if (!current_campaign) return; - try { - await toast.promise( - resumeLead.mutateAsync({ campaignId: current_campaign.id, contactId }), - { - loading: "Resuming lead…", - success: "Lead resumed", - error: (err: AppError) => buildError(err), - }, - ); - } catch { - /* toast.promise already surfaced it */ - } + const run = async () => { + try { + await toast.promise( + resumeLead.mutateAsync({ campaignId: current_campaign.id, contactId }), + { + loading: "Resuming lead…", + success: "Lead resumed", + error: (err: AppError) => buildError(err), + }, + ); + } catch { + /* toast.promise already surfaced it */ + } + }; + if (copied) confirm.show(CC_RESUME_CONFIRM, run); + else void run(); }, - [current_campaign, resumeLead], + [current_campaign, resumeLead, confirm], ); // Leads-view scope chips write straight into the search request, so the @@ -1185,7 +1190,7 @@ function ContactsTableBody({ // member without campaign write access, which takes the control off the // row rather than offering one that fails. onPauseLead?: (id: string, name: string) => void; - onResumeLead?: (id: string) => void; + onResumeLead?: (id: string, copied?: boolean) => void; emptyTitle: string; emptyBody: string; emptyCta: React.ReactNode; @@ -1405,7 +1410,7 @@ function ContactsTableBody({ type="button" aria-label="Resume lead" title={`${holdSummary(lead.hold)}. Resume now`} - onClick={() => onResumeLead(c.id)} + onClick={() => onResumeLead(c.id, lead.hold?.source === "cc")} className="size-6 rounded text-violet-500 hover:text-violet-700 hover:bg-violet-50 flex items-center justify-center transition-colors" > diff --git a/web/src/components/app/contacts/LeadHoldButtons.tsx b/web/src/components/app/contacts/LeadHoldButtons.tsx index 50602601f..96706df95 100644 --- a/web/src/components/app/contacts/LeadHoldButtons.tsx +++ b/web/src/components/app/contacts/LeadHoldButtons.tsx @@ -5,6 +5,7 @@ import { Loader2Icon, PauseIcon, PlayIcon } from "lucide-react"; import toast from "react-hot-toast"; import PauseLeadDialog from "./PauseLeadDialog"; import { useResumeLead } from "@/lib/api/hooks/app/campaigns/useLeadHold"; +import { useConfirm } from "@/hooks/context/confirm"; import type { AppError } from "@/lib/api/client/normalizeError"; import buildError from "@/lib/helper/buildError"; @@ -40,14 +41,18 @@ export function ResumeLeadButton({ label = "Resume", disabled = false, onBusyChange, + confirmText, }: { campaignId: string; contactId: string; label?: string; disabled?: boolean; onBusyChange?: (busy: boolean) => void; + // Asked first when resuming has a consequence worth a second look. + confirmText?: string; }) { const resume = useResumeLead(); + const confirm = useConfirm(); async function run() { onBusyChange?.(true); try { @@ -65,7 +70,7 @@ export function ResumeLeadButton({ return ( + )} + + ); + })} + {canAdd && ( + { + setOpen(false); + void save([...ids, id], "Copied on every email to this lead"); + }} + /> + )} + + ); +} + +function CCPicker({ + open, + setOpen, + campaignId, + contactId, + contactName, + exclude, + busy, + empty, + onPick, +}: { + open: boolean; + setOpen: (v: boolean) => void; + campaignId: string; + contactId: string; + contactName: string; + exclude: string[]; + busy: boolean; + empty: boolean; + onPick: (id: string) => void; +}) { + const ref = React.useRef(null); + const triggerRef = React.useRef(null); + const [query, setQuery] = React.useState(""); + useClickOutside(ref, () => setOpen(false)); + const placement = useFlipPlacement(triggerRef, open, 300); + + const q = useDebouncedValue(query.trim(), 250); + const suggestions = useLeadCCSuggestions(campaignId, contactId, open); + const search = useSearchContacts({ + options: { + query: q, + custom_field_filters: [], + campaign_ids: [], + sort_by: "updated_at", + reverse: false, + }, + limit: 8, + enabled: open && q.length > 0, + keepPrevious: true, + }); + + const skip = React.useMemo(() => new Set([contactId, ...exclude]), [contactId, exclude]); + const candidates: Candidate[] = React.useMemo(() => { + if (q.length > 0) { + return (search.contacts ?? []) + .filter((c) => !skip.has(c.id)) + .map((c) => ({ + id: c.id, + email: c.email, + name: `${c.first_name ?? ""} ${c.last_name ?? ""}`.trim() || c.email, + company: c.company, + })); + } + return (suggestions.data?.data ?? []) + .filter((s) => !skip.has(s.contact_id)) + .map((s) => ({ + id: s.contact_id, + email: s.email, + name: leadCCName(s), + company: s.company, + tag: s.reason === "company" ? "Same company" : "Same domain", + })); + }, [q, search.contacts, suggestions.data, skip]); + + const loading = q.length > 0 ? search.isFetching && !search.contacts : suggestions.isLoading; + + return ( +
{ + if (e.key === "Escape" && open) { + e.stopPropagation(); + setOpen(false); + } + }} + > + + + {open && ( + +
+ + setQuery(e.target.value)} + placeholder="Search contacts…" + autoFocus + className="w-full h-5 bg-transparent text-[12px] text-slate-900 placeholder:text-slate-400 outline-none" + /> +
+
+ {loading ? ( +
+ +
+ ) : candidates.length === 0 ? ( +
+ {q.length > 0 ? "No matching contacts." : "No colleagues found. Search for a contact."} +
+ ) : ( + candidates.map((c) => ( + + )) + )} +
+
+ Copied on every email to {contactName} in this campaign, up to {LEAD_CC_MAX}. A reply from + anyone counts as {contactName}'s reply. +
+
+ )} +
+
+ ); +} diff --git a/web/src/components/app/placement/tests/SeedChooser.tsx b/web/src/components/app/placement/tests/SeedChooser.tsx index b8c63c617..194388d84 100644 --- a/web/src/components/app/placement/tests/SeedChooser.tsx +++ b/web/src/components/app/placement/tests/SeedChooser.tsx @@ -103,7 +103,6 @@ export default function SeedChooser({ )} > toggle(s.email_account_id)} diff --git a/web/src/components/app/unibox/ContactContextPanel.tsx b/web/src/components/app/unibox/ContactContextPanel.tsx index 073a046b5..50ddc4166 100644 --- a/web/src/components/app/unibox/ContactContextPanel.tsx +++ b/web/src/components/app/unibox/ContactContextPanel.tsx @@ -46,7 +46,7 @@ import useContactCampaignStates from "@/lib/api/hooks/app/contacts/useContactCam import type ContactCampaignState from "@/lib/api/models/app/contacts/ContactCampaignState"; import type MiniCampaign from "@/lib/api/models/app/campaigns/MiniCampaign"; import { holdSummary } from "@/lib/api/models/app/contacts/Contact"; -import { leadCanBePaused } from "@/lib/leadHold"; +import { CC_RESUME_CONFIRM, leadCanBePaused } from "@/lib/leadHold"; import { usePermission } from "@/hooks/usePermission"; import { PauseLeadButton, ResumeLeadButton } from "@/components/app/contacts/LeadHoldButtons"; import LeadStatusPill from "@/components/app/contacts/LeadStatusPill"; @@ -372,7 +372,11 @@ function CampaignsSection({ {line.text} {canWrite && s.hold ? ( - + ) : canWrite && leadCanBePaused(s) ? (