mirror of
https://github.com/warmbly/warmbly.git
synced 2026-10-03 16:02:02 +00:00
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
This commit is contained in:
@@ -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)
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
$$;
|
||||
|
||||
@@ -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:"-"`
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user