From df62ba70a1ea7d71e76d4b59b25502cf38fcd3af Mon Sep 17 00:00:00 2001 From: Matthew Meszaros Date: Wed, 23 Sep 2026 07:46:49 -0700 Subject: [PATCH] feat: write a message notification read with no email when its message is already read (locked against a concurrent read), cancel an email-only notification's pending email when its message is read, and only drop Slack and push when the message left the unibox rather than on any insert error --- internal/app/notification/service.go | 6 +- ...000204_notification_unibox_email_fk.up.sql | 9 ++- internal/models/notification.go | 2 + internal/repository/pg_notification.go | 24 ++++++- .../unibox_notification_live_test.go | 64 ++++++++++++++++++- 5 files changed, 96 insertions(+), 9 deletions(-) diff --git a/internal/app/notification/service.go b/internal/app/notification/service.go index c03b1d6ff..d44e44465 100644 --- a/internal/app/notification/service.go +++ b/internal/app/notification/service.go @@ -7,6 +7,7 @@ package notification import ( "context" + "errors" "time" "github.com/google/uuid" @@ -232,9 +233,12 @@ func (s *service) notifyOne(ctx context.Context, userID uuid.UUID, orgID *uuid.U n.EmailDueAt = &due } created, cerr := s.repo.Create(ctx, n) - if cerr != nil && uniboxEmailID != nil { + if errors.Is(cerr, repository.ErrNotificationMessageGone) { return false // the message left the unibox first; nothing to announce } + if cerr == nil && created != nil && created.MessageSeen { + return false // already read where it arrived; the row is the record + } if cerr == nil && created != nil && cat.Channels.InApp && s.publisher != nil { s.publisher.PublishNotificationCreated(ctx, userID.String(), created.ID.String(), string(category), title, link) } diff --git a/internal/infrastructure/db/migrations/000204_notification_unibox_email_fk.up.sql b/internal/infrastructure/db/migrations/000204_notification_unibox_email_fk.up.sql index 32fc0c5fd..b8c627473 100644 --- a/internal/infrastructure/db/migrations/000204_notification_unibox_email_fk.up.sql +++ b/internal/infrastructure/db/migrations/000204_notification_unibox_email_fk.up.sql @@ -11,10 +11,13 @@ CREATE FUNCTION public.notifications_read_with_message() RETURNS trigger LANGUAGE plpgsql AS $$ BEGIN + -- An email-only notification is born read, so its pending email is + -- cancelled whatever read_at says. UPDATE public.notifications - SET read_at = now(), - email_state = CASE WHEN email_state = 'pending' THEN 'skipped' ELSE email_state END - WHERE unibox_email_id = NEW.id AND read_at IS NULL; + SET read_at = COALESCE(read_at, now()), + email_state = CASE WHEN email_state = 'pending' THEN 'skipped' ELSE email_state END, + email_due_at = CASE WHEN email_state = 'pending' THEN NULL ELSE email_due_at END + WHERE unibox_email_id = NEW.id AND (read_at IS NULL OR email_state = 'pending'); RETURN NULL; END; $$; diff --git a/internal/models/notification.go b/internal/models/notification.go index 16fe62322..7d307a3e5 100644 --- a/internal/models/notification.go +++ b/internal/models/notification.go @@ -145,6 +145,8 @@ type Notification struct { // loop can coalesce it into one email with every recipient in To. GroupKey string `json:"-"` UniboxEmailID *uuid.UUID `json:"-"` + // MessageSeen reports, on create, that the message was already read. + MessageSeen bool `json:"-"` EmailState string `json:"-"` EmailDueAt *time.Time `json:"-"` EmailAttempts int `json:"-"` diff --git a/internal/repository/pg_notification.go b/internal/repository/pg_notification.go index 7fa0d326f..31aa0e580 100644 --- a/internal/repository/pg_notification.go +++ b/internal/repository/pg_notification.go @@ -8,6 +8,7 @@ import ( "github.com/google/uuid" "github.com/jackc/pgx/v5" + "github.com/jackc/pgx/v5/pgconn" "github.com/jackc/pgx/v5/pgxpool" "github.com/warmbly/warmbly/internal/config" @@ -103,19 +104,36 @@ func (r *notificationRepository) Create(ctx context.Context, n *models.Notificat if n.GroupKey != "" { groupKey = &n.GroupKey } + // A message already read is announced as read, with no email. FOR SHARE + // waits out a read in flight, so the read trigger cannot miss this row. err := r.db.QueryRow(ctx, ` + WITH msg AS (SELECT seen FROM unibox_emails WHERE id = $13 FOR SHARE), + seen AS (SELECT COALESCE((SELECT seen FROM msg), false) AS v) INSERT INTO notifications (id, user_id, organization_id, category, title, body, link, metadata, group_key, email_state, email_due_at, read_at, unibox_email_id) - VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, CASE WHEN $12 THEN now() END, $13) - RETURNING created_at`, + SELECT $1, $2, $3, $4, $5, $6, $7, $8, $9, + CASE WHEN seen.v AND $10 = 'pending' THEN 'skipped' ELSE $10 END, + CASE WHEN seen.v AND $10 = 'pending' THEN NULL ELSE $11::timestamptz END, + CASE WHEN $12 OR seen.v THEN now() END, + $13 + FROM seen + RETURNING created_at, (SELECT v FROM seen)`, n.ID, n.UserID, n.OrganizationID, n.Category, n.Title, n.Body, n.Link, meta, - groupKey, n.EmailState, n.EmailDueAt, n.PreRead, n.UniboxEmailID).Scan(&n.CreatedAt) + groupKey, n.EmailState, n.EmailDueAt, n.PreRead, n.UniboxEmailID).Scan(&n.CreatedAt, &n.MessageSeen) if err != nil { + var pgErr *pgconn.PgError + if errors.As(err, &pgErr) && pgErr.ConstraintName == "notifications_unibox_email_id_fkey" { + return nil, ErrNotificationMessageGone + } return nil, err } return n, nil } +// ErrNotificationMessageGone is a message notification whose message left the +// unibox before the notification was written. +var ErrNotificationMessageGone = errors.New("notification: message no longer in the unibox") + func (r *notificationRepository) List(ctx context.Context, userID uuid.UUID, limit int, unreadOnly bool) ([]models.Notification, error) { if limit <= 0 || limit > 100 { limit = 50 diff --git a/internal/repository/unibox_notification_live_test.go b/internal/repository/unibox_notification_live_test.go index c125156d3..42a672f66 100644 --- a/internal/repository/unibox_notification_live_test.go +++ b/internal/repository/unibox_notification_live_test.go @@ -2,6 +2,7 @@ package repository import ( "context" + "errors" "testing" "time" @@ -142,7 +143,66 @@ func TestLiveUniboxNotificationLeavesWithTheMessage(t *testing.T) { if _, err := notifs.Create(ctx, &models.Notification{ UserID: f.user, OrganizationID: &f.org, Category: models.NotifInboundReply, Title: "New reply", UniboxEmailID: &id, - }); err == nil { - t.Fatal("created a notification about a message no longer in the unibox") + }); !errors.Is(err, ErrNotificationMessageGone) { + t.Fatalf("error = %v, want ErrNotificationMessageGone", err) + } +} + +// A message read before its notification is written (read in the mail client +// before the sync, or in a race with the reader) is announced as read, and +// an email-only notification is cancelled when its message is read. +func TestLiveUniboxNotificationRespectsTheMessagesReadState(t *testing.T) { + handle := liveUniboxFolderDB(t) + f := newUniboxFolderFixture(t, handle.Pool) + repo := NewUniboxRepository(handle) + notifs := NewNotificationRepository(handle.Pool) + ctx := context.Background() + now := time.Now().UTC() + due := now.Add(time.Hour) + + state := func(id uuid.UUID) (bool, string) { + var read bool + var email string + if err := handle.Pool.QueryRow(ctx, + `SELECT read_at IS NOT NULL, email_state FROM notifications WHERE id = $1`, id).Scan(&read, &email); err != nil { + t.Fatalf("read notification: %v", err) + } + return read, email + } + + already := f.scopedMessage(t, repo, "thread-already-read", "them@example.com", models.FolderInbox, now) + if _, err := repo.MarkSeenByThreads(ctx, f.org, []string{"thread-already-read"}, true); err != nil { + t.Fatalf("MarkSeenByThreads: %v", err) + } + n, err := notifs.Create(ctx, &models.Notification{ + UserID: f.user, OrganizationID: &f.org, Category: models.NotifInboundReply, + Title: "New reply", UniboxEmailID: &already, EmailState: "pending", EmailDueAt: &due, + }) + if err != nil { + t.Fatalf("Create: %v", err) + } + if !n.MessageSeen { + t.Error("Create did not report the message as already read") + } + if read, email := state(n.ID); !read || email != "skipped" { + t.Errorf("already-read message: read=%v email=%q, want read with no email", read, email) + } + + emailOnly := f.scopedMessage(t, repo, "thread-email-only", "them@example.com", models.FolderInbox, now) + n, err = notifs.Create(ctx, &models.Notification{ + UserID: f.user, OrganizationID: &f.org, Category: models.NotifInboundReply, + Title: "New reply", UniboxEmailID: &emailOnly, EmailState: "pending", EmailDueAt: &due, PreRead: true, + }) + if err != nil { + t.Fatalf("Create email-only: %v", err) + } + if n.MessageSeen { + t.Error("an unread message was reported as read") + } + if _, err := repo.MarkSeenByThreads(ctx, f.org, []string{"thread-email-only"}, true); err != nil { + t.Fatalf("MarkSeenByThreads: %v", err) + } + if read, email := state(n.ID); !read || email != "skipped" { + t.Errorf("email-only notification: read=%v email=%q, want its email cancelled by the read", read, email) } }