diff --git a/docs/content/docs/api/reference/contacts.mdx b/docs/content/docs/api/reference/contacts.mdx index f69031b0..6b174571 100644 --- a/docs/content/docs/api/reference/contacts.mdx +++ b/docs/content/docs/api/reference/contacts.mdx @@ -15,7 +15,7 @@ Auth: **Scope** `READ_CONTACTS` · **Org permission** `view_contacts` | Parameter | In | Type | Description | | --- | --- | --- | --- | -| `cursor` | query | string | Opaque pagination cursor from the previous page's `pagination.next_cursor`. | +| `cursor` | query | string | Opaque pagination cursor from the previous page's `pagination.next_cursor`. It carries the exact position of the next page under the ordering it was issued for, so rows that share a sort value are never skipped or repeated. A malformed cursor, or one replayed with a different `sort_by` or `reverse` than it was issued under, is a `400`. | | `limit` | query | string | Page size (numeric string). | | `category` | query | string | Convenience filter for a single category ID. | @@ -40,8 +40,8 @@ Every field is optional; an empty body matches all contacts in the organization. | `created_before` | string (RFC 3339) | No | Created on or before this time. | | `updated_after` | string (RFC 3339) | No | Updated on or after this time. | | `updated_before` | string (RFC 3339) | No | Updated on or before this time. | -| `sort_by` | string | No | Sort column, e.g. `first_name`, `campaign_count`. | -| `reverse` | boolean | No | Descending when true. | +| `sort_by` | string | No | Sort column: `created_at` (the default), `updated_at`, `first_name`, `last_name`, `email`, or `campaign_count`. Any other value falls back to `created_at`. | +| `reverse` | boolean | No | Ascending when true. The default is descending. | ```json { @@ -56,7 +56,7 @@ Every field is optional; an empty body matches all contacts in the organization. ### Response -Returns a `data` array of contacts plus a `pagination` envelope. +Returns a `data` array of contacts plus a `pagination` envelope. Paginate by passing `pagination.next_cursor` back as `cursor` until `has_more` is `false`; `pagination.total` is the count for the whole filtered set and is only sent on the first page. ```json { @@ -86,7 +86,7 @@ Returns a `data` array of contacts plus a `pagination` envelope. ], "pagination": { "total": 1280, - "next_cursor": "c1_b3BhcXVlLWN1cnNvcg", + "next_cursor": "s1_b3BhcXVlLWN1cnNvcg", "has_more": true } } diff --git a/docs/public/openapi.json b/docs/public/openapi.json index e457030d..c9e26bf3 100644 --- a/docs/public/openapi.json +++ b/docs/public/openapi.json @@ -5689,7 +5689,7 @@ "name": "cursor", "in": "query", "required": false, - "description": "Opaque pagination cursor from the previous page's pagination.next_cursor.", + "description": "Opaque pagination cursor from the previous page's pagination.next_cursor. It carries the exact position of the next page under the ordering it was issued for, so rows that share a sort value are never skipped or repeated. A malformed cursor, or one replayed with a different sort_by or reverse than it was issued under, is a 400.", "schema": { "type": "string" } @@ -21786,11 +21786,11 @@ }, "sort_by": { "type": "string", - "description": "Sort column, e.g. first_name, campaign_count." + "description": "Sort column: created_at (the default), updated_at, first_name, last_name, email, or campaign_count. Any other value falls back to created_at." }, "reverse": { "type": "boolean", - "description": "Descending when true." + "description": "Ascending when true. The default is descending." } } }, diff --git a/internal/app/contact/handler.go b/internal/app/contact/handler.go index 77e9af32..ef525648 100644 --- a/internal/app/contact/handler.go +++ b/internal/app/contact/handler.go @@ -97,7 +97,7 @@ func (s *contactService) checkSegmentTargets(ctx context.Context, orgID uuid.UUI } func (s *contactService) Search(ctx context.Context, orgID, cursor, category, limit string, filters models.SearchContacts) (*models.ContactsResult, *errx.Error) { - cursorId, err := paging.DecodeCursor(cursor) + cursorPos, err := paging.DecodeSortCursor(cursor) if err != nil { return nil, err } @@ -115,7 +115,7 @@ func (s *contactService) Search(ctx context.Context, orgID, cursor, category, li return nil, err } - return s.contactRepository.Search(ctx, orgID, categoryId, cursorId, filters, limitN) + return s.contactRepository.Search(ctx, orgID, categoryId, cursorPos, filters, limitN) } // validateLeadFilters gates the single-campaign Leads-view filters: an unknown diff --git a/internal/infrastructure/db/migrations/000136_contacts_org_created_index.down.sql b/internal/infrastructure/db/migrations/000136_contacts_org_created_index.down.sql new file mode 100644 index 00000000..fdf39cbd --- /dev/null +++ b/internal/infrastructure/db/migrations/000136_contacts_org_created_index.down.sql @@ -0,0 +1 @@ +DROP INDEX CONCURRENTLY IF EXISTS idx_contacts_org_created; diff --git a/internal/infrastructure/db/migrations/000136_contacts_org_created_index.up.sql b/internal/infrastructure/db/migrations/000136_contacts_org_created_index.up.sql new file mode 100644 index 00000000..83efa373 --- /dev/null +++ b/internal/infrastructure/db/migrations/000136_contacts_org_created_index.up.sql @@ -0,0 +1,8 @@ +-- The contacts list and every campaign's Leads tab page through one ordering: +-- an organization's contacts newest first, with the row id breaking ties +-- (issue #382). Without this the keyset has to read the whole organization on +-- every page. Descending on both columns so a forward scan serves the default +-- order and a backward scan serves the reversed one. Built concurrently, on its +-- own, because contacts is a live table. +CREATE INDEX CONCURRENTLY IF NOT EXISTS idx_contacts_org_created + ON contacts (organization_id, created_at DESC, id DESC); diff --git a/internal/repository/contact_pagination_live_test.go b/internal/repository/contact_pagination_live_test.go new file mode 100644 index 00000000..3ea47bf5 --- /dev/null +++ b/internal/repository/contact_pagination_live_test.go @@ -0,0 +1,255 @@ +package repository + +import ( + "context" + "fmt" + "testing" + + "github.com/google/uuid" + "github.com/jackc/pgx/v5/pgxpool" + + "github.com/warmbly/warmbly/internal/models" + "github.com/warmbly/warmbly/internal/utils/paging" +) + +// Regression cover for issue #382: a campaign's Leads tab stopped loading after +// a few pages, showing 198 of ~600 leads with no "Load more" left. +// +// The keyset cursor compared the sort column against a subquery that named the +// OUTER row's alias (`SELECT c.created_at FROM contacts WHERE id = $n`), which +// Postgres resolves as a correlated reference to the row being tested. The +// whole boundary collapsed to `c.id >= `, so each page returned +// the newest rows of a randomly shrinking id range: duplicates on screen and a +// list that ran out long before the leads did. The same loop backs ExportAll, +// so exports were truncated the same way. +// +// Run against the dev stack: +// +// WARMBLY_TEST_DB=postgres://warmbly:warmbly@localhost:15432/warmbly_dev?sslmode=disable \ +// go test ./internal/repository/ -run LiveContactPagination -v + +// pagedOrgFixture is one organization whose single campaign holds every contact +// as a lead, sized past several pages. +type pagedOrgFixture struct { + org uuid.UUID + user uuid.UUID + campaign uuid.UUID + total int +} + +func newPagedOrgFixture(t *testing.T, pool *pgxpool.Pool, total int) *pagedOrgFixture { + t.Helper() + ctx := context.Background() + f := &pagedOrgFixture{org: uuid.New(), user: uuid.New(), campaign: uuid.New(), total: total} + + exec := func(sql string, args ...any) { + t.Helper() + if _, err := pool.Exec(ctx, sql, args...); err != nil { + t.Fatalf("fixture %q: %v", sql[:min(60, len(sql))], err) + } + } + exec(`INSERT INTO users (id, first_name, last_name, email, password_hash) + VALUES ($1, 'Paged', 'Live', $2, 'x')`, f.user, "i382-"+f.user.String()[:8]+"@test.local") + exec(`INSERT INTO organizations (id, name, slug, owner_user_id) + VALUES ($1, 'Issue 382', $2, $3)`, f.org, "i382-"+f.org.String()[:8], f.user) + exec(`INSERT INTO organization_members (organization_id, user_id, role, accepted_at) + VALUES ($1, $2, 'owner', NOW())`, f.org, f.user) + exec(`INSERT INTO campaigns (id, user_id, organization_id, name, description, days, updated_at, created_at) + VALUES ($1, $2, $3, 'Q3 cold outreach', '', 62, NOW(), NOW())`, f.campaign, f.user, f.org) + + // A bulk import stamps a whole batch with the same created_at, so most of + // the list is one big tie the id tiebreak has to carry. Names repeat for the + // same reason; the sort must still be a total order. + exec(` + INSERT INTO contacts (id, user_id, organization_id, email, first_name, last_name, company, phone, custom_fields, updated_at, created_at) + SELECT gen_random_uuid(), $1, $2, + 'i382-' || i || '-' || $3 || '@test.local', + 'Lead' || (i % 7), 'Batch' || (i % 3), '', '', '{}'::jsonb, + NOW() - ((i % 5) || ' minutes')::interval, + NOW() - ((i / 100) || ' hours')::interval + FROM generate_series(1, $4) AS i`, + f.user, f.org, f.org.String()[:8], total) + exec(`INSERT INTO campaign_leads (campaign_id, contact_id) + SELECT $1, id FROM contacts WHERE organization_id = $2`, f.campaign, f.org) + + t.Cleanup(func() { + c := context.Background() + for _, step := range []struct { + sql string + arg any + }{ + {`DELETE FROM campaign_leads WHERE campaign_id IN (SELECT id FROM campaigns WHERE organization_id = $1)`, f.org}, + {`DELETE FROM campaigns WHERE organization_id = $1`, f.org}, + {`DELETE FROM contacts WHERE organization_id = $1`, f.org}, + {`DELETE FROM organization_members WHERE organization_id = $1`, f.org}, + {`DELETE FROM organizations WHERE id = $1`, f.org}, + {`DELETE FROM users WHERE id = $1`, f.user}, + } { + if _, err := pool.Exec(c, step.sql, step.arg); err != nil { + t.Errorf("cleanup %q: %v", step.sql, err) + } + } + }) + return f +} + +// pageThrough walks the list the way "Load more" does and returns the ids in +// the order they were served, plus the page count. +func pageThrough(t *testing.T, repo ContactRepository, orgID string, filters models.SearchContacts, limit int32) ([]uuid.UUID, int) { + t.Helper() + ctx := context.Background() + var ( + seen []uuid.UUID + cursor *paging.SortCursor + pages int + ) + for { + page, xerr := repo.Search(ctx, orgID, nil, cursor, filters, limit) + if xerr != nil { + t.Fatalf("page %d: %v", pages, xerr) + } + pages++ + for _, c := range page.Data { + seen = append(seen, c.ID) + } + if !page.Pagination.HasMore { + if page.Pagination.NextCursor != nil { + t.Fatalf("page %d says no more but still handed out a cursor", pages) + } + return seen, pages + } + if page.Pagination.NextCursor == nil { + t.Fatalf("page %d says there is more but handed out no cursor", pages) + } + next, xerr := paging.DecodeSortCursor(*page.Pagination.NextCursor) + if xerr != nil { + t.Fatalf("page %d cursor: %v", pages, xerr) + } + cursor = next + if pages > 200 { + t.Fatal("pagination did not terminate") + } + } +} + +// The issue as reported: ~600 leads, 50 a page, and the list has to hand back +// every lead exactly once. +func TestLiveContactPaginationWalksEveryLead(t *testing.T) { + handle, pool := liveContactDB(t) + f := newPagedOrgFixture(t, pool, 617) + repo := NewContactRepostory(handle) + + seen, pages := pageThrough(t, repo, f.org.String(), models.SearchContacts{ + CampaignIDs: []string{f.campaign.String()}, + }, 50) + + if pages != 13 { + t.Errorf("walked %d pages, want 13 (617 leads at 50 a page)", pages) + } + assertExactlyOnce(t, seen, f.total) +} + +// Every sort the dashboard offers, in both directions. campaign_count sorts on +// a computed column that WHERE cannot reach by its SELECT alias, so a cursor +// there used to be a hard SQL error rather than a short list. +func TestLiveContactPaginationCoversEverySort(t *testing.T) { + handle, pool := liveContactDB(t) + f := newPagedOrgFixture(t, pool, 137) + repo := NewContactRepostory(handle) + + for _, sortBy := range []string{"", "first_name", "last_name", "email", "created_at", "updated_at", "campaign_count"} { + for _, reverse := range []bool{false, true} { + t.Run(fmt.Sprintf("%s/reverse=%v", sortBy, reverse), func(t *testing.T) { + filters := models.SearchContacts{ + CampaignIDs: []string{f.campaign.String()}, + SortBy: sortBy, + Reverse: reverse, + } + // One page holding everything is the order to check against. + whole, xerr := repo.Search(context.Background(), f.org.String(), nil, nil, filters, int32(f.total)) + if xerr != nil { + t.Fatalf("single page: %v", xerr) + } + want := make([]uuid.UUID, 0, f.total) + for _, c := range whole.Data { + want = append(want, c.ID) + } + + seen, _ := pageThrough(t, repo, f.org.String(), filters, 10) + assertExactlyOnce(t, seen, f.total) + for i := range want { + if i < len(seen) && seen[i] != want[i] { + t.Fatalf("paged order diverges from the whole-list order at %d", i) + } + } + }) + } + } +} + +// The org-wide contacts list (no campaign scope) pages the same way, and the +// export walks the same cursor loop, so it has to reach every row too. +func TestLiveContactPaginationExportsEveryRow(t *testing.T) { + handle, pool := liveContactDB(t) + f := newPagedOrgFixture(t, pool, 617) + repo := NewContactRepostory(handle) + + seen, _ := pageThrough(t, repo, f.org.String(), models.SearchContacts{}, 50) + assertExactlyOnce(t, seen, f.total) + + filters := models.SearchContacts{CampaignIDs: []string{f.campaign.String()}} + rows, xerr := repo.ExportAll(context.Background(), f.org.String(), &filters, nil, 10000) + if xerr != nil { + t.Fatalf("export: %v", xerr) + } + ids := make([]uuid.UUID, 0, len(rows)) + for _, c := range rows { + ids = append(ids, c.ID) + } + assertExactlyOnce(t, ids, f.total) +} + +// A token minted under one ordering describes a position that does not exist +// under another, and a token from the old id-only format is not a position at +// all. Both are client contract errors, not silently wrong pages. +func TestLiveContactPaginationRejectsAForeignCursor(t *testing.T) { + handle, pool := liveContactDB(t) + f := newPagedOrgFixture(t, pool, 20) + repo := NewContactRepostory(handle) + ctx := context.Background() + + first, xerr := repo.Search(ctx, f.org.String(), nil, nil, models.SearchContacts{}, 5) + if xerr != nil { + t.Fatalf("first page: %v", xerr) + } + cursor, xerr := paging.DecodeSortCursor(*first.Pagination.NextCursor) + if xerr != nil { + t.Fatalf("decode: %v", xerr) + } + if _, xerr := repo.Search(ctx, f.org.String(), nil, cursor, models.SearchContacts{SortBy: "email"}, 5); xerr == nil { + t.Fatal("a created_at cursor must not be accepted on an email sort") + } + if _, xerr := repo.Search(ctx, f.org.String(), nil, &paging.SortCursor{Sort: "created_at:desc", ID: cursor.ID}, models.SearchContacts{}, 5); xerr == nil { + t.Fatal("a NULL boundary on a NOT NULL column must not be accepted") + } +} + +func assertExactlyOnce(t *testing.T, seen []uuid.UUID, total int) { + t.Helper() + unique := make(map[uuid.UUID]int, len(seen)) + for _, id := range seen { + unique[id]++ + } + dupes := 0 + for _, n := range unique { + if n > 1 { + dupes++ + } + } + if dupes > 0 { + t.Errorf("%d contacts were served more than once", dupes) + } + if len(unique) != total { + t.Errorf("saw %d distinct contacts (%d rows), want %d", len(unique), len(seen), total) + } +} diff --git a/internal/repository/pg_contact.go b/internal/repository/pg_contact.go index b2ce5163..c35f31d7 100644 --- a/internal/repository/pg_contact.go +++ b/internal/repository/pg_contact.go @@ -65,7 +65,7 @@ type ContactRepository interface { // to the caller's category IDs, creating the ones that don't exist yet. // Keys of the returned map are the lowercased titles. ResolveCategoryNames(ctx context.Context, userID uuid.UUID, names []string) (map[string]uuid.UUID, *errx.Error) - Search(ctx context.Context, userID string, category, cursor *string, filters models.SearchContacts, limit int32) (*models.ContactsResult, *errx.Error) + Search(ctx context.Context, userID string, category *string, cursor *paging.SortCursor, filters models.SearchContacts, limit int32) (*models.ContactsResult, *errx.Error) // SearchCounts returns org-wide contact facet totals for the browse // sidebar (independent of any search filters), mirroring campaigns-overview. SearchCounts(ctx context.Context, orgID string) (*models.ContactsCounts, *errx.Error) @@ -917,11 +917,40 @@ func (r *contactRepository) GetByIDsAndOrganization(ctx context.Context, organiz return out, nil } +// contactSort describes one sortable column of the contacts list: the SQL to +// order and compare on, the type a cursor's text value casts back to, and +// whether the column admits NULL (which decides where the NULL block sits in +// the order). Every column here is currently NOT NULL; flipping `nullable` +// turns on the NULL-aware keyset branches so a nullable sort cannot silently +// truncate the list. +type contactSort struct { + expr string + cast string + nullable bool +} + +// campaignCountLateral counts one contact's campaign memberships, joined only +// when a filter or the sort actually needs it. +const campaignCountLateral = `LEFT JOIN LATERAL ( + SELECT COUNT(*) AS campaign_count + FROM campaign_leads cl1 + WHERE cl1.contact_id = c.id + ) cl ON TRUE` + +var contactSorts = map[string]contactSort{ + "first_name": {expr: "c.first_name", cast: "text"}, + "last_name": {expr: "c.last_name", cast: "text"}, + "email": {expr: "c.email", cast: "text"}, + "created_at": {expr: "c.created_at", cast: "timestamp"}, + "updated_at": {expr: "c.updated_at", cast: "timestamp"}, + "campaign_count": {expr: "COALESCE(cl.campaign_count,0)", cast: "bigint"}, +} + func (r *contactRepository) Search( ctx context.Context, orgID string, - category, - cursor *string, + category *string, + cursor *paging.SortCursor, filters models.SearchContacts, limit int32, ) (*models.ContactsResult, *errx.Error) { @@ -1041,16 +1070,14 @@ func (r *contactRepository) Search( if len(filters.CampaignIDs) == 1 { singleCampaignPlaceholder = placeholders[0] } - campaignClause := fmt.Sprintf(` - c.id IN ( - SELECT contact_id - FROM campaign_leads - WHERE campaign_id IN (%s) - GROUP BY contact_id - HAVING COUNT(DISTINCT campaign_id) = %d - ) - `, strings.Join(placeholders, ","), len(filters.CampaignIDs)) - whereClauses = append(whereClauses, campaignClause) + // One EXISTS per campaign ("in ALL of them"), each a primary-key probe + // the planner can run either way round: driving from the ordered + // contacts index for a broad campaign, or from the leads for a narrow + // one. A GROUP BY/HAVING over campaign_leads forces the second. + for _, ph := range placeholders { + whereClauses = append(whereClauses, fmt.Sprintf( + "EXISTS (SELECT 1 FROM campaign_leads cle WHERE cle.campaign_id = %s AND cle.contact_id = c.id)", ph)) + } } // ----------------------------- @@ -1084,16 +1111,10 @@ func (r *contactRepository) Search( args = append(args, id) argIndex++ } - categoryClause := fmt.Sprintf(` - c.id IN ( - SELECT contact_id - FROM contact_categories - WHERE category_id IN (%s) - GROUP BY contact_id - HAVING COUNT(DISTINCT category_id) = %d - ) - `, strings.Join(placeholders, ","), len(filters.CategoryIDs)) - whereClauses = append(whereClauses, categoryClause) + for _, ph := range placeholders { + whereClauses = append(whereClauses, fmt.Sprintf( + "EXISTS (SELECT 1 FROM contact_categories cce WHERE cce.category_id = %s AND cce.contact_id = c.id)", ph)) + } } // ----------------------------- @@ -1123,50 +1144,67 @@ func (r *contactRepository) Search( // ----------------------------- // Sort logic // ----------------------------- - sortBy := "c.created_at" + // campaign_count is a computed column, so the cursor compares against the + // expression rather than the SELECT alias, which WHERE cannot see. + sortName := "created_at" + if _, ok := contactSorts[filters.SortBy]; ok { + sortName = filters.SortBy + } + spec := contactSorts[sortName] direction := "DESC" - allowedSorts := map[string]bool{ - "first_name": true, - "last_name": true, - "email": true, - "created_at": true, - "updated_at": true, - "campaign_count": true, - } - - if filters.SortBy != "" && allowedSorts[filters.SortBy] { - if filters.SortBy == "campaign_count" { - sortBy = "campaign_count" - } else { - sortBy = "c." + filters.SortBy - } - } + nulls := "NULLS FIRST" if filters.Reverse { direction = "ASC" - } else { - direction = "DESC" + nulls = "NULLS LAST" } + sortBy := spec.expr + // The ordering the cursor is taken under. A token minted under a different + // one points at a position that does not exist here. + sortKey := sortName + ":" + strings.ToLower(direction) // ----------------------------- // Cursor pagination // ----------------------------- - if cursor != nil && *cursor != "" { - cursorOp := ">" - if direction == "DESC" { - cursorOp = "<" + // Keyset, not offset: the token carries the boundary row's own sort value + // alongside its id, so a row deleted or re-sorted between pages cannot move + // the boundary. + if cursor != nil { + if cursor.Sort != sortKey { + return nil, errx.New(errx.BadRequest, "invalid cursor") } - sortSub := fmt.Sprintf("(SELECT %s FROM contacts WHERE id = $%d)", sortBy, argIndex) - args = append(args, *cursor) + if cursor.Value == nil && !spec.nullable { + return nil, errx.New(errx.BadRequest, "invalid cursor") + } + bound := fmt.Sprintf("$%d::%s", argIndex, spec.cast) + args = append(args, cursor.Value) + argIndex++ + idArg := fmt.Sprintf("$%d", argIndex) + args = append(args, cursor.ID) argIndex++ - whereClauses = append(whereClauses, fmt.Sprintf(` - ( - (%s %s %s) - OR (%s = %s AND c.id >= $%d) - ) - `, sortBy, cursorOp, sortSub, sortBy, sortSub, argIndex)) - args = append(args, *cursor) - argIndex++ + // The tiebreak follows the sort direction, so one index serves both ways + // round; the boundary row itself is included because it is this page's + // first row. + cmp, tie := "<", "<=" + if direction == "ASC" { + cmp, tie = ">", ">=" + } + after := fmt.Sprintf("(%[1]s %[2]s %[3]s OR (%[1]s = %[3]s AND c.id %[5]s %[4]s))", sortBy, cmp, bound, idArg, tie) + if spec.nullable { + switch { + case cursor.Value == nil: + // The boundary sits in the NULL block: the rest of that block by + // id, plus every non-NULL row when NULLs come first. + after = fmt.Sprintf("(%s IS NULL AND c.id %s %s)", sortBy, tie, idArg) + if nulls == "NULLS FIRST" { + after += fmt.Sprintf(" OR %s IS NOT NULL", sortBy) + } + case nulls == "NULLS LAST": + // Past the non-NULL rows, the NULL block still follows. + after += fmt.Sprintf(" OR %s IS NULL", sortBy) + } + } + whereClauses = append(whereClauses, "("+after+")") } // ----------------------------- @@ -1264,6 +1302,16 @@ func (r *contactRepository) Search( )`, singleCampaignPlaceholder, config.CampaignSendMaxAttempts, undeliverableClause(singleCampaignPlaceholder)) } + // campaign_count is only ever read by the min/max filters and the + // campaign_count sort; the response carries the campaign list itself. So the + // count is a lateral computed for the rows that survive, and it is left out + // entirely when nothing asks for it. Aggregating the whole campaign_leads + // table on every search was the list's dominant cost. + campaignCountJoin := "" + if filters.MinCampaigns != nil || filters.MaxCampaigns != nil || sortName == "campaign_count" { + campaignCountJoin = campaignCountLateral + } + // Main query. // // Both the `campaigns` and `categories` agg subqueries need the @@ -1277,7 +1325,6 @@ func (r *contactRepository) Search( c.custom_fields, c.subscribed, c.updated_at, c.created_at, c.verification_status, c.verification_reason, c.is_catch_all, c.verification_checked_at, c.verification_source, c.verification_provider, c.verification_sub_status, c.verification_confidence, - COALESCE(cl.campaign_count,0) AS campaign_count, COALESCE( ( SELECT json_agg(json_build_object('id', cam.id, 'name', cam.name)) @@ -1296,33 +1343,30 @@ func (r *contactRepository) Search( AND cat.user_id = $%d ), '[]'::json ) AS categories, - %s AS lead_progress + %s AS lead_progress, + (%s)::text AS sort_value FROM contacts c - LEFT JOIN ( - SELECT contact_id, COUNT(campaign_id) AS campaign_count - FROM campaign_leads - GROUP BY contact_id - ) cl ON c.id = cl.contact_id %s - ORDER BY %s %s, c.id ASC + %s + ORDER BY %s %s %s, c.id %s LIMIT $%d - `, argIndex, argIndex, leadProgressSelect, whereSQL, sortBy, direction, argIndex+1) + `, argIndex, argIndex, leadProgressSelect, sortBy, campaignCountJoin, whereSQL, sortBy, direction, nulls, direction, argIndex+1) args = append(args, orgID, limit+1) // Skip total count if cursor exists var totalCount *int64 - if cursor == nil || *cursor == "" { + if cursor == nil { + countJoin := "" + if filters.MinCampaigns != nil || filters.MaxCampaigns != nil { + countJoin = campaignCountLateral + } countQuery := fmt.Sprintf(` SELECT COUNT(*) FROM contacts c - LEFT JOIN ( - SELECT contact_id, COUNT(campaign_id) AS campaign_count - FROM campaign_leads - GROUP BY contact_id - ) cl ON c.id = cl.contact_id %s - `, whereSQL) + %s + `, countJoin, whereSQL) var tmp int64 if err := r.DB.QueryRow(ctx, countQuery, args[:argIndex-1]...).Scan(&tmp); err != nil { db.CaptureError(err, "countQuery", args, "queryrow") @@ -1346,12 +1390,15 @@ func (r *contactRepository) Search( // then produces [null], which crashes any downstream `.subscribed` // access. Always return an array. contacts := make([]models.Contact, 0, limit+1) + // The sort value of each row, kept alongside so the next cursor carries the + // boundary instead of re-reading it from a row that may be gone by then. + sortValues := make([]*string, 0, limit+1) for rows.Next() { var c models.Contact - var campaignCount int var campaignsJSON []byte var categoriesJSON []byte var leadProgressJSON []byte + var sortValue *string if err := rows.Scan( &c.ID, &c.FirstName, &c.LastName, &c.Email, @@ -1359,7 +1406,8 @@ func (r *contactRepository) Search( &c.UpdatedAt, &c.CreatedAt, &c.VerificationStatus, &c.VerificationReason, &c.IsCatchAll, &c.VerificationCheckedAt, &c.VerificationSource, &c.VerificationProvider, &c.VerificationSubStatus, &c.VerificationConfidence, - &campaignCount, &campaignsJSON, &categoriesJSON, &leadProgressJSON, + &campaignsJSON, &categoriesJSON, &leadProgressJSON, + &sortValue, ); err != nil { db.CaptureError(err, "", nil, "scan") return nil, errx.InternalError() @@ -1458,15 +1506,16 @@ func (r *contactRepository) Search( } contacts = append(contacts, c) + sortValues = append(sortValues, sortValue) } - // Next cursor + // Next cursor. The (limit+1)-th row is the first row of the NEXT page, so + // its position is the boundary and the id comparison is inclusive. var nextCursor *string var hasMore bool if len(contacts) > int(limit) { hasMore = true - nextID := contacts[limit].ID - nextCursor = paging.EncodeUUID(nextID) + nextCursor = paging.EncodeSort(sortKey, sortValues[limit], contacts[limit].ID) contacts = contacts[:limit] } @@ -2664,7 +2713,7 @@ func (r *contactRepository) ExportAll(ctx context.Context, orgID string, filters } out := make([]models.Contact, 0, 256) - var cursor *string + var cursor *paging.SortCursor pageSize := int32(500) for { page, xerr := r.Search(ctx, orgID, nil, cursor, search, pageSize) @@ -2685,14 +2734,13 @@ func (r *contactRepository) ExportAll(ctx context.Context, orgID string, filters if !page.Pagination.HasMore || page.Pagination.NextCursor == nil { break } - // NextCursor is now an opaque token; decode it back to the id the next - // Search call keys on. - id, derr := paging.DecodeUUID(*page.Pagination.NextCursor) - if derr != nil { + // NextCursor is an opaque token; decode it back to the keyset boundary + // the next Search call resumes from. + next, derr := paging.DecodeSortCursor(*page.Pagination.NextCursor) + if derr != nil || next == nil { break } - s := id.String() - cursor = &s + cursor = next } return out, nil } diff --git a/internal/utils/paging/paging.go b/internal/utils/paging/paging.go index 09df4688..0e8ff4a1 100644 --- a/internal/utils/paging/paging.go +++ b/internal/utils/paging/paging.go @@ -210,3 +210,81 @@ func DecodeMergedCursor(token string) (time.Time, int, uuid.UUID, *errx.Error) { } return at, source, id, nil } + +// sortPrefix versions the composite sort-keyset token: (sort key, sort value, +// id). Used by lists whose sort column is caller-chosen, nullable, or computed, +// where the boundary cannot be re-read from the row later: the row may be gone +// by the next page, and a mutable column (updated_at) may have moved. +const sortPrefix = "s1_" + +// SortCursor is the keyset boundary of a caller-sorted list: the first row of +// the next page, described by the ordering it was taken under, that row's sort +// value (nil when the column is NULL there) and its id, which breaks ties. +type SortCursor struct { + // Sort names the ordering (":"). A token replayed under a + // different ordering describes a position that does not exist there, so the + // caller rejects the mismatch instead of paging into nonsense. + Sort string + // Value is the boundary row's sort key rendered as text. The caller casts it + // back to the column's type, so the rendering has to round-trip. + Value *string + ID uuid.UUID +} + +// EncodeSort wraps a (sort key, sort value, id) keyset position in an opaque +// token. Returns nil for the zero id ("no next page") so the JSON field +// serializes as null. +func EncodeSort(sort string, value *string, id uuid.UUID) *string { + if id == uuid.Nil { + return nil + } + // Flag byte, then the fixed-width id, then the sort key; the value goes last + // so it keeps every byte it had, separators included. + flag := "0" + raw := "" + if value != nil { + flag = "1" + raw = *value + } + payload := flag + id.String() + "|" + sort + "|" + raw + tok := sortPrefix + base64.RawURLEncoding.EncodeToString([]byte(payload)) + return &tok +} + +// DecodeSortCursor reverses EncodeSort. An empty token yields (nil, nil) (start +// from the beginning); an invalid token returns a 400. +func DecodeSortCursor(token string) (*SortCursor, *errx.Error) { + if token == "" { + return nil, nil + } + invalid := errx.New(errx.BadRequest, "invalid cursor") + if !strings.HasPrefix(token, sortPrefix) { + return nil, invalid + } + raw, err := base64.RawURLEncoding.DecodeString(strings.TrimPrefix(token, sortPrefix)) + if err != nil { + return nil, invalid + } + head, rest, ok := strings.Cut(string(raw), "|") + if !ok || len(head) != 37 { + return nil, invalid + } + id, err := uuid.Parse(head[1:]) + if err != nil { + return nil, invalid + } + sort, value, ok := strings.Cut(rest, "|") + if !ok || sort == "" { + return nil, invalid + } + switch head[0] { + case '0': + if value != "" { + return nil, invalid + } + return &SortCursor{Sort: sort, ID: id}, nil + case '1': + return &SortCursor{Sort: sort, Value: &value, ID: id}, nil + } + return nil, invalid +} diff --git a/internal/utils/paging/paging_test.go b/internal/utils/paging/paging_test.go index 42abc95f..f6c82ae1 100644 --- a/internal/utils/paging/paging_test.go +++ b/internal/utils/paging/paging_test.go @@ -1,6 +1,7 @@ package paging import ( + "encoding/base64" "testing" "time" @@ -41,3 +42,79 @@ func TestMergedCursorRejectsMalformedTokens(t *testing.T) { } } } + +func TestSortCursorRoundTrip(t *testing.T) { + id := uuid.New() + val := "2026-06-09 11:42:00.123456" + tok := EncodeSort("created_at:desc", &val, id) + if tok == nil { + t.Fatal("want a token for a real position") + } + got, xerr := DecodeSortCursor(*tok) + if xerr != nil { + t.Fatalf("decode: %v", xerr) + } + if got.Sort != "created_at:desc" || got.ID != id || got.Value == nil || *got.Value != val { + t.Fatalf("round trip lost data: %+v", got) + } + if EncodeSort("created_at:desc", &val, uuid.Nil) != nil { + t.Fatal("the zero id means no next page and must encode as nil") + } +} + +// A NULL sort value is a real keyset position (the NULL block of a nullable +// column), and must survive the round trip as NULL rather than as "". +func TestSortCursorCarriesANullValue(t *testing.T) { + id := uuid.New() + got, xerr := DecodeSortCursor(*EncodeSort("first_name:asc", nil, id)) + if xerr != nil { + t.Fatalf("decode: %v", xerr) + } + if got.Value != nil || got.ID != id || got.Sort != "first_name:asc" { + t.Fatalf("want a NULL boundary, got %+v", got) + } + empty := "" + got, xerr = DecodeSortCursor(*EncodeSort("first_name:asc", &empty, id)) + if xerr != nil { + t.Fatalf("decode empty: %v", xerr) + } + if got.Value == nil || *got.Value != "" { + t.Fatalf("an empty string is not NULL, got %+v", got) + } +} + +// Values are user data: separators in them must not shift the fields. +func TestSortCursorKeepsSeparatorsInTheValue(t *testing.T) { + id := uuid.New() + val := "a|b|c" + got, xerr := DecodeSortCursor(*EncodeSort("last_name:asc", &val, id)) + if xerr != nil { + t.Fatalf("decode: %v", xerr) + } + if got.Value == nil || *got.Value != val { + t.Fatalf("value = %v, want %q", got.Value, val) + } +} + +func TestSortCursorRejectsMalformedTokens(t *testing.T) { + if c, xerr := DecodeSortCursor(""); xerr != nil || c != nil { + t.Fatalf("an empty token is the first page, got %v %v", c, xerr) + } + for _, tok := range []string{ + "s1_", // no payload + uuid.New().String(), // a bare id is not a cursor + "c1_" + (*EncodeUUID(uuid.New()))[3:], // the id-only format + "s1_!!!", // not base64 + "s1_" + base64Raw("1"+uuid.New().String()), // no sort key + "s1_" + base64Raw("2"+uuid.New().String()+"|email:asc|x"), // unknown flag + "s1_" + base64Raw("1nope|email:asc|x"), // not an id + } { + if _, xerr := DecodeSortCursor(tok); xerr == nil { + t.Fatalf("%q must be rejected", tok) + } + } +} + +func base64Raw(s string) string { + return base64.RawURLEncoding.EncodeToString([]byte(s)) +} diff --git a/web/src/components/app/contacts/ContactsTable.tsx b/web/src/components/app/contacts/ContactsTable.tsx index bb35be32..3ce26e31 100644 --- a/web/src/components/app/contacts/ContactsTable.tsx +++ b/web/src/components/app/contacts/ContactsTable.tsx @@ -544,6 +544,8 @@ export default function ContactsTable({ hasNextPage={!!contactsData.hasNextPage} isFetchingNextPage={contactsData.isFetchingNextPage} onLoadMore={() => contactsData.fetchNextPage()} + loadedCount={contacts?.length ?? 0} + totalCount={total} /> ); @@ -958,6 +960,8 @@ function ContactsTableBody({ hasNextPage, isFetchingNextPage, onLoadMore, + loadedCount, + totalCount, }: { embedded?: boolean; isLoading: boolean; @@ -999,6 +1003,9 @@ function ContactsTableBody({ hasNextPage: boolean; isFetchingNextPage: boolean; onLoadMore: () => void; + // How far through the list we are, so "Load more" says how much is left. + loadedCount: number; + totalCount: number; }) { if (isLoading) { return ( @@ -1015,7 +1022,10 @@ function ContactsTableBody({ ); } - if (isError) { + // A page that fails after rows are already on screen is reported in the + // footer instead; throwing away 500 loaded leads because page 13 failed is + // worse than the failure. + if (isError && contacts.length === 0) { return (
@@ -1297,7 +1307,30 @@ function ContactsTableBody({ })} - {hasNextPage && ( + {isError ? ( +
+

+ + {errorMessage} +

+ +
+ ) : hasNextPage ? (
- )} + ) : null} ); }