diff --git a/internal/api/handler/contact_detail.go b/internal/api/handler/contact_detail.go new file mode 100644 index 00000000..dcb208ba --- /dev/null +++ b/internal/api/handler/contact_detail.go @@ -0,0 +1,127 @@ +package handler + +import ( + "net/http" + "strconv" + "time" + + "github.com/gin-gonic/gin" + "github.com/google/uuid" + "github.com/warmbly/warmbly/internal/api/middleware" + "github.com/warmbly/warmbly/internal/errx" +) + +// GetContact returns the hydrated contact 360 payload. Used by the +// slide-over on open to render the Overview / engagement-summary tab +// in a single round-trip. Organization is optional — when the caller +// has none selected we still return core + engagement, just without +// org-scoped suppression / complaint counts. +func (h *Handler) GetContact(c *gin.Context) { + userID, err := middleware.GetUserUUID(c) + if err != nil { + errx.Handle(c, errx.ErrAuth) + return + } + contactID, err := uuid.Parse(c.Param("id")) + if err != nil { + errx.Handle(c, errx.ErrUuid) + return + } + + orgID := middleware.GetOrganizationID(c) + + detail, xerr := h.ContactService.GetDetail(c.Request.Context(), userID, orgID, contactID) + if xerr != nil { + errx.Handle(c, xerr) + return + } + + c.JSON(http.StatusOK, detail) +} + +// ListContactEmails returns one row per email we sent (or tried to +// send) to the contact. Cursor pagination keyed on the task ID. +func (h *Handler) ListContactEmails(c *gin.Context) { + userID, err := middleware.GetUserUUID(c) + if err != nil { + errx.Handle(c, errx.ErrAuth) + return + } + contactID, err := uuid.Parse(c.Param("id")) + if err != nil { + errx.Handle(c, errx.ErrUuid) + return + } + + limit := 50 + if l, err := strconv.Atoi(c.Query("limit")); err == nil && l > 0 && l <= 200 { + limit = l + } + + // The cursor for sent-email pagination is the (created_at, task_id) + // pair from the last row of the previous page. Both must be present + // or we skip the cursor and start fresh. + var beforeAt *time.Time + var beforeID *uuid.UUID + if v := c.Query("before_at"); v != "" { + if t, perr := time.Parse(time.RFC3339Nano, v); perr == nil { + beforeAt = &t + } + } + if v := c.Query("before_id"); v != "" { + if id, perr := uuid.Parse(v); perr == nil { + beforeID = &id + } + } + + res, xerr := h.ContactService.ListSentEmails(c.Request.Context(), userID, contactID, limit, beforeAt, beforeID) + if xerr != nil { + errx.Handle(c, xerr) + return + } + + c.JSON(http.StatusOK, res) +} + +// ListContactTimeline returns the merged activity feed for a contact. +// Suppression / deliverability / reply / note events are org-scoped +// and require a selected organization; we 400 if there isn't one to +// avoid silently returning a misleadingly thin feed. +func (h *Handler) ListContactTimeline(c *gin.Context) { + userID, err := middleware.GetUserUUID(c) + if err != nil { + errx.Handle(c, errx.ErrAuth) + return + } + contactID, err := uuid.Parse(c.Param("id")) + if err != nil { + errx.Handle(c, errx.ErrUuid) + return + } + + orgID := middleware.GetOrganizationID(c) + if orgID == nil { + errx.Handle(c, errx.New(errx.BadRequest, "no organization selected")) + return + } + + limit := 50 + if l, err := strconv.Atoi(c.Query("limit")); err == nil && l > 0 && l <= 200 { + limit = l + } + + var before *time.Time + if v := c.Query("before"); v != "" { + if t, perr := time.Parse(time.RFC3339Nano, v); perr == nil { + before = &t + } + } + + res, xerr := h.ContactService.ListTimeline(c.Request.Context(), userID, orgID, contactID, limit, before) + if xerr != nil { + errx.Handle(c, xerr) + return + } + + c.JSON(http.StatusOK, res) +} diff --git a/internal/api/routes.go b/internal/api/routes.go index a211099d..b6015f75 100644 --- a/internal/api/routes.go +++ b/internal/api/routes.go @@ -179,6 +179,12 @@ func Run( contacts.PATCH("/:id", m.RequireAccess(models.PermManageContacts, models.APIPermWriteContacts), h.UpdateContact) contacts.DELETE("/:id", m.RequireAccess(models.PermManageContacts, models.APIPermWriteContacts), h.DeleteContact) + // Contact 360 view: hydrated detail, every email sent to + // the contact, and the merged activity timeline. + contacts.GET("/:id", m.RequireAccess(models.PermViewContacts, models.APIPermReadContacts), h.GetContact) + contacts.GET("/:id/emails", m.RequireAccess(models.PermViewContacts, models.APIPermReadContacts), h.ListContactEmails) + contacts.GET("/:id/timeline", m.RequireAccess(models.PermViewContacts, models.APIPermReadContacts), h.ListContactTimeline) + // CRM: Notes & Activities (under contacts) contacts.GET("/:id/notes", m.RequireAccess(models.PermViewContacts, models.APIPermReadContacts), h.ListContactNotes) contacts.POST("/:id/notes", m.RequireAccess(models.PermManageContacts, models.APIPermWriteContacts), h.CreateContactNote) diff --git a/internal/app/contact/handler.go b/internal/app/contact/handler.go index fddcd705..6e22bf4a 100644 --- a/internal/app/contact/handler.go +++ b/internal/app/contact/handler.go @@ -2,6 +2,7 @@ package contact import ( "context" + "time" "github.com/google/uuid" "github.com/warmbly/warmbly/internal/errx" @@ -66,3 +67,15 @@ func (s *contactService) BulkDelete(ctx context.Context, userID string, contactI func (s *contactService) Delete(ctx context.Context, userID string, contactID string) *errx.Error { return s.contactRepository.Delete(ctx, userID, contactID) } + +func (s *contactService) GetDetail(ctx context.Context, userID uuid.UUID, orgID *uuid.UUID, contactID uuid.UUID) (*models.ContactDetail, *errx.Error) { + return s.contactRepository.GetDetail(ctx, userID, orgID, contactID) +} + +func (s *contactService) ListSentEmails(ctx context.Context, userID, contactID uuid.UUID, limit int, beforeSentAt *time.Time, beforeTaskID *uuid.UUID) (*models.ContactSentEmailsResult, *errx.Error) { + return s.contactRepository.ListSentEmails(ctx, userID, contactID, limit, beforeSentAt, beforeTaskID) +} + +func (s *contactService) ListTimeline(ctx context.Context, userID uuid.UUID, orgID *uuid.UUID, contactID uuid.UUID, limit int, before *time.Time) (*models.ContactTimelineResult, *errx.Error) { + return s.contactRepository.ListTimeline(ctx, userID, orgID, contactID, limit, before) +} diff --git a/internal/app/contact/service.go b/internal/app/contact/service.go index 8cea4df0..06c3d1ba 100644 --- a/internal/app/contact/service.go +++ b/internal/app/contact/service.go @@ -3,7 +3,9 @@ package contact import ( "context" "io" + "time" + "github.com/google/uuid" "github.com/warmbly/warmbly/internal/errx" "github.com/warmbly/warmbly/internal/models" "github.com/warmbly/warmbly/internal/repository" @@ -30,6 +32,18 @@ type ContactService interface { // and performs the upsert / skip / dedup work. Returns per-row // result counts plus a list of rows that failed (with reasons). ImportCommit(ctx context.Context, userID string, file io.Reader, filename string, opts *models.ContactImportCommit) (*models.ContactImportResult, *errx.Error) + + // GetDetail returns the 360 read model used by the contact + // slide-over: hydrated contact + engagement summary + suppression. + GetDetail(ctx context.Context, userID uuid.UUID, orgID *uuid.UUID, contactID uuid.UUID) (*models.ContactDetail, *errx.Error) + + // ListSentEmails enumerates every send (or attempted send) we made + // to the contact, newest first. + ListSentEmails(ctx context.Context, userID, contactID uuid.UUID, limit int, beforeSentAt *time.Time, beforeTaskID *uuid.UUID) (*models.ContactSentEmailsResult, *errx.Error) + + // ListTimeline returns a merged, reverse-chronological feed of all + // engagement + CRM events for the contact. + ListTimeline(ctx context.Context, userID uuid.UUID, orgID *uuid.UUID, contactID uuid.UUID, limit int, before *time.Time) (*models.ContactTimelineResult, *errx.Error) } type contactService struct { diff --git a/internal/models/contact.go b/internal/models/contact.go index f8626b2d..8190d8c2 100644 --- a/internal/models/contact.go +++ b/internal/models/contact.go @@ -39,6 +39,133 @@ type ContactsResult struct { Pagination Pagination `json:"pagination"` } +// ContactEngagement summarises every email touchpoint we have for a +// single contact. It's denormalised on read so the contact 360 view +// can render counts and "last X" timestamps in a single round-trip. +type ContactEngagement struct { + TotalSent int `json:"total_sent"` + TotalOpened int `json:"total_opened"` + TotalClicked int `json:"total_clicked"` + TotalReplied int `json:"total_replied"` + TotalBounced int `json:"total_bounced"` + TotalComplained int `json:"total_complained"` + + LastSentAt *time.Time `json:"last_sent_at,omitempty"` + LastOpenedAt *time.Time `json:"last_opened_at,omitempty"` + LastClickedAt *time.Time `json:"last_clicked_at,omitempty"` + LastRepliedAt *time.Time `json:"last_replied_at,omitempty"` + LastBouncedAt *time.Time `json:"last_bounced_at,omitempty"` +} + +// ContactSuppression mirrors a row from suppressed_recipients for the +// contact's email. Null on the wire when the contact is not suppressed. +type ContactSuppression struct { + Reason string `json:"reason"` + Source string `json:"source"` // bounce | complaint | unsubscribe + ExpiresAt *time.Time `json:"expires_at,omitempty"` + CreatedAt time.Time `json:"created_at"` +} + +// ContactDetail is the hydrated read model returned by GET /contacts/:id. +// It bundles everything the slide-over needs in one payload so the UI +// doesn't have to fan out a half-dozen requests on open. +type ContactDetail struct { + Contact + Engagement ContactEngagement `json:"engagement"` + Suppression *ContactSuppression `json:"suppression,omitempty"` +} + +// ContactSentEmail is one row in the "Emails sent to this contact" +// list. Each row corresponds to a single delivered (or attempted) +// task. Engagement timestamps come from campaign_contact_progress +// when present; some legacy rows may have nil progress. +type ContactSentEmail struct { + TaskID uuid.UUID `json:"task_id"` + Status string `json:"status"` + MessageID string `json:"message_id"` + Subject string `json:"subject"` + SentAt time.Time `json:"sent_at"` + + // Sender mailbox + EmailAccountID *uuid.UUID `json:"email_account_id,omitempty"` + EmailAccountEmail *string `json:"email_account_email,omitempty"` + EmailAccountName *string `json:"email_account_name,omitempty"` + + // Campaign + sequence context + CampaignID *uuid.UUID `json:"campaign_id,omitempty"` + CampaignName *string `json:"campaign_name,omitempty"` + SequenceID *uuid.UUID `json:"sequence_id,omitempty"` + SequenceName *string `json:"sequence_name,omitempty"` + + // Engagement (from campaign_contact_progress, may be nil). + OpenedAt *time.Time `json:"opened_at,omitempty"` + ClickedAt *time.Time `json:"clicked_at,omitempty"` + RepliedAt *time.Time `json:"replied_at,omitempty"` + BouncedAt *time.Time `json:"bounced_at,omitempty"` +} + +type ContactSentEmailsResult struct { + Data []ContactSentEmail `json:"data"` + Pagination Pagination `json:"pagination"` +} + +// ContactTimelineEventType is a closed enum so the frontend can pick +// the right icon / colour without parsing free text. +type ContactTimelineEventType string + +const ( + TimelineEmailSent ContactTimelineEventType = "email_sent" + TimelineEmailOpened ContactTimelineEventType = "email_opened" + TimelineEmailClicked ContactTimelineEventType = "email_clicked" + TimelineEmailReplied ContactTimelineEventType = "email_replied" + TimelineEmailBounced ContactTimelineEventType = "email_bounced" + TimelineReplyReceived ContactTimelineEventType = "reply_received" + TimelineDeliverability ContactTimelineEventType = "deliverability" + TimelineSuppressed ContactTimelineEventType = "suppressed" + TimelineNote ContactTimelineEventType = "note" +) + +// ContactTimelineEvent is one entry in the merged activity feed. The +// optional fields are tagged with omitempty so the JSON stays compact +// for event types that don't carry that data. +type ContactTimelineEvent struct { + Type ContactTimelineEventType `json:"type"` + At time.Time `json:"at"` + + // Mailbox sender (email_sent / opened / clicked / replied / bounced). + EmailAccountID *uuid.UUID `json:"email_account_id,omitempty"` + EmailAccountEmail *string `json:"email_account_email,omitempty"` + EmailAccountName *string `json:"email_account_name,omitempty"` + + // Campaign / sequence linkage. Optional because notes, suppression, + // and out-of-campaign reply intents don't always have one. + CampaignID *uuid.UUID `json:"campaign_id,omitempty"` + CampaignName *string `json:"campaign_name,omitempty"` + SequenceID *uuid.UUID `json:"sequence_id,omitempty"` + SequenceName *string `json:"sequence_name,omitempty"` + + // Task linkage for engagement events. + TaskID *uuid.UUID `json:"task_id,omitempty"` + Subject *string `json:"subject,omitempty"` + + // Type-specific. + Reason *string `json:"reason,omitempty"` // deliverability / suppression + Source *string `json:"source,omitempty"` // suppression: bounce/complaint/unsubscribe + Provider *string `json:"provider,omitempty"` // deliverability provider + Intent *string `json:"intent,omitempty"` // reply_intent classification + Content *string `json:"content,omitempty"` // note body + + // Author (notes). + UserID *uuid.UUID `json:"user_id,omitempty"` +} + +type ContactTimelineResult struct { + Data []ContactTimelineEvent `json:"data"` + // True if we hit the per-call cap and the caller should paginate + // via the `before` query param. + HasMore bool `json:"has_more"` +} + type UpdateContact struct { FirstName *string `json:"first_name"` LastName *string `json:"last_name"` diff --git a/internal/repository/pg_contact.go b/internal/repository/pg_contact.go index 1d463bd0..c621b661 100644 --- a/internal/repository/pg_contact.go +++ b/internal/repository/pg_contact.go @@ -4,7 +4,9 @@ import ( "context" "encoding/json" "fmt" + "sort" "strings" + "time" "github.com/getsentry/sentry-go" "github.com/google/uuid" @@ -30,6 +32,12 @@ type ContactRepository interface { BulkDelete(ctx context.Context, userID string, contactIDs []string) *errx.Error Delete(ctx context.Context, userID string, contactID string) *errx.Error GetContactCount(ctx context.Context, userID string) (int, *errx.Error) + + // 360 view read paths. orgID is optional — when nil, the suppression + // + deliverability + reply joins are skipped (they're org-scoped). + GetDetail(ctx context.Context, userID uuid.UUID, orgID *uuid.UUID, contactID uuid.UUID) (*models.ContactDetail, *errx.Error) + ListSentEmails(ctx context.Context, userID, contactID uuid.UUID, limit int, beforeSentAt *time.Time, beforeTaskID *uuid.UUID) (*models.ContactSentEmailsResult, *errx.Error) + ListTimeline(ctx context.Context, userID uuid.UUID, orgID *uuid.UUID, contactID uuid.UUID, limit int, before *time.Time) (*models.ContactTimelineResult, *errx.Error) } type contactRepository struct { @@ -1483,3 +1491,508 @@ func (r *contactRepository) GetContactCount(ctx context.Context, userID string) } return count, nil } + +// GetDetail loads the contact 360 payload: core fields + categories + +// campaigns + engagement counts + suppression. Single round-trip via a +// few separate queries (one main select + a couple of small aggregates) +// so the query plans stay simple and cheap to reason about. +// +// orgID is optional because not every caller has an org context (e.g. +// an API key scoped to a user without a selected org). When nil we +// skip the org-scoped joins (suppression, deliverability) and return +// zeros for those fields. +func (r *contactRepository) GetDetail(ctx context.Context, userID uuid.UUID, orgID *uuid.UUID, contactID uuid.UUID) (*models.ContactDetail, *errx.Error) { + // 1. Core contact + categories + campaigns. Same shape as Search + // so the UI gets identical fields back. + var detail models.ContactDetail + var campaignsJSON, categoriesJSON []byte + mainQuery := ` + 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, + COALESCE( + ( + SELECT json_agg(json_build_object('id', cam.id, 'name', cam.name)) + FROM campaign_leads cl + JOIN campaigns cam ON cam.id = cl.campaign_id + WHERE cl.contact_id = c.id AND cam.user_id = $1 + ), '[]'::json + ) AS campaigns, + COALESCE( + ( + SELECT json_agg(json_build_object('id', cat.id, 'title', cat.title, 'color', cat.color) ORDER BY cat.position ASC, cat.title ASC) + FROM contact_categories cc + JOIN categories cat ON cat.id = cc.category_id + WHERE cc.contact_id = c.id AND cat.user_id = $1 + ), '[]'::json + ) AS categories + FROM contacts c + WHERE c.id = $2 AND c.user_id = $1 + ` + err := r.DB.QueryRow(ctx, mainQuery, userID, contactID).Scan( + &detail.ID, &detail.FirstName, &detail.LastName, &detail.Email, + &detail.Company, &detail.Phone, &detail.CustomFields, &detail.Subscribed, + &detail.UpdatedAt, &detail.CreatedAt, &campaignsJSON, &categoriesJSON, + ) + if err != nil { + if err == pgx.ErrNoRows { + return nil, errx.ErrNotFound + } + db.CaptureError(err, mainQuery, []any{userID, contactID}, "GetDetail main") + return nil, errx.InternalError() + } + if detail.CustomFields == nil { + detail.CustomFields = map[string]string{} + } + detail.Campaigns = []models.MiniCampaign{} + if len(campaignsJSON) > 0 { + var raw []struct { + ID string `json:"id"` + Name string `json:"name"` + } + if err := json.Unmarshal(campaignsJSON, &raw); err != nil { + sentry.CaptureException(err) + return nil, errx.InternalError() + } + detail.Campaigns = make([]models.MiniCampaign, len(raw)) + for i, m := range raw { + detail.Campaigns[i] = models.MiniCampaign{ID: m.ID, Name: m.Name} + } + } + detail.Categories = []models.MiniCategory{} + if len(categoriesJSON) > 0 { + if err := json.Unmarshal(categoriesJSON, &detail.Categories); err != nil { + sentry.CaptureException(err) + return nil, errx.InternalError() + } + } + + // 2. Engagement aggregates. campaign_contact_progress is the canonical + // sent/opened/clicked/replied/bounced ledger keyed by (campaign, + // contact, sequence). Counts come from non-null timestamp columns, + // "last X" comes from MAX() of each. + engQuery := ` + SELECT + COUNT(*) FILTER (WHERE sent_at IS NOT NULL) AS sent, + COUNT(*) FILTER (WHERE opened_at IS NOT NULL) AS opened, + COUNT(*) FILTER (WHERE clicked_at IS NOT NULL) AS clicked, + COUNT(*) FILTER (WHERE replied_at IS NOT NULL) AS replied, + COUNT(*) FILTER (WHERE bounced_at IS NOT NULL) AS bounced, + MAX(sent_at), MAX(opened_at), MAX(clicked_at), MAX(replied_at), MAX(bounced_at) + FROM campaign_contact_progress + WHERE contact_id = $1 + ` + if err := r.DB.QueryRow(ctx, engQuery, contactID).Scan( + &detail.Engagement.TotalSent, &detail.Engagement.TotalOpened, + &detail.Engagement.TotalClicked, &detail.Engagement.TotalReplied, + &detail.Engagement.TotalBounced, + &detail.Engagement.LastSentAt, &detail.Engagement.LastOpenedAt, + &detail.Engagement.LastClickedAt, &detail.Engagement.LastRepliedAt, + &detail.Engagement.LastBouncedAt, + ); err != nil { + db.CaptureError(err, engQuery, []any{contactID}, "GetDetail engagement") + return nil, errx.InternalError() + } + + // 3. Org-scoped extras. Only run when we have an org id. + if orgID != nil { + // Complaints don't live in campaign_contact_progress — they + // arrive via deliverability_events. Count rows of type + // "complaint" pointing at this contact (either by contact_id + // or by recipient_email fallback for older rows). + complaintQuery := ` + SELECT COUNT(*) + FROM deliverability_events + WHERE organization_id = $1 + AND event_type = 'complaint' + AND (contact_id = $2 OR LOWER(recipient_email) = LOWER($3)) + ` + if err := r.DB.QueryRow(ctx, complaintQuery, *orgID, contactID, detail.Email).Scan( + &detail.Engagement.TotalComplained, + ); err != nil { + db.CaptureError(err, complaintQuery, []any{*orgID, contactID, detail.Email}, "GetDetail complaints") + return nil, errx.InternalError() + } + + // Suppression — there's at most one row per (org, email) + // thanks to the unique constraint. + suppQuery := ` + SELECT reason, source, expires_at, created_at + FROM suppressed_recipients + WHERE organization_id = $1 AND LOWER(email) = LOWER($2) + ` + var s models.ContactSuppression + err := r.DB.QueryRow(ctx, suppQuery, *orgID, detail.Email).Scan( + &s.Reason, &s.Source, &s.ExpiresAt, &s.CreatedAt, + ) + switch { + case err == nil: + detail.Suppression = &s + case err == pgx.ErrNoRows: + // not suppressed; leave nil + default: + db.CaptureError(err, suppQuery, []any{*orgID, detail.Email}, "GetDetail suppression") + return nil, errx.InternalError() + } + } + + return &detail, nil +} + +// ListSentEmails returns one row per task we sent (or attempted to +// send) to the contact, ordered by sent time DESC. Uses keyset +// pagination on (created_at, task_id) so we can scroll through the +// full history without blowing up offset. +// +// We deliberately scope by the contact's owning user via the +// campaign join — this keeps multi-tenant safety even though the +// tasks table itself has no user_id column. +func (r *contactRepository) ListSentEmails(ctx context.Context, userID, contactID uuid.UUID, limit int, beforeSentAt *time.Time, beforeTaskID *uuid.UUID) (*models.ContactSentEmailsResult, *errx.Error) { + if limit <= 0 || limit > 200 { + limit = 50 + } + + args := []any{userID, contactID} + cursorClause := "" + if beforeSentAt != nil && beforeTaskID != nil { + cursorClause = "AND (t.created_at, t.id) < ($3, $4)" + args = append(args, *beforeSentAt, *beforeTaskID) + } + args = append(args, limit+1) + + query := fmt.Sprintf(` + SELECT + t.id, t.status::text, t.message_id, t.created_at, + ea.id, ea.email, ea.name, + cam.id, cam.name, + seq.id, seq.name, + COALESCE(et.subject, seq.subject, '') AS subject, + ccp.opened_at, ccp.clicked_at, ccp.replied_at, ccp.bounced_at + FROM tasks t + JOIN campaign_tasks ct ON ct.task_id = t.id + LEFT JOIN email_accounts ea ON ea.id = t.email_account_id + LEFT JOIN campaigns cam ON cam.id = ct.campaign_id + LEFT JOIN sequences seq ON seq.id = ct.sequence_id + LEFT JOIN email_tasks et ON et.task_id = t.id + LEFT JOIN campaign_contact_progress ccp + ON ccp.campaign_id = ct.campaign_id + AND ccp.contact_id = ct.contact_id + AND ccp.sequence_id = ct.sequence_id + WHERE ct.contact_id = $2 + AND cam.user_id = $1 + %s + ORDER BY t.created_at DESC, t.id DESC + LIMIT $%d + `, cursorClause, len(args)) + + rows, err := r.DB.Query(ctx, query, args...) + if err != nil { + db.CaptureError(err, query, args, "ListSentEmails") + return nil, errx.InternalError() + } + defer rows.Close() + + out := make([]models.ContactSentEmail, 0, limit) + for rows.Next() { + var e models.ContactSentEmail + if err := rows.Scan( + &e.TaskID, &e.Status, &e.MessageID, &e.SentAt, + &e.EmailAccountID, &e.EmailAccountEmail, &e.EmailAccountName, + &e.CampaignID, &e.CampaignName, + &e.SequenceID, &e.SequenceName, + &e.Subject, + &e.OpenedAt, &e.ClickedAt, &e.RepliedAt, &e.BouncedAt, + ); err != nil { + db.CaptureError(err, "", nil, "ListSentEmails scan") + return nil, errx.InternalError() + } + out = append(out, e) + } + + hasMore := false + var nextCursor *uuid.UUID + if len(out) > limit { + hasMore = true + nextCursor = &out[limit].TaskID + out = out[:limit] + } + + return &models.ContactSentEmailsResult{ + Data: out, + Pagination: models.Pagination{ + NextCursor: nextCursor, + HasMore: hasMore, + }, + }, nil +} + +// ListTimeline merges per-contact events from several source tables +// into a single, reverse-chronological feed. +// +// Sources: +// - campaign_contact_progress → sent / opened / clicked / replied / bounced +// - reply_intents → received replies (with intent classification) +// - deliverability_events → bounce / complaint +// - suppressed_recipients → suppression added +// - contact_notes → CRM notes +// +// We pull up to (limit) candidates from each source ordered by time +// DESC, then merge-sort in Go. This avoids a 5-way UNION with +// matching column lists (each source has a different shape), and the +// per-source limit caps the read at roughly 5*limit rows. +// +// The `before` cursor is a wall-clock time; everything strictly older +// than it is eligible. The caller paginates by setting `before` to +// the oldest returned event's `At` on the next call. +func (r *contactRepository) ListTimeline(ctx context.Context, userID uuid.UUID, orgID *uuid.UUID, contactID uuid.UUID, limit int, before *time.Time) (*models.ContactTimelineResult, *errx.Error) { + if limit <= 0 || limit > 200 { + limit = 50 + } + + // We resolve the contact's email up front because some org-scoped + // joins (suppression, deliverability fallback, reply_intents) key + // off email rather than contact_id. + var contactEmail string + if err := r.DB.QueryRow(ctx, + `SELECT email FROM contacts WHERE id = $1 AND user_id = $2`, + contactID, userID, + ).Scan(&contactEmail); err != nil { + if err == pgx.ErrNoRows { + return nil, errx.ErrNotFound + } + db.CaptureError(err, "", nil, "ListTimeline contact email") + return nil, errx.InternalError() + } + + // "before" defaults to "now + 1 minute" so the first page picks + // up everything. Using a future bound keeps the SQL uniform — every + // query passes the same predicate. + bound := time.Now().Add(time.Minute) + if before != nil { + bound = *before + } + + events := make([]models.ContactTimelineEvent, 0, limit*2) + + // 1. Engagement events from campaign_contact_progress. One progress + // row can emit up to 5 events (sent/opened/clicked/replied/bounced). + progressQuery := ` + SELECT + ccp.sent_at, ccp.opened_at, ccp.clicked_at, ccp.replied_at, ccp.bounced_at, + cam.id, cam.name, + seq.id, seq.name, seq.subject, + ea.id, ea.email, ea.name + FROM campaign_contact_progress ccp + JOIN campaigns cam ON cam.id = ccp.campaign_id + JOIN sequences seq ON seq.id = ccp.sequence_id + LEFT JOIN LATERAL ( + SELECT ea.id, ea.email, ea.name + FROM tasks t + JOIN campaign_tasks ct ON ct.task_id = t.id + JOIN email_accounts ea ON ea.id = t.email_account_id + WHERE ct.campaign_id = ccp.campaign_id + AND ct.contact_id = ccp.contact_id + AND ct.sequence_id = ccp.sequence_id + ORDER BY t.created_at DESC + LIMIT 1 + ) ea ON TRUE + WHERE ccp.contact_id = $1 + AND cam.user_id = $2 + AND COALESCE(ccp.sent_at, ccp.opened_at, ccp.clicked_at, ccp.replied_at, ccp.bounced_at) < $3 + ORDER BY GREATEST( + COALESCE(ccp.sent_at, 'epoch'), + COALESCE(ccp.opened_at, 'epoch'), + COALESCE(ccp.clicked_at, 'epoch'), + COALESCE(ccp.replied_at, 'epoch'), + COALESCE(ccp.bounced_at, 'epoch') + ) DESC + LIMIT $4 + ` + prows, err := r.DB.Query(ctx, progressQuery, contactID, userID, bound, limit) + if err != nil { + db.CaptureError(err, progressQuery, []any{contactID, userID, bound, limit}, "ListTimeline progress") + return nil, errx.InternalError() + } + for prows.Next() { + var sentAt, openedAt, clickedAt, repliedAt, bouncedAt *time.Time + var campID, seqID, eaID *uuid.UUID + var campName, seqName, seqSubject, eaEmail, eaName *string + if err := prows.Scan( + &sentAt, &openedAt, &clickedAt, &repliedAt, &bouncedAt, + &campID, &campName, + &seqID, &seqName, &seqSubject, + &eaID, &eaEmail, &eaName, + ); err != nil { + prows.Close() + db.CaptureError(err, "", nil, "ListTimeline progress scan") + return nil, errx.InternalError() + } + baseSubject := seqSubject + makeEvent := func(t *time.Time, ty models.ContactTimelineEventType) { + if t == nil || !t.Before(bound) { + return + } + ev := models.ContactTimelineEvent{ + Type: ty, + At: *t, + EmailAccountID: eaID, + EmailAccountEmail: eaEmail, + EmailAccountName: eaName, + CampaignID: campID, + CampaignName: campName, + SequenceID: seqID, + SequenceName: seqName, + } + if baseSubject != nil && *baseSubject != "" { + ev.Subject = baseSubject + } + events = append(events, ev) + } + makeEvent(sentAt, models.TimelineEmailSent) + makeEvent(openedAt, models.TimelineEmailOpened) + makeEvent(clickedAt, models.TimelineEmailClicked) + makeEvent(repliedAt, models.TimelineEmailReplied) + makeEvent(bouncedAt, models.TimelineEmailBounced) + } + prows.Close() + + if orgID != nil { + // 2. Reply intents (inbound replies with classification). + replyQuery := ` + SELECT ri.created_at, ri.intent, ri.campaign_id, cam.name, ri.task_id + FROM reply_intents ri + LEFT JOIN campaigns cam ON cam.id = ri.campaign_id + WHERE ri.organization_id = $1 + AND LOWER(ri.contact_email) = LOWER($2) + AND ri.created_at < $3 + ORDER BY ri.created_at DESC + LIMIT $4 + ` + rrows, err := r.DB.Query(ctx, replyQuery, *orgID, contactEmail, bound, limit) + if err != nil { + db.CaptureError(err, replyQuery, nil, "ListTimeline replies") + return nil, errx.InternalError() + } + for rrows.Next() { + var ev models.ContactTimelineEvent + var intent string + if err := rrows.Scan(&ev.At, &intent, &ev.CampaignID, &ev.CampaignName, &ev.TaskID); err != nil { + rrows.Close() + db.CaptureError(err, "", nil, "ListTimeline replies scan") + return nil, errx.InternalError() + } + ev.Type = models.TimelineReplyReceived + ev.Intent = &intent + events = append(events, ev) + } + rrows.Close() + + // 3. Deliverability events (bounce / complaint / unsubscribe). + delivQuery := ` + SELECT de.created_at, de.event_type, de.provider, de.reason, + de.campaign_id, cam.name, de.task_id + FROM deliverability_events de + LEFT JOIN campaigns cam ON cam.id = de.campaign_id + WHERE de.organization_id = $1 + AND (de.contact_id = $2 OR LOWER(de.recipient_email) = LOWER($3)) + AND de.created_at < $4 + ORDER BY de.created_at DESC + LIMIT $5 + ` + drows, err := r.DB.Query(ctx, delivQuery, *orgID, contactID, contactEmail, bound, limit) + if err != nil { + db.CaptureError(err, delivQuery, nil, "ListTimeline deliv") + return nil, errx.InternalError() + } + for drows.Next() { + var ev models.ContactTimelineEvent + var eventType, provider, reason string + if err := drows.Scan(&ev.At, &eventType, &provider, &reason, &ev.CampaignID, &ev.CampaignName, &ev.TaskID); err != nil { + drows.Close() + db.CaptureError(err, "", nil, "ListTimeline deliv scan") + return nil, errx.InternalError() + } + ev.Type = models.TimelineDeliverability + ev.Source = &eventType + ev.Provider = &provider + if reason != "" { + ev.Reason = &reason + } + events = append(events, ev) + } + drows.Close() + + // 4. Suppression — emit one event at create time. We treat + // later updates as the same event for now. + suppQuery := ` + SELECT created_at, reason, source + FROM suppressed_recipients + WHERE organization_id = $1 + AND LOWER(email) = LOWER($2) + AND created_at < $3 + ORDER BY created_at DESC + LIMIT 1 + ` + var sAt time.Time + var sReason, sSource string + if err := r.DB.QueryRow(ctx, suppQuery, *orgID, contactEmail, bound).Scan(&sAt, &sReason, &sSource); err == nil { + ev := models.ContactTimelineEvent{ + Type: models.TimelineSuppressed, + At: sAt, + Source: &sSource, + } + if sReason != "" { + ev.Reason = &sReason + } + events = append(events, ev) + } else if err != pgx.ErrNoRows { + db.CaptureError(err, suppQuery, nil, "ListTimeline suppression") + return nil, errx.InternalError() + } + + // 5. Notes. + notesQuery := ` + SELECT created_at, user_id, content + FROM contact_notes + WHERE contact_id = $1 + AND organization_id = $2 + AND created_at < $3 + ORDER BY created_at DESC + LIMIT $4 + ` + nrows, err := r.DB.Query(ctx, notesQuery, contactID, *orgID, bound, limit) + if err != nil { + db.CaptureError(err, notesQuery, nil, "ListTimeline notes") + return nil, errx.InternalError() + } + for nrows.Next() { + var ev models.ContactTimelineEvent + var uid uuid.UUID + var content string + if err := nrows.Scan(&ev.At, &uid, &content); err != nil { + nrows.Close() + db.CaptureError(err, "", nil, "ListTimeline notes scan") + return nil, errx.InternalError() + } + ev.Type = models.TimelineNote + ev.UserID = &uid + ev.Content = &content + events = append(events, ev) + } + nrows.Close() + } + + // Merge sort: newest first. + sort.Slice(events, func(i, j int) bool { return events[i].At.After(events[j].At) }) + + hasMore := false + if len(events) > limit { + hasMore = true + events = events[:limit] + } + + return &models.ContactTimelineResult{ + Data: events, + HasMore: hasMore, + }, nil +}