diff --git a/docs/content/docs/api/error-codes.mdx b/docs/content/docs/api/error-codes.mdx index 1629ab9f1..15f38443b 100644 --- a/docs/content/docs/api/error-codes.mdx +++ b/docs/content/docs/api/error-codes.mdx @@ -235,6 +235,12 @@ Returned when the request conflicts with existing data. } ``` +One conflict carries its own `code`. A contact's email address has to be free, so changing one to an address another contact already holds is refused rather than merging the two: + +| `code` | Status | Meaning | +|--------|--------|---------| +| `contact_email_taken` | 409 | The address given to [update a contact](/api/reference/contacts/#update-a-contact) already belongs to another contact | + ### 422 Unprocessable Returned when validation fails on the request data. diff --git a/docs/content/docs/api/reference/contacts.mdx b/docs/content/docs/api/reference/contacts.mdx index 269f4628c..6b7cd07fa 100644 --- a/docs/content/docs/api/reference/contacts.mdx +++ b/docs/content/docs/api/reference/contacts.mdx @@ -146,7 +146,7 @@ A JSON array of contact objects (at least one, up to the per-request maximum; an | Field | Type | Required | Description | | --- | --- | --- | --- | -| `email` | string | Yes | Contact email address. | +| `email` | string | Yes | Contact email address. Stored lowercased; a display name (`Dana Reyes `) is reduced to the address inside it, and anything that is not an address answers `400`. | | `first_name` | string | No | First name. | | `last_name` | string | No | Last name. | | `company` | string | No | Company name. | @@ -585,6 +585,8 @@ When the contact is suppressed, `suppression` is an object: `{ "id", "kind", "va Partially updates a single contact. Only the fields present are changed. Campaign and category lists can be set wholesale or adjusted with diff-style add/remove. +`email` replaces the contact's address. It is stored lowercased, a display name (`Dana Reyes `) is reduced to the address inside it, and anything that is not an address answers `400`. The address has to be free: one another contact already holds answers `409` with `code` `contact_email_taken` rather than merging the two. A changed address drops the contact's verification verdict back to `unknown`, clears the delivery evidence behind it and forgets the cached `esp_provider`, because all three belonged to the old mailbox; the next verification pass checks the new address and the recipient provider is derived again from the new domain. Steps already sent went to the old address and keep their history. + Auth: **Scope** `WRITE_CONTACTS` · **Org permission** `manage_contacts` | Parameter | In | Type | Description | @@ -597,6 +599,7 @@ Auth: **Scope** `WRITE_CONTACTS` · **Org permission** `manage_contacts` | --- | --- | --- | --- | | `first_name` | string | No | First name. | | `last_name` | string | No | Last name. | +| `email` | string | No | New email address. Normalized to lowercase; must be a valid address and unused by another contact. | | `company` | string | No | Company. | | `phone` | string | No | Phone. | | `custom_fields` | object | No | Replaces the custom-fields map. | diff --git a/docs/content/docs/guides/contacts-crm.mdx b/docs/content/docs/guides/contacts-crm.mdx index 96103a4f9..2c304e386 100644 --- a/docs/content/docs/guides/contacts-crm.mdx +++ b/docs/content/docs/guides/contacts-crm.mdx @@ -27,7 +27,7 @@ You must map at least one column to **Email**. Skipping still adds the contact to the campaigns, segments, and categories the import targets: it leaves their fields alone, it does not leave them out of the list. If the same address appears twice in one file it becomes one contact, and the extra rows count as skipped. -Rows with a missing or invalid email are reported with a line number and reason rather than imported. A blank cell never erases a value you already have, so re-importing a partial export enriches contacts instead of wiping them. +Rows with a missing or invalid email are reported with a line number and reason rather than imported. Addresses are stored lowercased and stripped of any name around them, so a cell reading `Dana Reyes ` imports as `dana@acme.com` and both spellings of one address land on one contact. A blank cell never erases a value you already have, so re-importing a partial export enriches contacts instead of wiping them. Problems with the mapping itself, an unnamed custom field or a name Warmbly cannot use, are reported once before anything is written, so a single typo can never fail every row. @@ -87,6 +87,12 @@ Every bulk action reads the selection: **Edit**, **Segment**, **Remove from segm **Select all matching** is in the **From contacts** picker too, on a campaign's Leads tab and a segment page, so a whole search can be added as leads or members in one step. +## Editing one contact + +Clicking a row opens the contact's panel. Its **Details** tab edits the fields the contact is made of: name, email address, company, phone, subscription, campaigns, categories and custom fields. Nothing is sent until **Save**, **Discard** puts the panel back to the stored values, and closing it with unsaved edits asks first. + +The email address can be changed here, which is the right move when someone's address was mistyped on import or they moved to a new domain. It is stored lowercased and stripped of any display name, and it has to be free: an address another contact already holds is refused, because merging two people's campaign history is not something the edit could undo. A new address also clears the contact's [verification](/guides/deliverability/#address-verification) verdict and everything the platform had observed about the old mailbox, so the background check starts the new address from scratch on its next pass, and the recipient provider Warmbly matches senders against is worked out again from the new domain. Emails already sent went to the old address and stay in the timeline as they happened. + ## Editing many contacts at once Tick rows in any contact list, a segment's members or a campaign's leads, and the selection bar's **Edit** opens a bulk panel. It can add or remove campaigns, add or remove categories, force a subscription state, and queue custom-field operations (add, edit, delete, rename) that run on every selected contact, including a **Select all matching** selection that reaches past the loaded pages. diff --git a/docs/public/openapi.json b/docs/public/openapi.json index 6db8ee110..855df6f8e 100644 --- a/docs/public/openapi.json +++ b/docs/public/openapi.json @@ -6976,7 +6976,7 @@ ], "operationId": "contacts_update", "summary": "Update a contact", - "description": "Partially updates a single contact; only the fields present change. Category lists can be set wholesale or adjusted with diff-style add/remove. Scope `WRITE_CONTACTS`.", + "description": "Partially updates a single contact; only the fields present change. Category lists can be set wholesale or adjusted with diff-style add/remove. `email` replaces the contact's address and resets its verification verdict. Scope `WRITE_CONTACTS`.", "security": [ { "bearerAuth": [] @@ -7057,6 +7057,16 @@ } } }, + "409": { + "description": "Another contact already uses the supplied email address (code contact_email_taken).", + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/Error" + } + } + } + }, "429": { "description": "Rate limited.", "content": { @@ -23931,6 +23941,10 @@ "last_name": { "type": "string" }, + "email": { + "type": "string", + "description": "New email address. Stored lowercased and stripped of any display name. Must be free: an address another contact holds returns 409 contact_email_taken. Changing it resets the contact's verification verdict to unknown." + }, "company": { "type": "string" }, diff --git a/internal/app/advanced/service.go b/internal/app/advanced/service.go index a64dfda9a..3cdd7eb72 100644 --- a/internal/app/advanced/service.go +++ b/internal/app/advanced/service.go @@ -1163,7 +1163,7 @@ func (s *service) ProcessIncomingReply(ctx context.Context, emailAccountID uuid. if replyclassify.IsAutomated(replyResult.Class) { kind = "auto_replied" } - s.evidence.RecordEvidence(ctx, ctID, kind, msg.ID.String(), "") + s.evidence.RecordEvidence(ctx, ctID, models.Step(&cID, &sID), kind, msg.ID.String(), "") } if !replyclassify.IsAutomated(replyResult.Class) { _ = s.campaignProgressRepo.RecordEmailReplied(ctx, cID, ctID, sID) @@ -1510,7 +1510,10 @@ func (s *service) IngestDeliverabilityEvent(ctx context.Context, organizationID if emailverify.NamesRecipient(req.Reason) { kind = "bounced_recipient" } - s.evidence.RecordEvidence(ctx, *req.ContactID, kind, req.IdempotencyKey, req.Reason) + // Both halves of the step come from the resolved task, so the + // pair names one real row rather than a request-supplied + // campaign paired with a resolved sequence. + s.evidence.RecordEvidence(ctx, *req.ContactID, models.Step(campaignTask.CampaignID, campaignTask.SequenceID), kind, req.IdempotencyKey, req.Reason) } case models.DeliverabilityEventComplaint: _ = s.campaignProgressRepo.RecordEmailComplained(ctx, *req.CampaignID, *req.ContactID, *campaignTask.SequenceID) @@ -2381,7 +2384,7 @@ func (s *service) listQualityCheck(ctx context.Context, orgID, campaignID uuid.U // WireAudience attaches the launch-time list measurement. // EvidenceRecorder mirrors emailverify.EvidenceRecorder without importing it. type EvidenceRecorder interface { - RecordEvidence(ctx context.Context, contactID uuid.UUID, kind, ref, detail string) + RecordEvidence(ctx context.Context, contactID uuid.UUID, step models.EvidenceStep, kind, ref, detail string) } // EvidenceAware lets main hand the service the verification evidence ledger. diff --git a/internal/app/consumer/event_send_result.go b/internal/app/consumer/event_send_result.go index 10fde9970..056920daf 100644 --- a/internal/app/consumer/event_send_result.go +++ b/internal/app/consumer/event_send_result.go @@ -211,7 +211,7 @@ func (s *JobsService) failCampaignSend(ctx context.Context, task *repository.Tas // delivery route; a rejection of the sender, the session or the content // says nothing about the address. if s.Evidence != nil && ct.ContactID != nil && ct.SequenceID != nil && emailverify.NamesRecipient(reason) { - s.Evidence.RecordEvidence(ctx, *ct.ContactID, "bounced_recipient", "send:"+ct.SequenceID.String(), reason) + s.Evidence.RecordEvidence(ctx, *ct.ContactID, models.Step(&campaignID, ct.SequenceID), "bounced_recipient", "send:"+ct.SequenceID.String(), reason) } attempts, exhausted, rolledBack := 0, false, false diff --git a/internal/app/consumer/event_tracking.go b/internal/app/consumer/event_tracking.go index 561161359..d5d277f8b 100644 --- a/internal/app/consumer/event_tracking.go +++ b/internal/app/consumer/event_tracking.go @@ -359,7 +359,7 @@ func (tc *TrackingConsumer) HandleTrackingEvent(ctx context.Context, event *even // A human open proves the mailbox is live; a prefetch proves // only that a proxy fetched an image. if tc.evidence != nil { - tc.evidence.RecordEvidence(ctx, contactID, "opened", sequenceID.String(), "") + tc.evidence.RecordEvidence(ctx, contactID, models.Step(&campaignID, &sequenceID), "opened", sequenceID.String(), "") } } case events.EventTypeEmailClicked: @@ -379,7 +379,7 @@ func (tc *TrackingConsumer) HandleTrackingEvent(ctx context.Context, event *even } else if err == nil { instantKind = "click" if tc.evidence != nil { - tc.evidence.RecordEvidence(ctx, contactID, "clicked", sequenceID.String(), "") + tc.evidence.RecordEvidence(ctx, contactID, models.Step(&campaignID, &sequenceID), "clicked", sequenceID.String(), "") } } } @@ -461,7 +461,7 @@ func (tc *TrackingConsumer) finishHumanClick(task *repository.CampaignTask, even tc.publishTrackingEvent(ctx, task, event, true, label, origin) } else { if tc.evidence != nil { - tc.evidence.RecordEvidence(ctx, *task.ContactID, "clicked", task.SequenceID.String(), "") + tc.evidence.RecordEvidence(ctx, *task.ContactID, models.Step(task.CampaignID, task.SequenceID), "clicked", task.SequenceID.String(), "") } if tc.advancedService != nil { tc.advancedService.FireInstantActions(ctx, *task.CampaignID, *task.ContactID, *task.SequenceID, "click") diff --git a/internal/app/contact/import.go b/internal/app/contact/import.go index a310c3afb..bf17f4416 100644 --- a/internal/app/contact/import.go +++ b/internal/app/contact/import.go @@ -279,13 +279,17 @@ func (s *contactService) ImportCommit( parsed = append(parsed, p) continue } - contact.Email = strings.TrimSpace(contact.Email) - if contact.Email == "" || !email.IsValid(contact.Email) { + // Normalized rather than lowercased: a cell holding + // `Dana Reyes ` parses as an address and used to be + // imported whole as the recipient. The dedupe below keys on the result, + // so the two spellings of one address also collapse into one contact. + addr, ok := email.Normalize(contact.Email) + if !ok { p.errMsg = "missing or invalid email" parsed = append(parsed, p) continue } - contact.Email = strings.ToLower(contact.Email) + contact.Email = addr if prev, dup := firstByEmail[contact.Email]; dup { // Same address twice in one file. "skip" keeps the first row; diff --git a/internal/app/emailverify/evidence.go b/internal/app/emailverify/evidence.go index e509506a6..456c388ca 100644 --- a/internal/app/emailverify/evidence.go +++ b/internal/app/emailverify/evidence.go @@ -18,7 +18,7 @@ import ( // (contact, kind, ref) and best effort: a failure is logged, never returned // into the caller's path. type EvidenceRecorder interface { - RecordEvidence(ctx context.Context, contactID uuid.UUID, kind, ref, detail string) + RecordEvidence(ctx context.Context, contactID uuid.UUID, step models.EvidenceStep, kind, ref, detail string) } // Evidence scores contacts from the ledger. It is separate from Service so @@ -37,11 +37,11 @@ func NewEvidence(repo repository.VerificationEvidenceRepository) *Evidence { func (e *Evidence) SetOnChange(fn func(ctx context.Context, contactID uuid.UUID)) { e.onChange = fn } -func (e *Evidence) RecordEvidence(ctx context.Context, contactID uuid.UUID, kind, ref, detail string) { +func (e *Evidence) RecordEvidence(ctx context.Context, contactID uuid.UUID, step models.EvidenceStep, kind, ref, detail string) { if e == nil || e.repo == nil || contactID == uuid.Nil { return } - inserted, err := e.repo.Record(ctx, contactID, kind, ref, detail, time.Now().UTC()) + inserted, err := e.repo.Record(ctx, contactID, step, kind, ref, detail, time.Now().UTC()) if err != nil { log.Warn().Err(err).Str("contact_id", contactID.String()).Str("kind", kind).Msg("verification evidence not recorded") return diff --git a/internal/email/validation.go b/internal/email/validation.go index 3baf33a79..fe0c37930 100644 --- a/internal/email/validation.go +++ b/internal/email/validation.go @@ -1,8 +1,23 @@ package email -import "net/mail" +import ( + "net/mail" + "strings" +) func IsValid(email string) bool { _, err := mail.ParseAddress(email) return err == nil } + +// Normalize returns an address in the form contacts are stored in: trimmed, +// lowercased, and stripped of any display name, because mail.ParseAddress +// accepts `Dana Reyes ` and storing that whole string as the +// address would send to nobody. ok is false for anything it refuses. +func Normalize(addr string) (string, bool) { + a, err := mail.ParseAddress(strings.TrimSpace(addr)) + if err != nil { + return "", false + } + return strings.ToLower(a.Address), true +} diff --git a/internal/errx/common.go b/internal/errx/common.go index 8e3387dd5..33fb2d39d 100644 --- a/internal/errx/common.go +++ b/internal/errx/common.go @@ -195,6 +195,11 @@ var ( // Contact ErrContactSerialize = New(BadRequest, "Failed to serialize contact.") ErrContactSize = New(BadRequest, "Contact size cannot be bigger than 10KB.") + // A contact's address is unique within the workspace, so an edit that + // collides with another contact is refused rather than merged: merging two + // people's campaign progress is not something an edit can undo. + ErrContactEmailTaken = NewWithIdentifier(Conflict, "contact_email_taken", + "Another contact already uses this email address.") // Unibox ErrUniboxLimit = New(BadRequest, fmt.Sprintf("Limit must be between %d and %d.", config.UniboxLimitMin, config.UniboxLimitMax)) diff --git a/internal/infrastructure/db/migrations/000164_contact_evidence_reset.down.sql b/internal/infrastructure/db/migrations/000164_contact_evidence_reset.down.sql new file mode 100644 index 000000000..9b1085bbc --- /dev/null +++ b/internal/infrastructure/db/migrations/000164_contact_evidence_reset.down.sql @@ -0,0 +1,2 @@ +ALTER TABLE contacts + DROP COLUMN IF EXISTS verification_evidence_reset_at; diff --git a/internal/infrastructure/db/migrations/000164_contact_evidence_reset.up.sql b/internal/infrastructure/db/migrations/000164_contact_evidence_reset.up.sql new file mode 100644 index 000000000..0ff0a52e1 --- /dev/null +++ b/internal/infrastructure/db/migrations/000164_contact_evidence_reset.up.sql @@ -0,0 +1,14 @@ +-- When a contact's address last changed, so the delivery-credit job can tell +-- mail sent to the old mailbox from mail sent to this one. +-- +-- A contact's verdict and its evidence ledger describe an ADDRESS, so editing +-- the address wipes both. That wipe did not hold on its own: the credit job +-- re-derives a 'delivered' row from every step ever sent, with no lower bound +-- of its own, and its dedupe is exactly the set of rows the wipe deleted. The +-- next pass handed them straight back, and the new address inherited a verdict +-- earned by the old one, decisive enough to excuse it from ever being probed +-- (issue #511). +-- +-- NULL means the address has never changed, which is every contact today. +ALTER TABLE contacts + ADD COLUMN verification_evidence_reset_at timestamptz; diff --git a/internal/models/contact.go b/internal/models/contact.go index 0a653047c..db63781fd 100644 --- a/internal/models/contact.go +++ b/internal/models/contact.go @@ -659,9 +659,28 @@ type ContactTimelineResult struct { Pagination Pagination `json:"pagination"` } +// EvidenceStep names the campaign step a verification observation came from. +// It is what lets a late bounce or open be refused when the mail it describes +// left before the contact's address was edited: the observation is about the +// old mailbox, not the one the contact holds now. A zero value means the +// observation is not attributable to a step and is always recorded. +type EvidenceStep struct { + CampaignID *uuid.UUID + SequenceID *uuid.UUID +} + +// Step builds an EvidenceStep from ids a caller already holds. +func Step(campaignID, sequenceID *uuid.UUID) EvidenceStep { + return EvidenceStep{CampaignID: campaignID, SequenceID: sequenceID} +} + type UpdateContact struct { - FirstName *string `json:"first_name"` - LastName *string `json:"last_name"` + FirstName *string `json:"first_name"` + LastName *string `json:"last_name"` + // Email replaces the contact's address. It is the contact's identity, so + // changing it drops the verification verdict and the delivery evidence + // that belonged to the old mailbox. + Email *string `json:"email"` Company *string `json:"company"` Phone *string `json:"phone"` CustomFields *map[string]string `json:"custom_fields"` diff --git a/internal/repository/contact_campaign_live_test.go b/internal/repository/contact_campaign_live_test.go index 8aeeaade0..63bac19c7 100644 --- a/internal/repository/contact_campaign_live_test.go +++ b/internal/repository/contact_campaign_live_test.go @@ -50,7 +50,7 @@ func liveContactDB(t *testing.T) (*db.DB, *pgxpool.Pool) { if err := handle.Ping(context.Background()); err != nil { t.Fatalf("ping: %v", err) } - requireSchemaVersion(t, handle.Pool, 156) + requireSchemaVersion(t, handle.Pool, 164) return handle, handle.Pool } diff --git a/internal/repository/contact_email_update_live_test.go b/internal/repository/contact_email_update_live_test.go new file mode 100644 index 000000000..676d43398 --- /dev/null +++ b/internal/repository/contact_email_update_live_test.go @@ -0,0 +1,308 @@ +package repository + +import ( + "context" + "strings" + "testing" + "time" + + "github.com/google/uuid" + + "github.com/warmbly/warmbly/internal/errx" + "github.com/warmbly/warmbly/internal/models" +) + +// Regression cover for issue #511: editing a contact's email address saved +// nothing. models.UpdateContact carried no Email field, so the dashboard's +// `{"email": "..."}` was dropped by the JSON decode and the PATCH answered 200 +// with the old address. +// +// Run against the dev stack: +// +// WARMBLY_TEST_DB=postgres://warmbly:warmbly@localhost:15432/warmbly_dev?sslmode=disable \ +// go test ./internal/repository/ -run LiveContactEmail -v + +func TestLiveContactEmailIsSavedAndResetsVerification(t *testing.T) { + handle, pool := liveContactDB(t) + f := newSharedOrgFixture(t, pool) + repo := NewContactRepostory(handle) + ctx := context.Background() + mate := f.mate.String() + + // The contact arrives with a verdict and an observation, both of which + // belong to the address rather than to the person. + if _, err := pool.Exec(ctx, ` + UPDATE contacts SET verification_status = 'valid', verification_sub_status = 'role', + verification_reason = 'accepted', verification_source = 'probe', + verification_provider = 'builtin', is_catch_all = true, + verification_checked_at = NOW(), verification_confidence = 80, + verification_evidence_at = NOW(), esp_provider = 'gmail', esp_resolved_at = NOW() + WHERE id = $1`, f.contact); err != nil { + t.Fatalf("seed verdict: %v", err) + } + if _, err := pool.Exec(ctx, ` + INSERT INTO contact_verification_evidence (contact_id, kind, ref) VALUES ($1, 'delivered', 'i511')`, + f.contact); err != nil { + t.Fatalf("seed evidence: %v", err) + } + + type state struct { + addr, status, source, esp string + catchAll, checked bool + confidence int16 + evidenceAt, resetAt bool + evidence int + } + read := func() state { + t.Helper() + var st state + if err := pool.QueryRow(ctx, ` + SELECT c.email, c.verification_status, c.verification_source, c.esp_provider, c.is_catch_all, + c.verification_checked_at IS NOT NULL, c.verification_confidence, + c.verification_evidence_at IS NOT NULL, c.verification_evidence_reset_at IS NOT NULL, + (SELECT COUNT(*) FROM contact_verification_evidence e WHERE e.contact_id = c.id) + FROM contacts c WHERE c.id = $1`, f.contact). + Scan(&st.addr, &st.status, &st.source, &st.esp, &st.catchAll, &st.checked, &st.confidence, + &st.evidenceAt, &st.resetAt, &st.evidence); err != nil { + t.Fatalf("read contact: %v", err) + } + return st + } + + // A re-save of the address already on the row, in another case, is not a + // change: it must not throw away a verdict the address earned. + same := "I187-" + f.contact.String()[:8] + "@Test.Local" + if _, xerr := repo.Update(ctx, mate, f.contact.String(), f.org, &models.UpdateContact{Email: &same}); xerr != nil { + t.Fatalf("update same address: %v", xerr) + } + if st := read(); st.status != "valid" || st.evidence != 1 || st.resetAt { + t.Fatalf("re-saving the same address reset verification: %+v", st) + } + + // A stored address that predates normalization is rewritten in place, and + // that is still not a different mailbox. + if _, err := pool.Exec(ctx, `UPDATE contacts SET email = $2 WHERE id = $1`, f.contact, same); err != nil { + t.Fatalf("legacy casing: %v", err) + } + lower := strings.ToLower(same) + if _, xerr := repo.Update(ctx, mate, f.contact.String(), f.org, &models.UpdateContact{Email: &lower}); xerr != nil { + t.Fatalf("normalize casing: %v", xerr) + } + if st := read(); st.addr != lower || st.status != "valid" || st.evidence != 1 { + t.Fatalf("case-only save did not normalize in place: %+v", st) + } + + // The real edit. A display name and stray case are normalized away. + next := " Dana Reyes " + updated, xerr := repo.Update(ctx, mate, f.contact.String(), f.org, &models.UpdateContact{Email: &next}) + if xerr != nil { + t.Fatalf("update email: %v", xerr) + } + if updated.Email != "dana@acme.test" { + t.Fatalf("response email = %q, want dana@acme.test", updated.Email) + } + st := read() + if st.addr != "dana@acme.test" { + t.Fatalf("stored email = %q, want dana@acme.test", st.addr) + } + if st.status != "unknown" || st.source != "" || st.catchAll || st.checked || st.confidence != 0 || st.evidenceAt { + t.Fatalf("verification survived the address change: %+v", st) + } + if st.evidence != 0 { + t.Fatalf("evidence rows after the address change = %d, want 0", st.evidence) + } + if !st.resetAt { + t.Fatal("no evidence watermark, so the delivery credit job will hand the new address the old mailbox's record") + } + // esp_provider is a cache of the address domain and the scheduler only + // fills it when empty, so a stale one routes ESP-matched sends forever. + if st.esp != "" { + t.Fatalf("esp_provider after the address change = %q, want empty", st.esp) + } + + // An unrelated edit still reports the current address. + company := "Acme" + updated, xerr = repo.Update(ctx, mate, f.contact.String(), f.org, &models.UpdateContact{Company: &company}) + if xerr != nil { + t.Fatalf("update company: %v", xerr) + } + if updated.Email != "dana@acme.test" { + t.Fatalf("email after unrelated edit = %q", updated.Email) + } +} + +// The delivery-credit job re-derives 'delivered' evidence from every step ever +// sent. Without the watermark the rows the address change deletes come back on +// its next pass, and the new address inherits a verdict earned by the old one. +func TestLiveContactEmailChangeSurvivesTheDeliveryCreditJob(t *testing.T) { + handle, pool := liveContactDB(t) + f := newSharedOrgFixture(t, pool) + repo := NewContactRepostory(handle) + evidence := NewVerificationEvidenceRepository(handle) + ctx := context.Background() + + seq := uuid.New() + if _, err := pool.Exec(ctx, `INSERT INTO sequences (id, campaign_id, organization_id, name, subject, body_plain, body_html) + VALUES ($1, $2, $3, 'Email 1', 'Hi', 'Body', 'Body')`, seq, f.campaign, f.org); err != nil { + t.Fatalf("sequence: %v", err) + } + 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() - interval '30 days', NOW() - interval '30 days')`, f.campaign, f.contact, seq); err != nil { + t.Fatalf("progress: %v", err) + } + t.Cleanup(func() { + c := context.Background() + for _, sql := range []string{ + `DELETE FROM contact_verification_evidence WHERE contact_id IN (SELECT id FROM contacts WHERE organization_id = $1)`, + `DELETE FROM campaign_contact_progress WHERE campaign_id IN (SELECT id FROM campaigns WHERE organization_id = $1)`, + `DELETE FROM sequences WHERE organization_id = $1`, + } { + if _, err := pool.Exec(c, sql, f.org); err != nil { + t.Errorf("cleanup: %v", err) + } + } + }) + + // The old address earns its delivery. The credit job works a global + // backlog under one limit, so on a shared database this contact can sit + // behind other people's steps; drain until it comes back. + credited := false + for pass := 0; pass < 20 && !credited; pass++ { + ids, err := evidence.CreditCleanDeliveries(ctx, time.Hour, 500) + if err != nil { + t.Fatalf("credit: %v", err) + } + if len(ids) == 0 { + break + } + for _, id := range ids { + if id == f.contact { + credited = true + } + } + } + if !credited { + t.Fatal("the delivery to the old address was never credited, so the test proves nothing") + } + var rows int + count := func() int { + t.Helper() + if err := pool.QueryRow(ctx, `SELECT COUNT(*) FROM contact_verification_evidence WHERE contact_id = $1`, f.contact).Scan(&rows); err != nil { + t.Fatalf("count evidence: %v", err) + } + return rows + } + if count() != 1 { + t.Fatalf("evidence after the credit pass = %d, want 1", rows) + } + + next := "moved@acme.test" + if _, xerr := repo.Update(ctx, f.mate.String(), f.contact.String(), f.org, &models.UpdateContact{Email: &next}); xerr != nil { + t.Fatalf("update email: %v", xerr) + } + if count() != 0 { + t.Fatalf("evidence right after the address change = %d, want 0", rows) + } + + // The next pass must not hand it back. + if _, err := evidence.CreditCleanDeliveries(ctx, time.Hour, 1000); err != nil { + t.Fatalf("credit again: %v", err) + } + if count() != 0 { + t.Fatalf("the credit job re-derived %d evidence rows for the new address from mail sent to the old one", rows) + } + + // A send that was already on the bus when the address changed has its + // sent_at stamped by the worker's result afterwards, so reading sent_at + // alone would let a delivery to the OLD mailbox through the watermark. + // dispatched_at is when it left, and that is what the credit reads. + if _, err := pool.Exec(ctx, `UPDATE campaign_contact_progress SET sent_at = NOW() WHERE campaign_id = $1 AND contact_id = $2`, + f.campaign, f.contact); err != nil { + t.Fatalf("late stamp: %v", err) + } + if _, err := evidence.CreditCleanDeliveries(ctx, 0, 1000); err != nil { + t.Fatalf("credit in-flight: %v", err) + } + if count() != 0 { + t.Fatalf("a send dispatched before the address change was credited to the new address (%d rows)", rows) + } + + // Nor may a late event about that step: a hard bounce for the typo lands + // after the correction and would otherwise mark the new address invalid + // and stop every send to it. + step := models.Step(&f.campaign, &seq) + if inserted, err := evidence.Record(ctx, f.contact, step, "bounced_recipient", "late-bounce", "550 no such user", time.Now()); err != nil || inserted { + t.Fatalf("a bounce for the old address was recorded against the new one (inserted=%v, err=%v)", inserted, err) + } + if count() != 0 { + t.Fatalf("evidence after the late bounce = %d, want 0", rows) + } + + // A step dispatched after the correction is about this address, and an + // observation that names no step at all is always kept. + if _, err := pool.Exec(ctx, `UPDATE campaign_contact_progress SET dispatched_at = NOW(), sent_at = NOW() WHERE campaign_id = $1 AND contact_id = $2`, + f.campaign, f.contact); err != nil { + t.Fatalf("re-dispatch: %v", err) + } + if inserted, err := evidence.Record(ctx, f.contact, step, "opened", "fresh", "", time.Now()); err != nil || !inserted { + t.Fatalf("an observation about the new address was refused (inserted=%v, err=%v)", inserted, err) + } + if inserted, err := evidence.Record(ctx, f.contact, models.EvidenceStep{}, "replied", "no-step", "", time.Now()); err != nil || !inserted { + t.Fatalf("an observation with no step was refused (inserted=%v, err=%v)", inserted, err) + } + if count() != 2 { + t.Fatalf("evidence after the two allowed observations = %d, want 2", rows) + } +} + +func TestLiveContactEmailRefusesCollisionAndGarbage(t *testing.T) { + handle, pool := liveContactDB(t) + f := newSharedOrgFixture(t, pool) + repo := NewContactRepostory(handle) + ctx := context.Background() + mate := f.mate.String() + + other := uuid.New() + if _, err := pool.Exec(ctx, ` + INSERT INTO contacts (id, user_id, organization_id, email, first_name, last_name, company, phone, custom_fields, updated_at, created_at) + VALUES ($1, $2, $3, 'taken@acme.test', 'Sam', 'Ruiz', '', '', '{}'::jsonb, NOW(), NOW())`, + other, f.owner, f.org); err != nil { + t.Fatalf("second contact: %v", err) + } + + taken := "Taken@Acme.test" + _, xerr := repo.Update(ctx, mate, f.contact.String(), f.org, &models.UpdateContact{Email: &taken}) + if xerr == nil || xerr.Code != errx.Conflict { + t.Fatalf("colliding address = %v, want a 409", xerr) + } + + for _, bad := range []string{"", " ", "not-an-address", "dana@", "@acme.test"} { + v := bad + _, xerr := repo.Update(ctx, mate, f.contact.String(), f.org, &models.UpdateContact{Email: &v}) + if xerr == nil || xerr.Code != errx.BadRequest { + t.Fatalf("address %q = %v, want a 400", bad, xerr) + } + } + + // Creating one goes through the same normalizer, so the two paths cannot + // disagree about what an address is: mail.ParseAddress accepts a display + // name and the whole string used to be stored as the recipient. + created, xerr := repo.Add(ctx, f.owner.String(), f.org, []models.AddContact{{ + FirstName: "Dana", Email: " Dana Reyes ", + }}) + if xerr != nil || len(created) != 1 { + t.Fatalf("add: %v", xerr) + } + if created[0].Email != "dana@created.test" { + t.Fatalf("created contact email = %q, want dana@created.test", created[0].Email) + } + + // Nothing above may have written. + var addr string + if err := pool.QueryRow(ctx, `SELECT email FROM contacts WHERE id = $1`, f.contact).Scan(&addr); err != nil { + t.Fatalf("read contact: %v", err) + } + if want := "i187-" + f.contact.String()[:8] + "@test.local"; addr != want { + t.Fatalf("email = %q, want the original %q", addr, want) + } +} diff --git a/internal/repository/pg_contact.go b/internal/repository/pg_contact.go index 11fbde0a8..02f3cd17f 100644 --- a/internal/repository/pg_contact.go +++ b/internal/repository/pg_contact.go @@ -178,10 +178,16 @@ func (r *contactRepository) Add(ctx context.Context, userID string, orgID uuid.U categoryIDs := make([][]uuid.UUID, 0, len(contacts)) segmentIDs := make([][]uuid.UUID, 0, len(contacts)) for _, lead := range contacts { - lead.Email = strings.TrimSpace(lead.Email) - if !email.IsValid(lead.Email) { + // Normalize, not just trim: mail.ParseAddress accepts + // `Dana Reyes ` and the whole string used to be stored + // as the recipient address, which sends to nobody. The edit path + // normalizes the same way, so the two cannot disagree about what an + // address is. + addr, ok := email.Normalize(lead.Email) + if !ok { return nil, errx.ErrEmail } + lead.Email = addr lead.FirstName = strings.TrimSpace(lead.FirstName) lead.LastName = strings.TrimSpace(lead.LastName) lead.Company = strings.TrimSpace(lead.Company) @@ -2054,11 +2060,15 @@ func (r *contactRepository) Update(ctx context.Context, userID, contactID string // Validate contact existence and fetch current data var c models.Contact var campaignsJSON []byte + // The owner is read for the email-uniqueness check: the unique index is + // (user_id, lower(email)), which an org-scoped check alone does not cover + // for a member who owns contacts in more than one workspace. + var ownerID uuid.UUID query := ` SELECT c.id, c.first_name, c.last_name, c.email, c.company, c.phone, - c.custom_fields, c.subscribed, c.updated_at, c.created_at, + c.custom_fields, c.subscribed, c.updated_at, c.created_at, c.user_id, COALESCE( ( SELECT json_agg(json_build_object('id', cam.id, 'name', cam.name)) @@ -2084,7 +2094,7 @@ func (r *contactRepository) Update(ctx context.Context, userID, contactID string ).Scan( &c.ID, &c.FirstName, &c.LastName, &c.Email, &c.Company, &c.Phone, &c.CustomFields, &c.Subscribed, - &c.UpdatedAt, &c.CreatedAt, &campaignsJSON, + &c.UpdatedAt, &c.CreatedAt, &ownerID, &campaignsJSON, ) if err == pgx.ErrNoRows { return nil, errx.ErrNotFound @@ -2123,6 +2133,67 @@ func (r *contactRepository) Update(ctx context.Context, userID, contactID string var args []interface{} argIndex := 1 + // The address is the contact's identity: a changed one is checked for a + // collision here rather than left to the unique index, which surfaces as a + // 500, and it invalidates every verdict and observation the old mailbox + // earned (reset below, with the evidence rows dropped after the update). + emailChanged := false + if data.Email != nil { + next, ok := email.Normalize(*data.Email) + if !ok { + return nil, errx.ErrEmail + } + // Two different comparisons. A stored address that predates + // normalization can differ from `next` only in case, which is still a + // write (the row is normalized) but not a different mailbox, so it must + // not throw away a verdict the address earned. + if next != c.Email { + setClauses = append(setClauses, fmt.Sprintf("email = $%d", argIndex)) + args = append(args, next) + argIndex++ + } + if next != strings.ToLower(c.Email) { + var taken bool + dupQ := `SELECT EXISTS ( + SELECT 1 FROM contacts + WHERE LOWER(email) = $1 AND id <> $2 AND (organization_id = $3 OR user_id = $4) + )` + dupP := []any{next, contactID, orgID, ownerID} + if err := tx.QueryRow(ctx, dupQ, dupP...).Scan(&taken); err != nil { + db.CaptureError(err, dupQ, dupP, "queryrow") + return nil, errx.InternalError() + } + if taken { + return nil, errx.ErrContactEmailTaken + } + emailChanged = true + setClauses = append(setClauses, + "verification_status = 'unknown'", + "verification_sub_status = ''", + "verification_reason = ''", + "verification_source = ''", + "verification_provider = ''", + "is_catch_all = false", + "verification_checked_at = NULL", + "verification_confidence = 0", + "verification_evidence_at = NULL", + // The ledger the verdict is scored from is wiped below, and + // this is the watermark that keeps it wiped: the delivery + // credit job re-derives 'delivered' rows from every step ever + // sent, so without it the old mailbox's deliveries come back on + // the next pass and hand the new address a verdict it never + // earned. + "verification_evidence_reset_at = NOW()", + // esp_provider is derived from the address domain and cached + // forever: the scheduler only fills it when it is empty, so a + // gmail-to-outlook correction would keep routing ESP-matched + // sends by the old provider. + "esp_provider = ''", + "esp_resolved_at = NULL", + ) + } + } + // Update fields if provided if data.FirstName != nil { setClauses = append(setClauses, fmt.Sprintf("first_name = $%d", argIndex)) @@ -2193,6 +2264,13 @@ func (r *contactRepository) Update(ctx context.Context, userID, contactID string if err == pgx.ErrNoRows { return nil, errx.ErrNotFound } + // The collision check above is a read, so two edits moving two + // contacts onto one address can both pass it and the index + // decides. The loser gets the same answer it would have got a + // moment earlier rather than a 500. + if isUniqueViolation(err) { + return nil, errx.ErrContactEmailTaken + } db.CaptureError(err, query, args, "queryrow") return nil, errx.InternalError() } @@ -2200,6 +2278,16 @@ func (r *contactRepository) Update(ctx context.Context, userID, contactID string updatedContact = c // No fields updated, use existing contact } + // The evidence ledger is a record of what a mailbox did, so it follows the + // address rather than the row. Leaving it would let the scorer hand the new + // address a verdict earned by the old one. + if emailChanged { + if _, err := tx.Exec(ctx, `DELETE FROM contact_verification_evidence WHERE contact_id = $1`, contactID); err != nil { + db.CaptureError(err, "", nil, "verification evidence wipe") + return nil, errx.InternalError() + } + } + // Campaigns are organization assets: scoping membership by the caller made a // teammate's add/remove a silent no-op on campaigns they did not create // (issue #187). diff --git a/internal/repository/pg_verification_evidence.go b/internal/repository/pg_verification_evidence.go index 5f70e0e0c..416e84c23 100644 --- a/internal/repository/pg_verification_evidence.go +++ b/internal/repository/pg_verification_evidence.go @@ -18,8 +18,11 @@ import ( // verification verdict and the derived score on the contact row. type VerificationEvidenceRepository interface { // Record stores one observation. Idempotent on (contact, kind, ref); - // returns whether a new row was written. - Record(ctx context.Context, contactID uuid.UUID, kind, ref, detail string, observedAt time.Time) (bool, error) + // returns whether a new row was written. step names the campaign step the + // observation came from, so an observation about an address the contact no + // longer holds can be refused; a zero step means "not attributable to a + // step", which is always recorded. + Record(ctx context.Context, contactID uuid.UUID, step models.EvidenceStep, kind, ref, detail string, observedAt time.Time) (bool, error) // ListForContact returns the contact's evidence, newest first. ListForContact(ctx context.Context, contactID uuid.UUID) ([]models.ContactVerificationEvidence, error) // Verdict reads the contact's current check verdict for scoring. @@ -40,16 +43,28 @@ func NewVerificationEvidenceRepository(database *db.DB) VerificationEvidenceRepo return &verificationEvidenceRepository{DB: database} } -func (r *verificationEvidenceRepository) Record(ctx context.Context, contactID uuid.UUID, kind, ref, detail string, observedAt time.Time) (bool, error) { +func (r *verificationEvidenceRepository) Record(ctx context.Context, contactID uuid.UUID, step models.EvidenceStep, kind, ref, detail string, observedAt time.Time) (bool, error) { if observedAt.IsZero() { observedAt = time.Now().UTC() } + // A bounce or an open lands long after the mail that produced it, and + // editing a contact's address in between does not make the old mailbox's + // behaviour evidence about the new one: a delayed hard bounce for a typo + // would otherwise mark the corrected address invalid and stop every send + // to it. The refusal needs proof, so it only fires when the step is known + // AND provably left before the address changed; anything else is recorded. query := ` INSERT INTO contact_verification_evidence (contact_id, kind, ref, detail, observed_at) - VALUES ($1, $2, $3, $4, $5) + SELECT $1, $2, $3, $4, $5 + FROM contacts c + WHERE c.id = $1 + AND (c.verification_evidence_reset_at IS NULL OR NOT EXISTS ( + SELECT 1 FROM campaign_contact_progress p + WHERE p.campaign_id = $6 AND p.contact_id = $1 AND p.sequence_id = $7 + AND COALESCE(p.dispatched_at, p.sent_at) <= c.verification_evidence_reset_at)) ON CONFLICT (contact_id, kind, ref) DO NOTHING ` - params := []any{contactID, kind, ref, detail, observedAt} + params := []any{contactID, kind, ref, detail, observedAt, step.CampaignID, step.SequenceID} cmd, err := r.DB.Exec(ctx, query, params...) if err != nil { db.CaptureError(err, query, params, "exec") @@ -134,12 +149,21 @@ func (r *verificationEvidenceRepository) CreditCleanDeliveries(ctx context.Conte // One evidence row per sent step; the ref is the step so a re-run of the // job is a no-op, and a step that bounces later is excluded here and // recorded as a bounce by the deliverability path instead. + // The join to contacts is what keeps a corrected address from inheriting + // the old mailbox's record: this credit is derived from every step ever + // sent, with no lower bound of its own, so the evidence an address change + // deletes would come straight back on the next pass. Steps sent before the + // change went to a different mailbox and are not evidence about this one. + // dispatched_at, not sent_at: sent_at is stamped when the worker's result + // lands, which for a send already on the bus can be after the edit. query := ` WITH due AS ( SELECT p.contact_id, p.campaign_id, p.sequence_id, p.sent_at FROM campaign_contact_progress p + JOIN contacts c ON c.id = p.contact_id WHERE p.sent_at IS NOT NULL AND p.bounced_at IS NULL AND p.sent_at < NOW() - make_interval(secs => $1) + AND (c.verification_evidence_reset_at IS NULL OR COALESCE(p.dispatched_at, p.sent_at) > c.verification_evidence_reset_at) AND NOT EXISTS ( SELECT 1 FROM contact_verification_evidence e WHERE e.contact_id = p.contact_id AND e.kind = 'delivered' diff --git a/internal/repository/verification_evidence_live_test.go b/internal/repository/verification_evidence_live_test.go index 9583be5ab..6d191bf27 100644 --- a/internal/repository/verification_evidence_live_test.go +++ b/internal/repository/verification_evidence_live_test.go @@ -7,6 +7,7 @@ import ( "github.com/google/uuid" + "github.com/warmbly/warmbly/internal/models" "github.com/warmbly/warmbly/internal/pkg/emailverify" ) @@ -57,11 +58,11 @@ func TestLiveVerificationEvidenceLedger(t *testing.T) { } } - inserted, err := repo.Record(ctx, contact, emailverify.EvidenceReplied, "msg-1", "", time.Now()) + inserted, err := repo.Record(ctx, contact, models.EvidenceStep{}, emailverify.EvidenceReplied, "msg-1", "", time.Now()) if err != nil || !inserted { t.Fatalf("record: %v %v", inserted, err) } - if dup, _ := repo.Record(ctx, contact, emailverify.EvidenceReplied, "msg-1", "", time.Now()); dup { + if dup, _ := repo.Record(ctx, contact, models.EvidenceStep{}, emailverify.EvidenceReplied, "msg-1", "", time.Now()); dup { t.Fatal("same reply recorded twice") } rows, err := repo.ListForContact(ctx, contact)